YARN (Yet Another Resource Negotiator) is the part of Hadoop that decides who gets which slice of which machine. It arrived in Hadoop 2 to split the old MapReduce JobTracker into two jobs: cluster-wide resource arbitration, done by one ResourceManager, and per-job scheduling and recovery, done by an ApplicationMaster that every application brings with it. That split is why Spark, Tez, Flink, MapReduce and long-running services can all share one pool of machines.
This article is the operator's view. It explains the resource model YARN actually enforces, walks through how a request turns into a running container, works through node and container sizing with real numbers, compares the schedulers, and covers memory enforcement, high availability and the handful of commands and failure modes you will meet every week. The combined HDFS and YARN architecture and job flow are covered in HDFS and YARN architecture.
The three roles
ResourceManager (RM). One active RM per cluster holds the global view: every node's capacity, every queue, every application. It has two main parts. The ApplicationsManager accepts submissions, starts each application's first container (its ApplicationMaster) and restarts it if it fails. The Scheduler decides which application gets the next free resources, according to queue capacities and policies. The Scheduler deliberately does nothing else: it does not monitor tasks or retry them.
NodeManager (NM). One per worker. It reports the node's resources and health to the RM in a heartbeat (every 1,000 ms by default), launches containers when an ApplicationMaster asks, enforces their limits, runs auxiliary services such as the Spark and MapReduce shuffle services, and uploads logs when a container finishes. If the RM hears nothing from an NM for ten minutes (yarn.nm.liveness-monitor.expiry-interval-ms, 600,000 by default) it marks the node LOST and treats its containers as gone.
ApplicationMaster (AM). One per application, running in a container like any other. It asks the Scheduler for containers, decides what to run in each, tracks progress and handles task failures. Because the AM is application code, YARN does not need to understand MapReduce or Spark; it only hands out containers. The request protocol and locality handling are covered in the ApplicationMaster in depth.
The resource model: memory and vcores
A container is a promise of an amount of each resource type on one node. By default there are two types: memory in megabytes and virtual cores. Each NM advertises its totals through yarn.nodemanager.resource.memory-mb and yarn.nodemanager.resource.cpu-vcores. In Hadoop 3 both default to -1, which means the NM falls back to a fixed default unless hardware detection is switched on (yarn.nodemanager.resource.detect-hardware-capabilities, false by default). Do not rely on either: set both explicitly for your hardware. GPUs, FPGAs and custom types can be added through the resource-types framework.
Requests are bounded and rounded. The Scheduler rejects anything above yarn.scheduler.maximum-allocation-mb (8,192 by default) or maximum-allocation-vcores (4), and rounds memory up to a multiple of minimum-allocation-mb (1,024) in the Capacity Scheduler. A 9,011 MB request therefore consumes 9,216 MB of the node's budget. The defaults are small; any cluster running Spark executors with large heaps must raise the maximum.
The most important and least obvious setting is how the Capacity Scheduler counts. Its default DefaultResourceCalculator looks only at memory. Vcores are recorded but not used to decide whether a container fits, so a node can be promised far more cores than it has. Switching to DominantResourceCalculator makes it schedule on whichever resource is scarcer for each request, following the dominant resource fairness idea. Most clusters that run CPU-heavy work want it.
<!-- yarn-site.xml on a worker with 128 GiB RAM and 32 cores -->
<property><name>yarn.nodemanager.resource.memory-mb</name><value>114688</value></property>
<property><name>yarn.nodemanager.resource.cpu-vcores</name><value>30</value></property>
<property><name>yarn.scheduler.minimum-allocation-mb</name><value>1024</value></property>
<property><name>yarn.scheduler.maximum-allocation-mb</name><value>32768</value></property>
<property><name>yarn.scheduler.maximum-allocation-vcores</name><value>8</value></property>
<property><name>yarn.log-aggregation-enable</name><value>true</value></property>
<!-- capacity-scheduler.xml: make the scheduler count CPU as well as memory -->
<property>
<name>yarn.scheduler.capacity.resource-calculator</name>
<value>org.apache.hadoop.yarn.util.resource.DominantResourceCalculator</value>
</property>
From request to running container
- The client submits an application to the RM with a queue name and a launch context for the AM: command, environment, local resources such as jars on HDFS, and the AM's resource size.
- The RM checks the queue's ACLs and limits, admits the application and asks the Scheduler for one AM container. When a node's heartbeat reports space, the Scheduler allocates it, and the RM tells that NM to launch the AM.
- The AM registers with the RM and starts calling
allocate, listing the containers it wants by size, count and preferred nodes or racks. The same call returns containers granted since the last one, completed containers and preemption notices. - For each granted container the AM talks directly to that NM to launch its task. The NM localises files, sets up the working directory, applies limits and starts the process.
- When a container exits, the NM reports its status through its heartbeat and the RM passes it to the AM on its next allocate. When the AM finishes it unregisters, and the NMs aggregate logs if that is enabled.
If the AM itself dies, the RM starts a new attempt, up to yarn.resourcemanager.am.max-attempts (2 by default; frameworks can set a lower per-application value). Whether the new attempt keeps the old one's running containers depends on the framework; MapReduce and Spark behave differently here.
A worked example: sizing a node for Spark
Take a worker with 128 GiB of RAM and 32 cores that also runs an HDFS DataNode. Reserve memory for the operating system, the DataNode, the NodeManager itself and page cache: 16 GiB is a reasonable starting point, which leaves 112 GiB, or 114,688 MB, for containers. Reserve two cores for the same daemons and advertise 30 vcores. Now size Spark executors with an 8 GiB heap and 4 cores each.
import math
def normalise(mb, min_mb=1024):
# Round a request up to a multiple of the minimum allocation.
return math.ceil(mb / min_mb) * min_mb
def containers_per_node(node_mb, node_vcores, req_mb, req_vcores, count_cpu):
by_mem = node_mb // req_mb
return min(by_mem, node_vcores // req_vcores) if count_cpu else by_mem
executor_heap = 8 * 1024 # spark.executor.memory = 8g
overhead = max(384, int(0.10 * executor_heap)) # Spark's default overhead rule
req = normalise(executor_heap + overhead) # 9011 MB -> 9216 MB
print(req)
print(containers_per_node(114688, 30, req, 4, count_cpu=False)) # 12, 48 cores promised
print(containers_per_node(114688, 30, req, 4, count_cpu=True)) # 7Spark adds overhead for off-heap memory, by default the larger of 384 MB or 10% of the heap, so each executor asks YARN for 9,011 MB, which is rounded up to 9,216 MB. By memory alone, 12 executors fit on the node. With the default memory-only calculator the Scheduler will place all 12, promising 48 cores on a 30-core machine; tasks will then compete for CPU, garbage collection slows and stages run long for no visible reason. With DominantResourceCalculator the node takes 7 executors and roughly 50 GiB of memory sits idle. The right fix is to change the executor shape so the resources run out together: 30 vcores and 112 GiB is about 3.7 GiB per core. Three-core executors with a 10 GiB heap ask for 10,240 + 1,024 = 11,264 MB, already a multiple of 1,024, so ten of them use all 30 vcores and 110 of the 112 GiB, with both dimensions running out together.
The general rule: compute the node's memory-per-vcore ratio and shape containers to match it. Whatever does not match is stranded. Clusters with several hardware generations have several ratios; either give each node type its own settings and accept some stranding, or steer shape-sensitive workloads to matching nodes with node labels or placement constraints. Remember too that the ApplicationMaster needs a container of its own, so a job asking for ten executors actually occupies eleven containers.
Schedulers: Capacity versus Fair
The RM runs one pluggable scheduler for the whole cluster; the Apache default in Hadoop 3 is the Capacity Scheduler.
| Capacity Scheduler | Fair Scheduler | |
|---|---|---|
| Model | Hierarchical queues with guaranteed percentages or absolute amounts | Hierarchical queues sharing resources by weight |
| Idle capacity | Borrowed up to each queue's maximum capacity | Shared automatically among active queues |
| Within a queue | FIFO or fair ordering policy | Fair, FIFO or DRF policy per queue |
| Reclaiming borrowed capacity | Preemption, when enabled | Preemption, when enabled |
| Typical fit | Multi-tenant clusters with contractual shares per team | Mixed interactive and batch work where equal sharing is the goal |
Both let queues lend unused capacity, and both can take it back by preempting containers when the owner needs it; without preemption a queue that lent its share may wait for other jobs to finish. Details are in the Capacity Scheduler, the Fair Scheduler and YARN preemption.
Enforcement: why containers get killed
A container's size is a contract, and the NM enforces the memory half. By default it checks both physical memory (yarn.nodemanager.pmem-check-enabled, true) and virtual memory (vmem-check-enabled, true, with a ratio of 2.1 times physical). A process tree that exceeds either is killed, and the AM sees an exit status of KILLED_EXCEEDED_PMEM or KILLED_EXCEEDED_VMEM with a diagnostic such as "running beyond physical memory limits".
The virtual memory check causes many false kills, because modern JVMs and glibc reserve large address ranges they never touch. Many operators disable it or raise the ratio. Physical memory kills, by contrast, are real: the heap plus off-heap usage (direct buffers, thread stacks, native libraries, Python workers for PySpark) exceeded the container. Fix them by raising the overhead setting or shrinking the heap, not by disabling the check. CPU is not limited at all unless the Linux container executor with cgroups is configured; see YARN and cgroups.
High availability and restart
A single RM is a single point of failure for scheduling, though not for running containers. Production clusters run an active and a standby RM coordinated through ZooKeeper, and store application state in ZooKeeper so the standby can resume. RM recovery is off by default (yarn.resourcemanager.recovery.enabled, false). Once it is on, the recovery is work-preserving by default: the new active RM rebuilds its view from the state store and from NM re-registrations, and running containers keep running. NodeManager recovery, also off by default, lets an NM restart for an upgrade without killing its containers. It needs a local state directory, supervised mode so containers are not cleaned up when the NM exits, and a fixed port in yarn.nodemanager.address: the default ephemeral port would change across the restart and the NM could not reclaim its containers.
<property><name>yarn.resourcemanager.ha.enabled</name><value>true</value></property>
<property><name>yarn.resourcemanager.ha.rm-ids</name><value>rm1,rm2</value></property>
<property><name>yarn.resourcemanager.hostname.rm1</name><value>master1.example.com</value></property>
<property><name>yarn.resourcemanager.hostname.rm2</name><value>master2.example.com</value></property>
<property><name>yarn.resourcemanager.cluster-id</name><value>prod-yarn</value></property>
<property><name>hadoop.zk.address</name><value>zk1:2181,zk2:2181,zk3:2181</value></property>
<property><name>yarn.resourcemanager.recovery.enabled</name><value>true</value></property>
<property>
<name>yarn.resourcemanager.store.class</name>
<value>org.apache.hadoop.yarn.server.resourcemanager.recovery.ZKRMStateStore</value>
</property>
<property><name>yarn.nodemanager.recovery.enabled</name><value>true</value></property>
<property><name>yarn.nodemanager.recovery.dir</name><value>/var/lib/hadoop-yarn/nm-recovery</value></property>
<property><name>yarn.nodemanager.recovery.supervised</name><value>true</value></property>
<property><name>yarn.nodemanager.address</name><value>0.0.0.0:45454</value></property>
Operating it: commands and failure modes
yarn node -list -all # node states: RUNNING, UNHEALTHY, LOST, DECOMMISSIONED
yarn queue -status default # capacity, used capacity, state
yarn application -list -appStates RUNNING # what is running, in which queue
yarn application -status application_1727650000000_0042
yarn logs -applicationId application_1727650000000_0042 | less # needs log aggregation
yarn rmadmin -refreshQueues # reload capacity-scheduler.xml without restart
yarn rmadmin -getServiceState rm1 # active or standby
yarn top # live per-queue and per-app usage| Symptom | Likely cause | First check |
|---|---|---|
| Application stuck in ACCEPTED | Queue or user AM limit reached, or no node has room for the AM | Queue status and maximum-am-resource-percent; AM size versus free node capacity |
| Request rejected at submit | Container larger than the maximum allocation | Raise the scheduler maximum or shrink the request |
| Containers killed for memory | Off-heap usage beyond the container size | Container diagnostics; raise overhead |
| Nodes turn UNHEALTHY | Local or log disks full or failed | yarn node -list -all; NM disk health checker thresholds |
| Jobs slow, cluster "not full" | CPU oversubscribed under memory-only accounting | Switch calculator or reshape containers |
| No logs after a job fails | Log aggregation disabled (the default) | Enable it and set a retention period |
What to do next
- Set
resource.memory-mbandresource.cpu-vcoresexplicitly on every node type, after reserving memory and cores for daemons. - Compute each node type's memory-per-vcore ratio and reshape your largest workloads' containers to match it.
- Decide whether CPU should be scheduled; if yes, switch the Capacity Scheduler to
DominantResourceCalculator. - Raise the maximum allocation to fit your largest legitimate container, and no further.
- Enable log aggregation, RM HA with ZooKeeper, RM recovery and NM recovery.
- Review your queue tree and preemption settings against who is actually waiting, using
yarn topand queue metrics. - Read the ApplicationMaster article next if you write or debug YARN applications.