skip to content

In Spark, what files does a map task write during a sort-based shuffle?

level: middleimportance: must knowfreq 68%

answer

  1. count the files, not the partitions
  2. one task, two artifacts on local disk
  3. the reader has to know where to seek
  4. a data file plus an offset table

basics

~20 s

Each Spark map task writes two local files: one data file whose records are grouped and ordered by destination reduce partition, plus a small index file giving each partition's byte offset inside that data file.

solid answer

~40 s

Spark's sort-based shuffle writes **one data file and one index file per map task**, not one file per reduce partition. The writer buffers records in execution memory, sorts them by target partition id (and by key when map-side aggregation is needed), spills sorted runs to local disk under `spark.local.dir` when memory runs short, then merges those runs into the single data file. The `.index` file records where each reduce partition's bytes begin, so a fetch becomes a seek plus a ranged read. Three writer paths exist: `BypassMergeSortShuffleWriter` (no map-side combine and partition count at or below `spark.shuffle.sort.bypassMergeThreshold`, default 200), `UnsafeShuffleWriter` (the serialized Tungsten sort), and the general `SortShuffleWriter`. Output is compressed by default via `spark.shuffle.compress` using `spark.io.compression.codec`.

code

text · 3 lines
text
# one map task's output on the executor's local disk
/data/1/spark-local/blockmgr-6f2a.../0f/shuffle_2_137_0.data
/data/1/spark-local/blockmgr-6f2a.../3a/shuffle_2_137_0.index

go deeper

for a junior

Recall that a wide transformation makes each upstream task write its output to local disk first, and that the next stage reads those files. Knowing the handoff is file-based, not a live network stream, is enough here.

for a middle

Be ready to describe the data-file-plus-index layout, the sort-by-partition-then-spill-then-merge sequence, and why one file per task beats one file per partition. Name spark.shuffle.compress and the bypass threshold without guessing values you are unsure of.

for a senior

Connect the write path to what you see in production: disk pressure under spark.local.dir, spill metrics, and why a dead executor makes its shuffle files unreachable. Explain which writer path a given operation takes and what that costs.

for a principal

Own the platform consequence: shuffle write volume drives local-disk provisioning, node instance choice and the codec tradeoff across many tenants. Be ready to argue when to standardise on zstd or to move shuffle storage off the executor entirely.

