spark.sql.shuffle.partitions (200) fixes how many reduce tasks a shuffle produces, whatever the data size. Too few, and each task holds too much and spills; too many, and task overhead and tiny files dominate. The listing aggregates the order lines by customer under four settings.
import time
sc = spark.sparkContext
def run(aqe, n):
spark.conf.set("spark.sql.adaptive.enabled", aqe)
spark.conf.set("spark.sql.shuffle.partitions", n)
sc.setJobGroup(aqe + n, "shuffle tuning")
t0 = time.perf_counter()
(lines.groupBy("customer_id").agg(F.sum("qty"), F.count("*")) # 50,000 groups
.write.format("noop").mode("overwrite").save())
stages = [s for j in sc.statusTracker().getJobIdsForGroup(aqe + n)
for s in sc.statusTracker().getJobInfo(j).stageIds]
return time.perf_counter() - t0, sc.statusTracker().getStageInfo(max(stages)).numTasks
for aqe, n in [("true", "200"), ("false", "200"), ("false", "8"), ("false", "2000")]:
run(aqe, n) # warm-up
secs, tasks = run(aqe, n)
print(f"AQE {aqe:<5} shuffle.partitions {n:>4}: {secs:5.2f}s, {tasks:>4} reduce tasks")Output
AQE true shuffle.partitions 200: 2.93s, 2 reduce tasks AQE false shuffle.partitions 200: 4.59s, 200 reduce tasks AQE false shuffle.partitions 8: 0.94s, 8 reduce tasks AQE false shuffle.partitions 2000: 17.49s, 2000 reduce tasks
With AQE (the default) the 200 planned partitions were coalesced into 2. Without it, 200 tasks cost five times the hand-tuned 8, and 2,000 cost almost twenty times. The rule: keep AQE on, set spark.sql.shuffle.partitions high enough for your largest shuffle (about shuffle bytes divided by 100 to 200 MB), and let AQE coalesce the small ones; check advisoryPartitionSizeInBytes (64 MB) when tasks spill.