The memory model
Execution vs storage memory, overhead, and what really causes spills and out-of-memory errors.
On this page
You will learn
- How executor memory is divided
- How execution and storage memory borrow from each other
- What memory overhead is, and why containers get killed
- What spills are, and how to diagnose out-of-memory errors
Read first
- Driver, executors and the cluster manager · 4 min read
- Caching and persistence · 3 min read
Comfortable with these? Read on.
spark.memory.fraction) is a unified pool shared by execution (joins, sorts, aggregations) and storage (cache). Off-heap needs, including Python workers, come from memory overhead. Running out of the pool causes spills; running out of the container causes kills.How an executor's memory is laid out
| Region | Size (16 GB heap example) | Used for |
|---|---|---|
| Reserved | 300 MB | Spark internals |
| Unified pool | (16 GB - 300 MB) × 0.6 ≈ 9.4 GB | Execution and storage, shared |
| User memory | The remaining 40% ≈ 6.3 GB | Your objects, UDF data structures, internal metadata |
| Memory overhead (outside the heap) | max(384 MB, 10%) = 1.6 GB | JVM overhead, thread stacks, off-heap buffers, network, Python workers (unless set separately) |
The container the cluster manager allocates is heap + overhead (+ spark.executor.pyspark.memory if set, + off-heap memory if enabled).
Execution and storage share the pool
- Execution memory holds hash tables for joins and aggregations, sort buffers and shuffle buffers. It is needed now, for running tasks.
- Storage memory holds cached blocks and broadcast variables.
- Either can borrow free space from the other. Execution can evict cached blocks to reclaim space, but only down to
spark.memory.storageFraction(50% of the pool), which is protected. Storage can never evict execution memory. - Execution memory is split among the running tasks: with 4 cores, each task gets between 1/8 and 1/4 of the execution memory available.
Spills
When a task's sort or hash table does not fit its share of execution memory, Spark writes part of it to local disk and merges later. That is a spill: slower, but not a failure. The Spark UI shows "Spill (memory)" (deserialized size) and "Spill (disk)" (serialized size) per stage. Occasional spills are normal; large spills on a few tasks mean skew; large spills everywhere mean partitions are too big or memory too small.
Diagnosing out-of-memory errors
| Symptom | Likely cause | Fix |
|---|---|---|
java.lang.OutOfMemoryError on the driver | collect(), toPandas(), large broadcasts, huge plans or file listings | Do not collect; reduce broadcast size; raise spark.driver.memory |
| OOM in an executor task | A huge partition (skew), a huge collect_list group, too few shuffle partitions | More partitions, fix skew, avoid giant arrays |
| Container killed by YARN / Kubernetes for exceeding memory limits | Off-heap use beyond overhead: Python workers, native libraries, many shuffle connections | Raise spark.executor.memoryOverhead or spark.executor.pyspark.memory |
| Long GC pauses, tasks slow then fail | Heap nearly full, often very large executors or caching too much | Smaller executors, less caching |
spark.executor.memory does nothing for them; the container limit comes from overhead or spark.executor.pyspark.memory.Rules of thumb
- Leave
spark.memory.fractionandstorageFractionat their defaults unless you have measured a reason. - More, smaller partitions are usually a better fix for spills than more memory.
- Executors of about 4 to 5 cores and 16 to 32 GB heap avoid very long GC pauses.
Common mistakes
Raising executor memory for container kills
Caching large DataFrames during heavy joins
Treating every spill as a problem
Key takeaways
- The unified pool is 60% of (heap - 300 MB), shared by execution and storage.
- Execution can evict cache down to the protected storage fraction; storage cannot evict execution.
- Spills write to disk when execution memory runs out; they are slow, not fatal.
- Container kills mean off-heap usage: raise memoryOverhead or PySpark memory, not the heap.
Check yourself
3 questions1. With a 10.3 GB heap and defaults, how large is the unified pool?
Show the answer
About 6 GB. (10.3 GB - 0.3 GB) × 0.6 = 6 GB.
2. Can storage (cached data) evict execution memory?
Show the answer
No; execution can evict storage, not the reverse. Running tasks take priority; cache gives way.
3. Python UDF workers run out of memory and YARN kills the container. What do you raise?
Show the answer
spark.executor.memoryOverhead or spark.executor.pyspark.memory. Python workers live outside the JVM heap.
Go deeper
Caching and persistenceSpark internals
Data skew and saltingSpark internals
Driver, executors and the cluster manager
Primary sources: Tuning: memory management · Configuration: memory