## What a shuffle has to accomplish A wide transformation — a join, a `groupBy`, a `repartition` — needs every row to land on the executor responsible for its target partition. Spark implements that as a two-sided file exchange between two stages. The upstream stage's tasks (map tasks, in shuffle vocabulary) each write their entire output to **local disk** on the machine where they ran. The downstream stage's tasks (reduce tasks) then pull the byte ranges addressed to them. Nothing flows through the driver, and nothing is written to HDFS or object storage: shuffle data is local, temporary, and recomputable from lineage if lost. ## Two files per map task Since sort-based shuffle became the only shuffle manager (Spark 2.0 removed the old hash shuffle), each map task leaves exactly two files behind, regardless of how many reduce partitions the shuffle has: - `shuffle_<shuffleId>_<mapId>_0.data` — all of that task's output records, laid out contiguously, partition 0's bytes first, then partition 1's, and so on. - `shuffle_<shuffleId>_<mapId>_0.index` — a tiny file of offsets, one per reduce partition, marking where each partition's region starts. That pairing is the whole design. A reduce task that wants partition 57 from map task 137 does not need a file per partition; it needs the index entry for 57 and a ranged read of the data file. The historical alternative, hash shuffle, wrote one file per (map task, reduce partition) pair — M × R files — which exhausted file descriptors and turned the disk into a random-I/O workload on any large job. ## Inside the writer: sort, spill, merge The general path uses an `ExternalSorter`. Records are inserted into an in-memory buffer sized from the executor's execution memory pool. The sorter sorts by partition id, and additionally by key when the operation requires **map-side combine** (an aggregation like `reduceByKey`, which pre-aggregates on the map side so fewer bytes cross the network). When the sorter cannot acquire more execution memory, it sorts what it has, writes that run to disk as a spill file, and starts over. At the end, all spill files plus the in-memory remainder are merged in partition order into the one data file. Spill is normal on large shuffles, not a failure — it just costs extra write-then-read I/O. ## Three writer implementations Spark picks the writer per shuffle: - **Bypass merge sort** — chosen when there is no map-side combine and the reduce-partition count is at or below `spark.shuffle.sort.bypassMergeThreshold` (default 200). It opens one temporary file per reduce partition, appends records without sorting, then concatenates the temp files into the standard data-plus-index pair. Cheap when the partition count is small. - **Serialized (Tungsten) sort** — `UnsafeShuffleWriter`. Records are serialized once on insert, and the sorter sorts compact 8-byte pointers carrying the partition id rather than moving objects. It requires a serializer that supports relocation of serialized objects, no map-side combine, and a partition count under Spark's serialized-mode ceiling. DataFrame shuffles, which move `UnsafeRow` binary, land here often. - **General sort** — everything else, including all map-side-combine shuffles. ## Compression, buffering, serialization Shuffle output is compressed by default: `spark.shuffle.compress` is on and uses `spark.io.compression.codec`, whose default is lz4; snappy and zstd are the common alternatives, with zstd trading CPU for a smaller footprint. Spill files have their own switch, `spark.shuffle.spill.compress`. Each open output stream buffers writes through `spark.shuffle.file.buffer` (default 32k) so the writer is not issuing tiny syscalls. Everything crossing a shuffle is serialized — that is why an unserializable object in a closure fails at a shuffle boundary and not before. ## Lifetime, and why it matters operationally Shuffle files live in the executor's block-manager directories under `spark.local.dir` and survive the stage that produced them: a later stage, a retry, or a re-run of a cached-and-lost partition may need them again. They are cleaned up when the shuffle is unregistered (Spark's `ContextCleaner`, once the owning RDD or query result is garbage collected) or when the application exits. Two consequences follow directly: shuffle-heavy jobs need real local disk headroom, and if the executor process dies, its files become unreachable unless an external shuffle service on the node serves them — which is exactly the situation that produces a fetch failure and a recomputed upstream stage. ## What interviewers are checking That you know the exchange is a disk-mediated file handoff, not a network stream between live tasks; that the file count scales with map tasks rather than with the product of map and reduce counts; and that you can connect "my job spilled" or "my disk filled up" to a concrete write path instead of treating the shuffle as a black box.

  • Why does Spark write one data file per map task instead of one file per reduce partition?
    The old hash shuffle did exactly that and produced map-tasks × reduce-partitions files, exhausting file descriptors and turning writes into random I/O. Sorting by partition id lets one file serve every reduce partition through an index. The bypass path still opens per-partition temp files, but only when the partition count is at or below `spark.shuffle.sort.bypassMergeThreshold` (default 200), and it concatenates them into the same two-file layout.
  • Is shuffle output compressed on disk, and can you change how?
    Yes. `spark.shuffle.compress` is enabled by default and uses `spark.io.compression.codec`, which defaults to lz4; snappy and zstd are the usual alternatives. Spill files are governed separately by `spark.shuffle.spill.compress`. Switching to zstd shrinks bytes written and fetched at the cost of CPU, which pays off on network- or disk-bound shuffles and hurts on CPU-bound ones.
  • What lets Spark use the serialized Tungsten shuffle writer?
    `UnsafeShuffleWriter` is eligible when the serializer supports relocation of serialized objects, the shuffle needs no map-side combine, and the reduce-partition count is under Spark's serialized-mode ceiling. It serializes each record once and sorts packed pointers that embed the partition id, so sorting never touches or copies the record bodies — much less CPU and GC pressure than sorting deserialized objects.

saying these in an interview costs you the question

  • Says each map task writes one file per reduce partition
  • Thinks shuffle data is routed through the driver
  • Claims shuffle output goes to HDFS or object storage
  • Believes shuffle blocks stay in memory until fetched
  • Assumes shuffle files are deleted the moment the stage ends

context