Shuffles: Advanced
Sort-Based Shuffle Internals
| Hash-based (old) | Sort-based (default) | |
|---|---|---|
| Files per map task | One per reduce partition (R) | One data + one index (2) |
| Total files | M x R (millions at scale) | 2M (linear in map tasks) |
| How reduce finds its data | Open its own file from each map | Seek by offset using the index |
| Why it matters | Collapsed under file-handle pressure | Scales to large clusters |
The External Shuffle Service
Pricing a Shuffle
| Lever on shuffle cost | What it changes | How to pull it |
|---|---|---|
| Bytes entering the shuffle | The base volume everything scales from | Filter and pre-aggregate before the wide op |
| Partition count | Per-partition size; spill risk | Size to roughly data over 128MB |
| Number of shuffles | How many times you pay the whole stack | Restructure to remove wide ops |
| Skew | Whether one partition stalls the barrier | Salt or isolate hot keys |
Shrink the bytes before the shuffle
(order_items
.filter(F.col("quantity") > 2)
.groupBy("product_id")
.agg(F.sum("unit_price").alias("total"))
.orderBy(F.col("total").desc()))Eliminating a Shuffle
- Every row crosses the network
- The shuffle moves the full dataset
- Expensive, and risks spill and skew
- No reduction before the network
- Partial-aggregate on each partition first
- Only the small partials shuffle
- Orders of magnitude less data moved
- The default for DataFrame aggregations
Broadcast the small side
> Join order_items to the small products table, broadcasting products so the large side never shuffles, and return total quantity per category, highest first. The broadcast is what removes the join shuffle.
(order_items .join((products), "product_id") .groupBy("category") .agg(F.("quantity").alias("units")) .orderBy(F.col("units").desc(), F.col("category").asc()))
The Shuffle Tuning Knobs
| Knob | What it controls | When to reach for it |
|---|---|---|
| Shuffle compression | Whether shuffle files are compressed | On by default; trades CPU for less disk and network |
| Shuffle buffer sizes | Memory used for the sort and fetch | Raise to reduce spill on large partitions |
| Fetch concurrency | How many buckets a reduce task pulls at once | Tune for network throughput vs memory |
| spark.sql.shuffle.partitions | Post-shuffle partition count | The first and biggest knob; size to the data |
- Eliminate a shuffle first: map-side combine, broadcast a small side, or bucket at write time.
- Reduce the bytes entering a shuffle by filtering and pre-aggregating before the wide op.
- Size the partition count so nothing spills before reaching for buffer knobs.
- Rely on the external shuffle service so a lost executor does not destroy shuffle output.
- Don't reach for buffer and compression knobs before trying to remove the shuffle.
- Don't groupByKey when a map-side-combine aggregation moves orders of magnitude less data.
- Don't shuffle a large table to join it when the other side is small enough to broadcast.
- Don't re-bucket a table on every read; pay the reorganisation once at write time.
> A pipeline joins a 500GB fact table to a 50MB dimension and then aggregates, and it shuffles both the join and the aggregation, missing its SLA. The same fact table is joined this way in several downstream jobs.
The cheapest shuffle is the one you engineered away.
- Category
- Spark
- Difficulty
- advanced
- Duration
- 15 minutes
- Challenges
- 2 hands-on challenges
Topics covered: Sort-Based Shuffle Internals, The External Shuffle Service, Pricing a Shuffle, Eliminating a Shuffle, The Shuffle Tuning Knobs
Lesson Sections
- Sort-Based Shuffle Internals (concepts: paShuffleOptimization)
Early Spark had a shuffle implementation that did not scale, and understanding why it was replaced explains the design you rely on today. The original hash-based shuffle had each map task write a separate file for each reduce partition. With M map tasks and R reduce partitions, that is M times R files, and on a large cluster M times R is millions of tiny files. The operating system buckled under the file handles and the random I/O of opening millions of small files crushed performance. Sort-base
- The External Shuffle Service (concepts: paShuffleOptimization)
There is a problem with shuffle files living on an executor's local disk: what happens to them when that executor dies or is taken away? The reduce tasks still need those buckets, but the process that wrote them is gone. Without a solution, losing an executor mid-shuffle would force the map tasks that ran on it to be recomputed, the kind of expensive cascade you want to avoid. The external shuffle service is that solution. The external shuffle service is a separate, long-lived process that runs
- Pricing a Shuffle (concepts: paShuffleOptimization)
A senior engineer can estimate a shuffle's cost before running it, and the estimate starts from one number: how many bytes the shuffle moves. That is roughly the size of the data entering the wide operation, possibly reduced if a pre-aggregation shrinks it first. Knowing the bytes, you can reason about the wall-clock, because the bytes have to be written to disk, sent over the network, and read back, and each of those has a rate you can ballpark. You are not computing an exact number. You are bu
- Eliminating a Shuffle (concepts: paBroadcastJoin)
Several ways exist to restructure a job so a shuffle you expected does not occur. These are the strongest optimisations in Spark, because removing a shuffle removes the entire stacked cost at once rather than a piece of it. 3 techniques cover most cases, and a strong candidate can name all 3. The first is map-side combine, and it is the classic reduceByKey versus groupByKey distinction. If you are aggregating, you can often reduce within each partition before the shuffle, so that only the partia
- The Shuffle Tuning Knobs (concepts: paShuffleOptimization)
When you cannot eliminate a shuffle, you tune it, and there is a small set of knobs beyond the partition count that shape how a shuffle performs. None of them is as powerful as removing the shuffle or sizing the partitions, but together they trim the edges, and knowing they exist is part of a complete picture. Shuffle compression is on by default and usually worth keeping, because the CPU cost of compressing is small next to the disk and network it saves; the data being shuffled is often compres