A shuffle is when a distributed engine must redistribute data across workers so that records with the same key land on the same machine. It is network + disk heavy.
Before shuffle (data local to workers): Worker A: (user=1, ...), (user=7, ...) Worker B: (user=1, ...), (user=3, ...) After shuffle for GROUP BY user_id: Worker A: all user=1 rows Worker B: all user=3,7 rows
Operations that trigger shuffles
GROUP BY,DISTINCT, many joinsORDER BY/ global sortrepartition(key)
Why shuffles hurt
Bytes cross the network, stages wait on the slowest worker, skew amplifies pain, and spills to disk if memory is tight.
How to reduce shuffled bytes
- Filter early
- Use columnar formats and only needed columns
- Broadcast small tables
- Prefer map-side combiners / partial aggregates
- Avoid unnecessary
repartition
Interview tip: "Shuffle = redistributing records by key across the cluster." Say it is often necessary for correctness; the skill is minimizing shuffled bytes and handling skew.