Skip to content
Great engineers know 4 min read · spark.memory.fraction

The memory model

Execution vs storage memory, overhead, and what really causes spills and out-of-memory errors.

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

Comfortable with these? Read on.

TL;DR An executor heap keeps 300 MB reserved; 60% of the rest (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

RegionSize (16 GB heap example)Used for
Reserved300 MBSpark internals
Unified pool(16 GB - 300 MB) × 0.6 ≈ 9.4 GBExecution and storage, shared
User memoryThe remaining 40% ≈ 6.3 GBYour objects, UDF data structures, internal metadata
Memory overhead (outside the heap)max(384 MB, 10%) = 1.6 GBJVM 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

SymptomLikely causeFix
java.lang.OutOfMemoryError on the drivercollect(), toPandas(), large broadcasts, huge plans or file listingsDo not collect; reduce broadcast size; raise spark.driver.memory
OOM in an executor taskA huge partition (skew), a huge collect_list group, too few shuffle partitionsMore partitions, fix skew, avoid giant arrays
Container killed by YARN / Kubernetes for exceeding memory limitsOff-heap use beyond overhead: Python workers, native libraries, many shuffle connectionsRaise spark.executor.memoryOverhead or spark.executor.pyspark.memory
Long GC pauses, tasks slow then failHeap nearly full, often very large executors or caching too muchSmaller executors, less caching
Why this matters: PySpark jobs with pandas UDFs or Python UDFs use memory in Python worker processes, which live outside the JVM heap. Raising spark.executor.memory does nothing for them; the container limit comes from overhead or spark.executor.pyspark.memory.

Rules of thumb

  • Leave spark.memory.fraction and storageFraction at 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

The limit being hit is overhead, not heap.

Caching large DataFrames during heavy joins

Storage pressure leaves less room for execution, causing spills.

Treating every spill as a problem

Small spills are normal; look at their distribution.

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 questions

1. 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

Primary sources: Tuning: memory management · Configuration: memory