skip to content

Why does YARN kill a Spark executor for exceeding memory limits when its JVM heap never filled up?

level: seniorimportance: should knowfreq 58%

answer

  1. the limit is not the heap
  2. the process is bigger than the JVM
  3. shuffle buffers and interpreters live outside
  4. overhead defaults to ten percent

basics

~20 s

The container limit covers the whole process, not just the heap: spark.executor.memory plus spark.executor.memoryOverhead. Off-heap consumers — shuffle network buffers, JVM metaspace and thread stacks, native libraries and Python workers — burst past the overhead allowance while the heap sits half empty.

solid answer

~40 s

YARN enforces a limit on the **container**, which Spark requests as `spark.executor.memory` + `spark.executor.memoryOverhead` (plus `spark.memory.offHeap.size` and `spark.executor.pyspark.memory` when set). The heap is only the first term. Everything else the process allocates — Netty's direct byte buffers for shuffle transfer, metaspace and code cache, thread stacks, native compression codecs, and in PySpark the separate Python worker processes — comes out of the overhead. `spark.executor.memoryOverhead` defaults to `max(384 MB, 0.1 x spark.executor.memory)`, which is thin for shuffle-heavy or Python jobs. So the JVM never OOMs; the operating system's cgroup accounting simply exceeds the container ceiling and YARN kills it. On Kubernetes the same cause appears as an `OOMKilled` pod with exit code 137. The fix is to raise the overhead, budget Python memory explicitly, and reduce concurrent tasks per executor.

code

bash · 7 lines
bash
spark-submit \
  --conf spark.executor.memory=8g \
  --conf spark.executor.memoryOverhead=3g \
  --conf spark.executor.pyspark.memory=2g \
  --conf spark.executor.cores=4 \
  job.py
# container request = 8g + 3g + 2g

go deeper

for a junior

Know that an executor's container is bigger than its Java heap, and that Spark asks the cluster manager for executor memory plus an overhead allowance on top.

for a middle

Name spark.executor.memoryOverhead and its default of the larger of 384 MB and ten percent of executor memory, and list what actually consumes memory outside the heap.

for a senior

Diagnose from evidence: correlate kills with the stage type, confirm the heap was not the constraint, and pick between raising overhead, budgeting PySpark memory, and lowering cores per executor rather than reflexively enlarging the heap.

for a principal

Own executor shapes as a platform decision — container size, cores per executor and overhead defaults interact with queue capacity, so a per-team overhead free-for-all quietly costs cluster throughput.

