In the Spark UI, one task reports Spill (Memory) of 48 GB on an executor with an 8 GB heap — what is that measuring?
answer
- not a peak, a running total
- one number counts bytes, the other counts objects
- every flush adds to the tally
- many small spills, one big merge
basics
~20 sSpill (Memory) is the cumulative in-memory, deserialized size of every buffer the task flushed to disk, summed over all spill events. It is not a peak, so a task that spills repeatedly can report far more than the executor heap ever held.
solid answer
~50 sSpark reports two spill numbers per task. **Spill (Memory)** is the deserialized, in-memory size of the data at each moment it was spilled, **summed across every spill event**. **Spill (Disk)** is the serialized, compressed bytes actually written. Because the memory figure is cumulative, a task that fills and flushes a 2 GB sort buffer twenty-four times reports 48 GB on an 8 GB heap — nothing overflowed the JVM. The ratio between the two numbers is just your serialization and compression factor, so a 6:1 gap is ordinary. Spilling is Spark working correctly under memory pressure, not an error: `ExternalSorter` and `ExternalAppendOnlyMap` flush to local disk when a task cannot acquire more execution memory. It is a cost signal — sort-merge and disk IO instead of an in-memory pass — and you attack it with more partitions, fewer concurrent tasks per executor, or less cached data crowding the pool.
code
text · 6 linesSummary Metrics for 200 Completed Tasks
Metric Min 25th Median 75th Max
Duration 12 s 14 s 15 s 17 s 4.2 min
Shuffle Read Size 180 MB 190 MB 195 MB 210 MB 9.8 GB
Spill (Memory) 0.0 B 0.0 B 0.0 B 1.1 GB 48.0 GB
Spill (Disk) 0.0 B 0.0 B 0.0 B 180 MB 7.9 GBgo deeper
Know where the Spill (Memory) and Spill (Disk) columns appear in the Spark UI's stage page and that spilling means data was written to local disk rather than kept in memory.
Explain that both numbers are cumulative across spill events, why the memory figure can exceed the executor heap, and which operators spill when they cannot acquire execution memory.
Diagnose from the task-level distribution: uniform spill points at partition sizing or cores per executor, while a couple of huge spilling tasks point at skew, and be able to say which lever you would pull first and why.
Frame spill as a cost signal inside a cluster budget rather than a defect to eliminate, and set platform guidance on partition sizing and cores per executor so teams stop trading memory for parallelism blindly.
## Two numbers, two units The stage detail page in the Spark UI shows, per task and aggregated per stage: - **Spill (Memory)** — the size the spilled records occupied *in memory*, deserialized, at the instant of each spill, summed over all spills by that task. - **Spill (Disk)** — the bytes actually written to the executor's local disk, after serialization and compression. They measure the same records through different lenses, which is why the memory number is routinely several times the disk number. That ratio is your serialization plus compression factor and is not, by itself, a problem. ## Why the memory number can exceed the heap This is the part that trips people up. Neither number is a peak. Each is a running total. A task whose sort buffer grows to 2 GB, spills, grows again, spills again, and repeats twenty-four times has never held more than 2 GB at once — but reports 48 GB of memory spill. Seeing a spill figure larger than `spark.executor.memory` tells you the task spilled *many times*, which is itself the useful signal: many small spill files mean an expensive merge phase afterwards. ## Who spills, and when Execution memory is requested incrementally by spillable operators. When a task asks the executor's execution memory pool for more space and cannot get it, the operator flushes its current contents to local disk and starts fresh. The main spillers are: - `ExternalSorter` — sort-based shuffle writes and sorts. - `ExternalAppendOnlyMap` — map-side combining and aggregations. - The Tungsten sort/aggregate operators, which spill their binary pages the same way. - Sort-merge join, which sorts both sides before merging. All of them are designed to spill. Correctness is never at risk; the cost is extra IO and a merge pass over the spill files. ## Why a task cannot use the whole pool The executor's execution memory is shared dynamically among its concurrently running tasks. A task is guaranteed at least `1/(2N)` of the pool and may acquire at most `1/N`, where N is the number of active tasks — that is, `spark.executor.cores` when the executor is saturated. So an 8-core executor with a 6 GB unified pool gives one task at most roughly 750 MB of execution memory, no matter how much heap the machine has. Doubling executor cores to get more parallelism halves each task's memory budget, and spill is usually the first symptom. The second squeeze is caching. Cached blocks occupy the storage side of the same unified pool. Once storage holds its protected floor (`spark.memory.storageFraction`), execution cannot evict any further, so a heavily cached job leaves execution a smaller pool to divide. ## Reading the evidence Open the stage's task table and look at the distribution, not the sum: - **Every task spills a similar amount** → the stage's per-partition data volume simply exceeds the per-task memory budget. Raise the partition count so each partition is smaller, or reduce concurrency per executor. - **One or two tasks spill enormously and run long while the rest finish fast** → skew. The fix is in the data distribution, not the memory settings. - **Spill appears only after a `cache()` was added upstream** → storage is crowding execution. Cache less, use a disk-backed or serialized level, or lower `spark.memory.storageFraction` so execution may evict more. ## What to change, in order 1. **More partitions.** For SQL shuffles this is `spark.sql.shuffle.partitions`; for RDD code it is the `numPartitions` argument or an explicit repartition. Smaller partitions mean smaller per-task working sets, and this is the lever that works most often. 2. **Fewer cores per executor.** Fewer concurrent tasks means a larger `1/N` slice each. Throughput per executor may drop slightly while total runtime improves. 3. **Less cached data**, or a cheaper storage level, to free the pool. 4. **A bigger heap**, last — it is the expensive lever and does nothing about skew. ## What spill is not Spill is not an out-of-memory error, and "zero spill" is not a tuning goal. A job that spills a modest amount and finishes on time needs no attention. Spill is also not the same as *shuffle write*: shuffle write bytes are the normal output of a map stage, written to disk by design for every shuffle. Confusing the two leads people to "fix" healthy shuffles. Finally, spilling in Spark is unrelated to a garbage-collection problem — high GC time is a separate metric in the same task table, and the two have different remedies.
- Why is a task's Spill (Disk) number usually much smaller than its Spill (Memory) number?Because they measure the same records in different forms. Spill (Memory) is the deserialized in-memory footprint, inflated by JVM object headers and pointers; Spill (Disk) is what was written after serialization and compression. A ratio of several times is completely ordinary and tells you about your serializer and codec, not about a problem.
- An executor has 8 cores and a 6 GB unified pool. Roughly how much execution memory can one task hold?At most about one eighth of the pool, so roughly 750 MB, and it is guaranteed no less than about one sixteenth. Execution memory is divided dynamically among active tasks, so raising executor cores to increase parallelism directly shrinks each task's budget. That is why adding cores to a spilling stage often makes the spilling worse rather than better.
- Should the goal of tuning be to eliminate spill entirely?No. Spilling is Spark's designed response to memory pressure and keeps jobs correct rather than failing them. Modest spill on a stage that meets its runtime target needs no action. Chase it only when the spill volume clearly dominates the stage time, when a handful of tasks spill far more than the rest, or when the merge of many small spill files is what is slow.
It is a water meter, not a tank gauge: it records every litre that ever passed through, so a small bucket emptied many times can register a huge total.
saying these in an interview costs you the question
- Reads Spill (Memory) as a peak rather than a cumulative total
- Says spill means the executor ran out of heap and will crash
- Confuses spill bytes with normal shuffle write bytes
- Adds executor cores to fix a spilling stage
- Treats zero spill as the tuning target