An executor's memory has two big parts: the JVM heap you set with --executor-memory, and memory outside the heap called overhead. Inside the heap, Spark splits the space into reserved memory, a unified region for execution and storage, and user memory.
The layout, with numbers
Take an executor with --executor-memory 10g.
Container (what YARN or Kubernetes sees) = heap + overhead
Heap 10 GB
Reserved 300 MB
Usable = 10 GB - 300 MB ~9.7 GB
Unified (execution+storage) 60% = ~5.8 GB (spark.memory.fraction = 0.6)
User memory 40% = ~3.9 GB
Overhead = max(384 MB, 10% of heap) = 1 GBWhat each part is for
- Reserved: a fixed 300 MB for Spark internals.
- Execution: shuffles, joins, sorts and aggregations use it for their buffers and hash tables.
- Storage: cached data and broadcast blocks.
- User memory: your own objects, such as data structures built inside UDFs or RDD code, and Spark's metadata.
- Overhead: everything outside the heap. Python worker processes, network buffers, thread stacks and off-heap allocations.
Execution and storage share the unified region and borrow from each other. Execution can evict cached blocks if it needs space, down to a protected share set by spark.memory.storageFraction (0.5 by default). Storage cannot evict running execution data.
Why overhead causes surprising failures
The resource manager sees heap plus overhead. If the total process memory goes above the container limit, YARN or Kubernetes kills the container. You will see "Container killed by YARN for exceeding memory limits" or an OOMKilled pod, not a Java OutOfMemoryError. PySpark jobs with Python UDFs or large pandas conversions hit this most. The usual fix is a larger spark.executor.memoryOverhead, not more heap.
One practical rule
When asked to tune memory, first say which part is full. Heap OOM in execution points to big partitions or skew. A killed container points to overhead. They have different fixes.