Data Engineering › Batch & Distributed Processing
Spill to Disk
Running out of memory and writing intermediate data to disk.
Also known as: disk spill, spilling, external sort, out of core
Spilling to disk is what an engine does when an operation needs more memory than it has: it writes intermediate data to disk and reads it back later. It’s a safety valve that lets a job finish instead of failing with out of memory.
The classic mistake is trusting a job that passed on a sample. A GROUP BY with millions of distinct keys, a sort bigger than the task’s memory, or a hash join that doesn’t fit will spill on full data even though it looked fine in a test. Spilling turns memory-speed work into disk-speed work and can make a job several times slower.
Where it happens
- Sorts that don’t fit in memory become external sorts, writing runs to disk and merging them.
- Hash joins and hash aggregations spill the build side when the table is too large.
- Shuffles (shuffle) write partition data to local disk before it’s fetched.
- Engines differ in when they spill and how they name it, but Spark, Hive/MapReduce, Presto/Trino and most analytical databases all do it.
How to see it
Engine UIs and metrics expose it directly — for example Spark reports “spill (memory)” and “spill (disk)” per task, and query plans show external sorts or hash operators. A single task with far more spill than its peers usually points to data skew: one oversized partition spills while the others don’t.
Reducing it
- Give tasks more memory, or split the work into more, smaller partitions.
- Pre-aggregate or filter earlier so less data reaches the heavy operation.
- Use a broadcast join when one side is small, so the big side isn’t shuffled.
- Reduce the number of distinct keys where you can, and fix skew.
- Rewrite the query if the query plan shows a plan that was never going to fit.
When it’s fine
Some spilling is normal and healthy: large sorts on a shared cluster are supposed to use disk rather than demand huge memory. Only treat it as a problem when it dominates runtime, fills the local disk, or happens for every task. Spilling is a symptom; the fix is usually a smaller working set per task, not just more hardware.