Shuffle internals
Map output files, block fetches, fetch failures and the external shuffle service.
On this page
The map side: writing shuffle files
-
Partitionpartition: A chunk of a DataFrame's rows. Spark processes each partition as one task, so partitions decide how much work runs in parallel. Learn more →
Each row gets a target reduce partition: hash(key) mod N
-
Buffer
Rows are buffered in memory, spillingspill: Writing data to local disk because it does not fit in memory. Slower, but the job keeps running. Learn more → sorted runs to disk if needed
-
Sort and merge
Runs are merged, ordered by partition id
-
Write
One .data file and one .index file of offsets per map tasktask: The work for one partition in one stage, run on one CPU core. Learn more →
Early Spark versions wrote one file per map task per reducer, so 1,000 map tasks × 1,000 reducers meant a million files. The sort-based shuffleshuffle: Moving rows between machines so that all rows with the same key end up together. Needed by joins, groupBy and sorting, and usually the most expensive step of a job. Learn more →, the only one since Spark 2.0, writes two files per map task no matter how many reducers there are. The index file records where each reducer's bytes start and end inside the data file.
| Writer | Used when | How |
|---|---|---|
| Bypass merge sort | No map-side combine and at most spark.shuffle.sort.bypassMergeThreshold (200) reduce partitions | Writes one temporary file per reducer, then concatenates them. DataFrame shuffles into 200 partitions often take this path. |
| Unsafe (serialized) shuffle | Serializer supports relocating serialized records (as Spark SQL's row serializer does), no map-side combine | Sorts compact pointers to serialized binary rows instead of objects: less memory and GC. |
| Sort shuffle | Everything else, such as RDD reduceByKey with a map-side combine | Sorts and combines deserialized objects, spilling as needed. |
The reduce side: fetching blocks
- 1The reducer asks the driverdriver: The process that runs your program, plans the work and sends tasks to executors. Learn more → where each map output lives.
- 2It reads local blocks directly from disk and requests remote blocks from other executorsexecutor: A worker process on a cluster machine that runs tasks and holds data in memory. Learn more →, or from the shuffle service on their node, keeping at most
spark.reducer.maxSizeInFlight(48 MB) of requests in flight. - 3Fetched blocks are decompressed and deserialized, then fed to the join, aggregation or sort that follows, which may itself spill.
With M map tasks and N reducers, that is up to M × N small blocks. When both are large, the blocks get small (a 1 TB shuffle with 10,000 × 10,000 tasks gives 10 KB blocks) and the shuffle becomes dominated by random disk reads and network round trips rather than bandwidth.
Fetch failures
If a reducer cannot fetch a block because the executor or node that wrote it is gone, the task fails with a FetchFailedException. That is not retried like an ordinary task failure: Spark marks the map output as lost, re-runs the map tasks that produced it, then re-runs the reduce stagestage: A group of steps Spark can run without moving data between machines. A new stage starts at every shuffle. Learn more →. After spark.stage.maxConsecutiveAttempts (4) consecutive failed attempts, the job fails. Lost executors under memory pressure, preempted spot instances and overloaded shuffle services are the usual causes.
The external shuffle service
By default an executor serves its own shuffle files, so when it dies, they become unreachable. The external shuffle service (spark.shuffle.service.enabled) is a separate long-running process on each node that serves the files instead. Executors can then be removed, by dynamic allocation or a crash, without losing finished map output, as long as the node and its disk survive.
spark.shuffle.push.enabled) has map tasks also push their blocks to remote shuffle services, which merge them into one larger file per reduce partition. Reducers then read a few large files instead of thousands of tiny blocks, and a second copy of the data exists if a node fails.Reading shuffle metrics
| Spark UI metric | What it tells you |
|---|---|
| Shuffle Write Size / Records | Bytes and rows the stage produced for the next stage. Compare with the input to see if filtering or pre-aggregation worked. |
| Shuffle Read Size (local / remote) | Bytes fetched by the stage; most should be remote on a real cluster. |
| Shuffle Read Fetch Wait Time | Time tasks were blocked waiting for blocks. High values point at network, disk or overloaded shuffle services. |
| Spill (memory / disk) | Data that did not fit during the shuffle's sort or the operator after it. |
| Skipped stages | Stages whose shuffle output already existed, so they did not run again. Common with AQEAQE: Adaptive Query Execution: Spark re-plans a running query using the real data sizes it has measured. Learn more →, which runs each query stage as its own job. |
Configs worth knowing
| Config | Default | Notes |
|---|---|---|
spark.shuffle.compress | true | Compress map output, with spark.io.compression.codec (lz4) |
spark.local.dir | /tmp | Where shuffle and spill files go. Point it at fast local disks; "No space left on device" errors are about this, not HDFS or S3. |
spark.reducer.maxSizeInFlight | 48m | Larger means fewer, bigger fetch rounds, at the cost of memory |
spark.shuffle.io.maxRetries / retryWait | 3 / 5s | Retries before a fetch is declared failed; raise on flaky networks or during long GC pauses |
Common mistakes
- Treating a FetchFailedException as the root cause — It is a symptom. Look for why the executor that wrote the file died: OOM kills, spot loss, decommissioning.
- Tiny shuffle partitions on huge shuffles — M × N tiny blocks make fetching slow. Let AQE coalesce, or size shuffle partitions to 100 to 200 MB.
- Small local disks — Shuffle and spill files fill them and tasks fail, even when the output goes to object storage.
What you learned
- What a map task writes during a shuffle, file by file
- How reducers fetch shuffle blocks, and what a fetch failure triggers
- The external shuffle service and push-based shuffle
- The shuffle metrics and configs worth knowing
Key takeaways
- Sort-based shuffle writes one data file and one index file per map task.
- Reducers fetch their byte ranges from every map output; M × N blocks in total.
- Fetch failures make Spark recompute lost map output, then retry the stage.
- The external shuffle service keeps shuffle files available after executors leave.
- Watch shuffle write size, fetch wait time and spill in the Spark UI.
Check yourself
3 questionsHow many files does a sort-based shuffle map task write, with 1,000 reducers?
Show the answer
One data file and one index file. The index file records each reducer's byte range in the single data file.
What does Spark do on a FetchFailedException?
Show the answer
Re-runs the map tasks whose output was lost, then the reduce stage. The data is gone, so it must be recomputed from lineage.
What does the external shuffle service provide?
Show the answer
Shuffle files served independently of the executor that wrote them. Files stay available when an executor is removed.
Keep going
Up next · lesson 22 of 30 · 4 min readHow aggregations execute
Partial and final aggregation, hash vs sort aggregation, and why COUNT(DISTINCT) costs more.
Related lessons
Narrow vs wide transformationsSpark internals · 4 min read
Dynamic allocationSpark internals · 3 min read
Shuffle partitions
Previous: The memory model
Primary sources: RDD guide: shuffle operations · Configuration: shuffle behavior