repartition(n) shuffles rows into n even partitions, repartition(n, col) hashes by a column so equal keys meet, and coalesce(n) merges partitions without a shuffle (it can only reduce their number).
import contextlib, io
def layout(label, df):
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
df.explain()
per_part = df.groupBy(F.spark_partition_id().alias("p")).count()
lo, hi, n = per_part.agg(F.min("count"), F.max("count"), F.count("*")).first()
print(f"{label:<23} {n:>2} partitions rows {lo:>9,} to {hi:>9,} "
f"shuffle: {'Exchange' in buf.getvalue()}")
layout("as read (4 files)", lines)
layout("repartition(8)", lines.repartition(8))
layout("coalesce(2)", lines.coalesce(2))
layout("repartition(8, book_id)", lines.repartition(8, "book_id"))Output
as read (4 files) 4 partitions rows 326,701 to 351,848 shuffle: False repartition(8) 8 partitions rows 172,586 to 172,588 shuffle: True coalesce(2) 2 partitions rows 677,916 to 702,778 shuffle: False repartition(8, book_id) 4 partitions rows 142,075 to 684,397 shuffle: True
Round-robin repartition(8) balanced rows to within two. Hashing six book IDs into eight partitions left four empty and one with 684,397 rows where two popular books collided: Diagnosing Data Skew's skew again. Use coalesce before writing small outputs, repartition(col) before grouping or writing by col; aim for 100 to 200 MB partitions, two or three per core.