A shuffle is the redistribution of data across partitions so that rows that belong together (same key, same range) land on the same executor partition. It is the expensive part of wide transformations.
What happens under the hood
Map side (write shuffle files) Reduce side (read + compute) ------------------------------ ---------------------------- Executor tasks hash/sort keys --> network/disk fetch --> tasks aggregate/join
1. Map tasks write shuffle blocks partitioned by key hash (or range). 2. Data moves over the network (and often spills to disk). 3. Reduce tasks read their assigned blocks and continue the plan.
Operations that shuffle
groupBy/agg, most joins, distinct, repartition(n), orderBy, window functions with partitioning, cogroup, etc.
Cost drivers
- Data volume and skew
- Number of shuffle partitions (
spark.sql.shuffle.partitions, AQE coalescing) - Serializer, compression, disk spill
- Network bandwidth
Example
# Shuffle on user_id
revenue = orders.groupBy("user_id").sum("amount")
# Often avoid shuffle with broadcast if one side is small
from pyspark.sql.functions import broadcast
joined = orders.join(broadcast(dim_users), "user_id")Interview closer
"Shuffle is necessary for correctness when keys must meet. The craft is cutting shuffled bytes and handling skew."