Repartition and Coalesce

Repartition, Coalesce and Partition Sizing

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).

What repartition and coalesce do to the order linesJavaScript
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.