Skip to content
Great engineers know 6 min read · spark.shuffle.service.enabled

Shuffle internals

Map output files, block fetches, fetch failures and the external shuffle service.

TL;DR Spark uses a sort-based shuffle: each map task sorts its output by target partition and writes one data file plus one index file to local disk. Each reduce task then fetches its slice from every map output over the network. If those files are lost, Spark re-runs the map stage that produced them, which is why shuffle files must outlive executors when executors come and go.
Read first: Narrow vs wide transformations · Jobs, stages and tasks (skip if you know them)

The map side: writing shuffle files

  1. 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

  2. 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

  3. Sort and merge

    Runs are merged, ordered by partition id

  4. 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.

WriterUsed whenHow
Bypass merge sortNo map-side combine and at most spark.shuffle.sort.bypassMergeThreshold (200) reduce partitionsWrites one temporary file per reducer, then concatenates them. DataFrame shuffles into 200 partitions often take this path.
Unsafe (serialized) shuffleSerializer supports relocating serialized records (as Spark SQL's row serializer does), no map-side combineSorts compact pointers to serialized binary rows instead of objects: less memory and GC.
Sort shuffleEverything else, such as RDD reduceByKey with a map-side combineSorts and combines deserialized objects, spilling as needed.

The reduce side: fetching blocks

  1. 1
    The 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.
  2. 2
    It 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.
  3. 3
    Fetched 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.

Note: Push-based shuffle (Spark 3.2+, YARN with the external shuffle service, 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 metricWhat it tells you
Shuffle Write Size / RecordsBytes 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 TimeTime 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 stagesStages 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

ConfigDefaultNotes
spark.shuffle.compresstrueCompress map output, with spark.io.compression.codec (lz4)
spark.local.dir/tmpWhere 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.maxSizeInFlight48mLarger means fewer, bigger fetch rounds, at the cost of memory
spark.shuffle.io.maxRetries / retryWait3 / 5sRetries 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 questions

How 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 read
How aggregations execute
Partial and final aggregation, hash vs sort aggregation, and why COUNT(DISTINCT) costs more.

Related lessons

Previous: The memory model

Primary sources: RDD guide: shuffle operations · Configuration: shuffle behavior