A shuffle redistributes data across the cluster so that records that must be processed together (same join key, same group key) land on the same partition/executor.
Before shuffle (partition by file): After shuffle (by user_id): exec1: u1,u7,u2 exec1: all u1,u4,... exec2: u4,u1,u9 ======> exec2: all u2,u7,... exec3: u2,u4 exec3: all u9,...
Why it hurts
- Network transfer of large bytes
- Disk spill when memory is tight
- Sort/hash overhead
- Sensitive to skew (one key owns 80% of rows)
DE habits
- Filter early before shuffle
- Broadcast small dimension tables to avoid big shuffles
- Salting / skew hints for hot keys
- Right-size partitions (
spark.sql.shuffle.partitions)
Interview tip: Define shuffle as "repartition by key across the network," then name join/groupBy as causes and skew as the failure mode.