Why does a YARN NodeManager kill a container for exceeding physical memory limits?
answer
- the grant is for the whole process tree
- heap is not the whole process
- read which limit the message names
- more cores means more tasks sharing memory
- one container dying is a data problem
basics
~20 sEach YARN container is granted a fixed memory amount, and the NodeManager monitors the container's whole process tree against it. When total resident memory crosses the grant the container is killed, because YARN protects the node from oversubscription rather than letting the machine swap or OOM.
solid answer
~50 sA container is a lease on a specific amount of memory, and the NodeManager's container monitor samples each container's *entire process tree* — JVM heap plus everything outside it — against that grant. Cross it and the container is killed with the familiar "running beyond physical memory limits" diagnostic. The check is governed by `yarn.nodemanager.pmem-check-enabled` (true by default), with a parallel virtual-memory check via `yarn.nodemanager.vmem-check-enabled` and `yarn.nodemanager.vmem-pmem-ratio` (2.1). The trap is that the JVM heap is only part of the total: thread stacks, metaspace, direct/off-heap buffers, native libraries and Python worker processes all count. For Spark on YARN the request is `spark.executor.memory` plus `spark.executor.memoryOverhead` — 10% of executor memory with a 384 MiB floor by default — so the fix is usually to raise the overhead or lower per-container concurrency, not to disable the check.
code
text · 3 linesContainer [pid=24817,containerID=container_1699034221_0007_01_000042] is
running beyond physical memory limits. Current usage: 8.4 GB of 8 GB physical
memory used; 12.1 GB of 16.8 GB virtual memory used. Killing container.go deeper
Know that a YARN container has a fixed memory grant and that exceeding it gets the container killed rather than slowed down, and that the grant covers more than the JVM heap.
Explain the enforcement mechanism — the NodeManager sampling the whole process tree — and account for what lives outside the heap: metaspace, thread stacks, off-heap buffers and child processes.
Diagnose a real occurrence: read the diagnostic to see which limit fired, distinguish steady undersizing from a stage-specific spike or skew, and pick the fix in the right order rather than reaching for more memory.
Own the platform policy: what container sizes the cluster is standardised on, whether the virtual-memory check stays enabled, and how container shape affects packing efficiency and utilisation across the fleet.
## What the container was promised When an ApplicationMaster requests a container it names a memory size in MB. The ResourceManager rounds that request **up** to a multiple of `yarn.scheduler.minimum-allocation-mb` (default 1024) and refuses anything above `yarn.scheduler.maximum-allocation-mb` (default 8192, and routinely raised on real clusters). The NodeManager then admits the container against what it advertises in `yarn.nodemanager.resource.memory-mb`. That advertised number is a fiction the whole cluster depends on: YARN packs containers onto a node assuming none will exceed its grant, so a container that does would push the machine into swapping or into a kernel OOM kill that could take down unrelated containers or the NodeManager itself. ## How enforcement works The NodeManager runs a containers-monitor thread that periodically walks each container's **process tree** and sums its resident set size. Two checks apply: - **Physical memory**, `yarn.nodemanager.pmem-check-enabled`, true by default. Exceeding the granted MB kills the container. - **Virtual memory**, `yarn.nodemanager.vmem-check-enabled`, true by default, with the allowance computed as granted memory times `yarn.nodemanager.vmem-pmem-ratio` (default **2.1**). The diagnostic looks like: `Container [pid=...] is running beyond physical memory limits. Current usage: 8.4 GB of 8 GB physical memory used; 12.1 GB of 16.8 GB virtual memory used. Killing container.` Read *which* limit the sentence names — physical and virtual overruns have different causes and different fixes. Where the LinuxContainerExecutor is configured with cgroups, CPU is enforced by the kernel as a genuine limit; memory enforcement in the common configuration remains this monitor-and-kill behaviour, which is why you see a killed container rather than a throttled one. ## Why the JVM heap is not the number that matters This is the heart of the question. A JVM launched with `-Xmx6g` inside a 6 GB container will be killed, because the process consumes far more than its heap: - **Metaspace and code cache** — hundreds of MB in a large application. - **Thread stacks** — roughly 1 MB per thread; a container running many concurrent tasks accumulates real memory here. - **Direct and off-heap buffers** — network buffers, memory-mapped shuffle files, off-heap execution memory. - **Native libraries** — compression codecs and native math libraries allocate outside the heap. - **Child processes** — PySpark launches Python workers whose memory belongs to the same process tree and counts against the same container. Spark models this explicitly: a YARN container for an executor is sized as `spark.executor.memory` (the heap) plus `spark.executor.memoryOverhead`, which defaults to 10% of executor memory with a **384 MiB** floor and is regularly too small for shuffle-heavy or Python workloads. PySpark additionally has `spark.executor.pyspark.memory`. The corresponding driver settings apply when the driver runs inside the ApplicationMaster container in cluster deploy mode. ## Diagnosing a real occurrence First, read the message and note the ratio between usage and limit — a container 5% over needs more overhead, one at 2x needs a design change. Second, decide whether the overrun is *steady* or *spiky*: steady overrun means the container is simply undersized; spikes at a particular stage point at a specific operation, usually a shuffle fetch, a large broadcast, a wide aggregation, or collecting too much data into the coordinator. Third, look at concurrency inside the container — Spark's `spark.executor.cores` multiplies the number of tasks sharing one container's memory, so halving cores per executor is often a more effective fix than adding memory. Fourth, check for skew: if exactly one container dies repeatedly while hundreds succeed, the memory setting is not the problem, the data distribution is. ## Fixes, in the order to try them 1. Raise the framework's overhead allowance so the container request matches real consumption. 2. Reduce per-container parallelism so fewer tasks share the same memory. 3. Address the workload — reduce partition sizes, avoid pulling large results into the coordinator, fix skew. 4. Only then raise total container memory, keeping in mind that fewer, larger containers pack worse onto nodes and are refused outright once they exceed `yarn.scheduler.maximum-allocation-mb`. ## The one thing not to do Disabling `yarn.nodemanager.pmem-check-enabled` is the tempting one-line fix and is almost always wrong: you have not made the container use less memory, you have removed the safety net that keeps one greedy container from destabilising the whole node. The defensible exception is the **virtual** memory check. Virtual memory overruns often reflect address-space reservations rather than real usage — glibc per-thread arenas are the classic cause — so many clusters legitimately turn `yarn.nodemanager.vmem-check-enabled` off while leaving the physical check on. Know the difference between those two decisions.
- Why is disabling the physical memory check the wrong first response?Because it removes the protection, not the consumption. YARN packs containers onto a node on the promise that none exceeds its grant; without the check an over-consuming container drives the machine into swap or a kernel OOM kill, which can take out unrelated containers or the NodeManager itself. The virtual-memory check is a different matter — it often fires on address-space reservations rather than real usage, and disabling that one alone is a defensible, common choice.
- One container out of four hundred is killed for memory, repeatedly, at the same stage. What do you suspect?Data skew, not undersized memory. A uniform shortage kills many containers; a single repeat offender means one partition holds far more rows than the rest, so raising memory just moves the threshold and wastes it on the other 399. Investigate the key distribution feeding that stage and fix the partitioning before touching container sizes.
- How does the number of cores per container relate to memory kills?Cores set how many tasks run concurrently inside one container, and those tasks share the container's single memory grant. Doubling cores roughly doubles peak concurrent working set for the same memory. Reducing cores per executor is therefore often a cheaper fix than adding memory: it lowers peak usage, and smaller containers also pack better onto nodes and are less likely to be unschedulable.
saying these in an interview costs you the question
- Assumes the container limit applies only to the JVM heap
- Disables the physical memory check as the first fix
- Ignores whether the message named physical or virtual memory
- Raises container memory when a single container dies from skew
- Forgets that Python worker processes count against the same container