3 · Under the hood · lesson 13 of 20
The Shuffle
The expensive redistribution of data across the network — and how to avoid it.
Any operation that needs to co-locate rows by key — groupBy, join, distinct, window, repartition — has to shuffle: every executor writes intermediate files, then every executor on the next stage reads the pieces it needs.
Python
# Watch for these operators in the query plan — each is a shuffle
df.groupBy("country").agg(...) # HashAggregate + Exchange
df.join(other, "id") # SortMergeJoin + Exchange
df.distinct() # Aggregate + Exchange
df.repartition(64, "user_id") # Exchange
# Broadcast small sides to skip the shuffle entirely
df.join(F.broadcast(small_dim), "key")WATCH OUT
Shuffles are usually the #1 performance cost. Cut them by broadcasting small tables, filtering early, and pre-partitioning frequently-joined data.
Loading 3D scene…
Key takeaways
- ✓Shuffle = disk write on the sender + network read on the receiver.
- ✓Broadcast joins and partition pruning are the biggest wins.
- ✓explain() shows every Exchange (= shuffle) in the plan.