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.