## The error and what it is really saying The YARN NodeManager message names *physical memory used* against the container's limit and suggests boosting the executor memory overhead. The important word is **container**. YARN measures the whole process tree, not the JVM heap. A heap that peaked at 4 GB out of 8 GB is entirely consistent with a container kill. ## What Spark asks for When Spark requests an executor container, the size is: ``` spark.executor.memory # the JVM -Xmx + spark.executor.memoryOverhead # default max(384 MB, 0.1 x executor memory) + spark.memory.offHeap.size # when off-heap is enabled + spark.executor.pyspark.memory # when set, for PySpark ``` Since Spark 3.3 the 0.1 multiplier itself is exposed as `spark.executor.memoryOverheadFactor`. Setting `spark.executor.memoryOverhead` explicitly overrides the computed value. On pre-2.3 clusters you will still see the old `spark.yarn.executor.memoryOverhead` name in scripts; it was superseded by the cluster-manager-neutral key. ## What lives outside the heap - **Netty direct buffers.** Shuffle blocks are transferred with off-heap direct byte buffers. A stage that fetches thousands of blocks concurrently across many cores can hold hundreds of megabytes here. - **JVM metaspace, code cache, thread stacks, GC structures.** Fixed-ish, but non-trivial: hundreds of threads at the default stack size adds up. - **Native libraries.** Compression codecs (Zstd, LZ4, Snappy), Arrow buffers, and native BLAS in ML workloads allocate outside the heap. - **Python workers.** This is the biggest one in practice. In PySpark, each concurrent task on an executor forks a Python process whose pandas frames, arrow batches and interpreter footprint are all outside the JVM. An executor with 8 cores can be running 8 Python interpreters at once. - **Off-heap Tungsten memory**, if `spark.memory.offHeap.enabled` is on. ## How to diagnose it 1. **Read which limit was breached.** Physical-memory kills point at real RSS; virtual-memory kills on older clusters were often an artifact of the JVM's address-space reservation and were usually addressed at the NodeManager's virtual-memory check rather than in Spark. 2. **Correlate with the stage.** If kills cluster on shuffle-read stages, suspect network buffers and the fetch concurrency. If they cluster on Python UDF stages, suspect the Python workers. 3. **Check the heap trend separately.** If the executor's heap is comfortably below `-Xmx` and GC time is low, the pressure is definitively outside the heap and raising `spark.executor.memory` will not help — it can even *reduce* headroom if the container size is fixed by a queue policy. 4. **Check task concurrency.** Off-heap consumption scales with concurrent tasks, so `spark.executor.cores` is a direct multiplier. ## The fixes, in the order I would try them - **Raise `spark.executor.memoryOverhead`.** For shuffle-heavy JVM jobs, going well above the 10% default is routine; for PySpark, larger still. - **Budget Python explicitly** with `spark.executor.pyspark.memory` so the request accounts for interpreters rather than hoping the overhead absorbs them, and cut the size of the pandas batches a UDF materializes. - **Lower `spark.executor.cores`.** Fewer concurrent tasks means fewer simultaneous Python workers and fewer in-flight fetch buffers. Often the cheapest real fix. - **Shrink partitions.** Smaller partitions mean smaller per-task native buffers and smaller pandas batches. More partitions is again the general-purpose lever. - **Only then raise the heap**, and remember it raises the container request too, which a constrained YARN queue may refuse to schedule. ## The same failure on other cluster managers On Kubernetes there is no YARN message: the pod is `OOMKilled` by the kernel and the executor exits with code **137**, because the pod's memory limit is computed from the same sum. The diagnosis is identical — the process tree exceeded a limit the heap alone does not explain. On standalone Spark with no cgroup enforcement, the same over-allocation instead shows up as host-level swapping or the Linux OOM killer taking a process, which is harder to attribute. ## The misdiagnosis to avoid The instinct is to increase `spark.executor.memory`. That grows the heap *and* the container, so the container ceiling moves and the symptom may disappear for the wrong reason — you paid for gigabytes of heap to buy a few hundred megabytes of overhead. Worse, on a cluster with fixed executor counts it reduces the number of executors your queue can hold. Diagnose which side of the boundary the memory is on before spending.

  • A PySpark job is killed for exceeding its container limit while the JVM heap is half free. What do you change?
    Budget the Python side explicitly with spark.executor.pyspark.memory and raise spark.executor.memoryOverhead, then reduce spark.executor.cores so fewer Python interpreters run concurrently on one executor. Also shrink the batch a pandas UDF materializes and increase the partition count, since a Python worker's footprint scales with the rows it holds at once. Raising spark.executor.memory alone buys heap you do not need.
  • How does this failure look on Kubernetes rather than YARN?
    The pod is OOMKilled by the kernel and the executor exits with code 137; there is no NodeManager message naming the overhead setting. The cause is the same, because the pod's memory limit is computed from spark.executor.memory plus spark.executor.memoryOverhead plus off-heap and PySpark memory. You diagnose it the same way: compare heap usage against the pod limit and look for off-heap consumers.
  • Why can raising spark.executor.memory make the situation worse rather than better?
    It grows the heap and the container request together, so you pay for gigabytes of heap to gain a fraction of that in overhead headroom. On a queue with a fixed memory budget it also reduces how many executors you can run, cutting parallelism. If the heap was never the constraint, the correct lever is the overhead setting, the PySpark budget, or fewer cores per executor.

saying these in an interview costs you the question

  • Says the JVM must have thrown OutOfMemoryError
  • Raises spark.executor.memory as the first response
  • Ignores Python worker processes in PySpark diagnosis
  • Thinks memoryOverhead is part of the heap
  • Assumes the default 10% overhead is adequate for shuffle-heavy jobs

context