skip to content

What makes a Spark shuffle task spill records to disk?

level: seniorimportance: should knowfreq 55%

answer

  1. it is a controlled retreat, not a crash
  2. the budget is smaller than you think
  3. count the tasks sharing one executor's heap
  4. sorted runs written out, merged later

basics

~20 s

A shuffle task spills when its sorter or aggregation map cannot acquire more execution memory for the records it is holding. Spark sorts what it has, writes it to local disk as a run, and merges the runs later.

solid answer

~50 s

Spill is what Spark does instead of failing when a task's in-memory structure outgrows its share of execution memory. On the write side, the `ExternalSorter` or `ShuffleExternalSorter` buffers records; when it cannot grow, it sorts the buffer, writes a sorted run to local disk, and continues. On the read side, aggregations and sort-merge joins use spilling structures the same way. The trigger is **per-task memory**, not cluster memory: the executor's execution pool is shared by all tasks running concurrently there, so more cores per executor means a smaller slice each. The stage UI shows *Spill (Memory)* — the estimated in-memory footprint of what was spilled — and *Spill (Disk)*, the serialized, compressed bytes actually written, which is why the disk number is smaller. Common cures: more shuffle partitions, fewer cores per executor, and pre-aggregation instead of `groupByKey`.

code

text · 5 lines
text
Stage 12 (sort-merge join)  tasks 800/800
Shuffle Read:        612.0 GB
Spill (Memory):      1,204.7 GB
Spill (Disk):        188.3 GB
spark.executor.cores = 8   spark.sql.shuffle.partitions = 200

go deeper

for a junior

Know that when data does not fit in memory Spark writes it to local disk and continues rather than failing, and that this makes the stage slower.

for a middle

Explain the sorter's buffer-sort-spill-merge cycle and that the memory limit applies per task inside one executor's execution pool, not cluster-wide. Distinguish the two spill metrics.

for a senior

Diagnose a spilling stage from its metrics and prescribe the right lever — partition count, cores per executor, pre-aggregation, projection — and justify why adding executors alone would not help.

for a principal

Weigh spill against cost: memory-rich executors that never spill versus cheaper ones that spill modestly, and the disk provisioning that implies. Set the default sizing conventions others inherit.

## Spill is a safety valve, not an error A Spark task processing a shuffle cannot assume its data fits in memory. Rather than failing, the memory-hungry structures — sorters, aggregation maps, join buffers — are *spillable*: when they cannot get more memory, they serialize their current contents to local disk in sorted order and start fresh. Spill costs a write and a re-read, plus a merge, so a spilling stage is slower, but it completes. A stage that spills a little is usually fine; a stage spilling tens of gigabytes is telling you a sizing decision is wrong. ## Where spill happens in a shuffle **Write side.** The shuffle writer inserts records into an `ExternalSorter` (or `ShuffleExternalSorter` on the serialized path), which sorts by partition id — and by key when the operation does map-side combine. Whenever it fails to acquire additional execution memory, it sorts the buffer, writes that run to a spill file under `spark.local.dir`, and keeps going. At task end, all runs merge into the map task's single data file. `spark.shuffle.spill.compress` (on by default) compresses those runs. **Read side.** Reduce tasks spill too. An aggregation whose per-key state does not fit uses a spilling append-only map; a sort-merge join sorts each side with a spilling sorter. This is why you can see spill in a stage that performs no shuffle write at all. ## Why free cluster memory does not prevent it The memory that matters is a **task's share of one executor's execution pool**. Spark's unified memory manager splits the executor heap into an execution region (shuffles, sorts, joins, aggregations) and a storage region (cached blocks); execution can borrow back from storage by evicting cached blocks, but it can never reach across to another executor. Within an executor, the pool is divided among the tasks running concurrently, which is `spark.executor.cores` of them. So an executor with a large heap and 8 cores gives each task roughly an eighth of the pool, and a job can spill heavily while the cluster-wide memory chart looks half empty. The corollary surprises people: *reducing* cores per executor often removes spill, because each remaining task gets a bigger slice. ## Reading the metrics The stage detail page reports two numbers per task, and they are not measuring the same thing: - **Spill (Memory)** — the estimated size of the records *as held in memory* before spilling: deserialized objects, with JVM overhead. - **Spill (Disk)** — the bytes actually written after serialization and compression. Disk is normally several times smaller than memory, and that ratio is a rough gauge of how expensive your in-memory representation is. Compare spill against shuffle read size for the same stage: spilling twice the data you read means it is being written and re-read repeatedly through multiple merge passes. ## Fixing it, in order of leverage 1. **More shuffle partitions.** Halving the bytes per partition halves what each task must hold. This is the first lever because it is cheap and reversible, and it is why an under-set partition count shows up as spill first. 2. **Fewer cores per executor** (or more memory per executor). Directly enlarges each task's slice of the execution pool. 3. **Pre-aggregate.** `reduceByKey` and `aggregateByKey` combine on the map side; `groupByKey` moves every raw record and materialises whole groups on the reader, which is the textbook way to force enormous spill. In DataFrame land the planner already does partial aggregation, but a wide `collect_list` reproduces the same problem. 4. **Shrink the rows.** Project away unused columns before the exchange; fewer bytes cross and fewer bytes must be held. 5. **Address skew** if only a handful of tasks spill while the rest do not — that is a distribution problem, not a sizing problem. Raising `spark.memory.fraction` is available but low-leverage and risks starving cached data; treat it as a last resort rather than a first answer. ## What interviewers are checking The distinguishing insight is *per-task* memory. A candidate who says "add executors" has not understood that spill is decided inside one task's memory budget, and more executors change neither the bytes per partition nor the slice per task unless the partition count moves with them. The second insight is that spill is a graceful degradation — the alternative would be an OOM, and tuning that eliminates all spill at the cost of huge executors is often the worse trade.

  • Why is Spill (Disk) usually much smaller than Spill (Memory) for the same task?
    They measure different representations. Spill (Memory) estimates the footprint of the records as live JVM objects, including per-object overhead; Spill (Disk) counts the bytes after serialization and compression by `spark.shuffle.spill.compress`. A large gap indicates an expensive in-memory representation, which is one reason the DataFrame API's compact binary rows spill less than equivalent RDD object pipelines.
  • Why can lowering spark.executor.cores remove spill without adding memory?
    An executor's execution memory pool is shared by the tasks running concurrently in it, one per core. Cutting cores from eight to four roughly doubles each remaining task's share, so structures that previously could not grow now fit. You trade some parallelism for fewer spill round trips, which frequently makes the stage faster overall.
  • Is spill always worth eliminating?
    No. Spill is graceful degradation; the alternative when memory truly runs out is an OOM and a lost executor. Sizing every executor so that nothing ever spills usually means over-provisioning memory for a worst-case partition. Aim to remove pathological spill — tens of gigabytes, repeated merge passes, or spill concentrated in a few skewed tasks — and accept modest spill.

saying these in an interview costs you the question

  • Says spill means the cluster is out of memory
  • Recommends adding executors without changing partition count
  • Treats any spill as a bug to be eliminated
  • Confuses Spill (Memory) with Spill (Disk) bytes
  • Thinks caching the DataFrame prevents shuffle spill

context