Every worker in a Hadoop YARN cluster runs one NodeManager. The ResourceManager decides where work goes, but it never touches a worker directly: everything it knows about a machine comes from that machine's NodeManager, and every process that runs there is started, watched and cleaned up by it. When a cluster looks smaller than the hardware, when jobs cluster on a few nodes, or when a node goes UNHEALTHY at 3 a.m., the explanation is almost always in the NodeManager's configuration or its reports.
This article treats the NodeManager as a daemon you operate. The container lifecycle itself (launch contexts, localization, container executors, exit codes) is covered in YARN Containers, in depth, and kernel-level isolation in YARN cgroups. Here the focus is what the NodeManager advertises, what it says in each heartbeat, how it decides it is sick, how it survives its own restart, and how to take it out of service without killing work. Property names and defaults were checked against Apache Hadoop's yarn-default.xml and NodeManager documentation; check them against your distribution's version before copying.
What the NodeManager owns, and what it does not
A NodeManager is a single Java process made of cooperating services. The NodeStatusUpdater registers the node with the ResourceManager and sends heartbeats. The ContainerManager serves the ContainerManagementProtocol: ApplicationMasters call it to start, stop and query containers. The ContainersMonitor samples each container's memory and CPU and kills those over their limits. The NodeHealthCheckerService combines a disk checker with optional administrator scripts. AuxServices hosts long-lived plug-ins such as the MapReduce shuffle. Log aggregation and a deletion service upload logs and clean local directories, and a web server on port 8042 exposes a UI and REST API.
The NodeManager does not schedule. It never decides which application gets its memory; it only reports capacity and enforces what it was told. A subtle consequence is shown in the diagram: the ResourceManager allocates a container to an application, returns it to the ApplicationMaster with a signed container token, and the ApplicationMaster then calls the NodeManager directly to launch it. The NodeManager verifies the token using a master key the ResourceManager distributed in heartbeat responses. If an ApplicationMaster never calls, the allocation expires at the ResourceManager and the node never runs anything.
Advertising capacity: the two numbers that size your cluster
The ResourceManager's view of a node's size is whatever the NodeManager registers, and two properties dominate: yarn.nodemanager.resource.memory-mb and yarn.nodemanager.resource.cpu-vcores. Both default to -1. With yarn.nodemanager.resource.detect-hardware-capabilities left at its default of false, -1 means 8,192 MB and 8 vcores regardless of the hardware. A 256 GB machine that nobody configured therefore contributes 8 GB to the cluster, which is the most common reason a new cluster looks tiny.
Hardware detection can compute the values, but explicit values are easier to reason about on nodes that also run a DataNode or other agents. Budget memory for the operating system, co-located daemons and the NodeManager heap, then give YARN the rest.
<!-- yarn-site.xml on a 256 GB / 48-core worker that also runs a DataNode -->
<property>
<name>yarn.nodemanager.resource.memory-mb</name>
<value>245760</value> <!-- 240 GB: 16 GB left for OS, DataNode, NM, page cache headroom -->
</property>
<property>
<name>yarn.nodemanager.resource.cpu-vcores</name>
<value>44</value> <!-- 4 cores kept for daemons and interrupts -->
</property>
<property>
<name>yarn.nodemanager.address</name>
<value>${yarn.nodemanager.hostname}:45454</value> <!-- fixed port: required for NM recovery -->
</property>Worked example. A team runs Spark executors with spark.executor.memory=8g and 4 cores. Spark's default overhead is the larger of 384 MB and 10 per cent of executor memory, so each container requests 8,192 + 819 = 9,011 MB. The scheduler normalizes requests up to a multiple of its minimum allocation; with 1,024 MB that becomes 9,216 MB. 245,760 / 9,216 = 26.7, so memory allows 26 executors per node.
Now CPU. The CapacityScheduler's default DefaultResourceCalculator looks only at memory, so it will place all 26 executors, asking for 104 cores on a node that advertised 44. The result is heavy CPU contention and tasks that run far slower than their sizing suggests. Switching the scheduler to DominantResourceCalculator makes vcores a real constraint: 44 / 4 = 11 executors per node, leaving about 140 GB of memory idle. The fix is to match executor shape to node shape: here, 4 cores with about 20 GB per executor uses both resources evenly.
Heartbeats: the only channel to the ResourceManager
At start-up the NodeStatusUpdater calls registerNodeManager with the node's address, HTTP port, total resources and the status of any containers it recovered. After that it heartbeats at the interval the ResourceManager dictates, by default yarn.resourcemanager.nodemanagers.heartbeat-interval-ms = 1,000 ms. The ResourceManager returns the next interval in every response, so it can slow nodes down under load.
Each heartbeat carries the status of every container on the node (running, or completed with an exit status and diagnostics), the node's health report, and resource utilization. The response carries lists of containers and applications to clean up, updated token master keys, and occasionally an action: RESYNC after a ResourceManager restart, telling the node to re-register and report its running containers, or SHUTDOWN when the node has been decommissioned or rejected.
Liveness is tracked from the other side. If the ResourceManager hears nothing for yarn.nm.liveness-monitor.expiry-interval-ms, 600,000 ms by default, it marks the node LOST and reports every container on it to its ApplicationMaster as completed with an aborted status. Ten minutes is deliberately long, so that a garbage-collection pause or brief partition does not reschedule hundreds of containers. The cost is that ApplicationMasters wait that long before rerunning a dead node's tasks; shortening it on a noisy network declares running containers lost, which is worse.
Health: the disk checker and your scripts
A node can be alive and still unfit for work. The NodeManager decides that with two inputs. When a node reports itself UNHEALTHY, the ResourceManager removes it from the scheduler, and the scheduler kills the node's running containers as lost; this is not just a stop on new placements. In-flight work is rerun elsewhere, which is why a false health alarm is expensive.
The disk checker runs every yarn.nodemanager.disk-health-checker.interval-ms (two minutes) over every directory in yarn.nodemanager.local-dirs and yarn.nodemanager.log-dirs. A directory fails if it cannot be written, if its disk is more than max-disk-utilization-per-disk-percentage full (default 90), or if free space drops below min-free-space-per-disk-mb (default 0). Failed directories are removed from use; if the healthy fraction falls below min-healthy-disks (default 0.25), the whole node is unhealthy. A node with one data disk is therefore one full disk away from leaving the cluster, and a job that writes enormous spill files can push nodes out one by one.
Health scripts cover everything else. yarn.nodemanager.health-checker.scripts lists keywords (default script), and each keyword gets its own .path and .opts properties. The rules are unusual: the node is marked unhealthy if the script prints a line beginning with ERROR, times out, or throws an exception. A non-zero exit code is not a failure. The defaults run every 10 minutes with a 20-minute timeout, which is slow for most purposes.
#!/usr/bin/env bash
# /etc/hadoop/nm-health.sh -- runs every interval-ms; any line starting with ERROR marks
# the node unhealthy. Exit codes are ignored, so print, don't exit 1.
# Only test things that are LOCAL to this node. Never test a shared dependency.
# 1. Is the kernel reporting memory hardware errors? (dmesg may need root when
# kernel.dmesg_restrict=1; if so, read a log the NM user can access instead)
if dmesg 2>/dev/null | tail -n 500 | grep -q "Hardware Error"; then
echo "ERROR hardware error in dmesg"
fi
# 2. Is the local Kerberos keytab still readable by the NM user?
if ! test -r /etc/security/keytabs/nm.service.keytab; then
echo "ERROR nm keytab unreadable"
fi
# 3. Are we about to run out of inodes on the root filesystem?
used=$(df -Pi / | awk 'NR==2 {gsub("%","",$5); print $5}')
if [ "${used:-0}" -ge 95 ]; then
echo "ERROR root inode usage ${used}%"
fi
echo "OK"<property><name>yarn.nodemanager.health-checker.scripts</name><value>script</value></property>
<property><name>yarn.nodemanager.health-checker.script.path</name><value>/etc/hadoop/nm-health.sh</value></property>
<property><name>yarn.nodemanager.health-checker.script.opts</name><value></value></property>
<property><name>yarn.nodemanager.health-checker.interval-ms</name><value>120000</value></property>
<property><name>yarn.nodemanager.health-checker.timeout-ms</name><value>60000</value></property>The most important rule is in the script's comment: test only local conditions. A script that checks an NFS mount, a DNS server or a metastore will, when that shared dependency fails, mark every node unhealthy at once, and the cluster runs nothing while the shared dependency is down.
Work-preserving restart
Without recovery, restarting a NodeManager kills every container on the node. Recovery changes that: the NodeManager records container, application and localization state in a LevelDB store under yarn.nodemanager.recovery.dir, and on restart it reacquires running containers instead of killing them. Containers are separate operating-system processes, and the launch script writes each container's exit code to a file, so a restarted NodeManager can learn whether a container finished while it was down.
<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>
<!-- plus the fixed yarn.nodemanager.address port shown earlier -->
<!-- shuffle for MapReduce and Spark's external shuffle service -->
<property><name>yarn.nodemanager.aux-services</name><value>mapreduce_shuffle,spark_shuffle</value></property>
<property><name>yarn.nodemanager.aux-services.mapreduce_shuffle.class</name>
<value>org.apache.hadoop.mapred.ShuffleHandler</value></property>
<property><name>yarn.nodemanager.aux-services.spark_shuffle.class</name>
<value>org.apache.spark.network.yarn.YarnShuffleService</value></property>Three details decide whether recovery actually works. First, yarn.nodemanager.address must use a fixed port. The default port of 0 picks a random port, so the restarted daemon would look like a different node and the ResourceManager and ApplicationMasters would lose track of its containers. Second, recovery.supervised=true tells the NodeManager that a supervisor will restart it, so a normal shutdown leaves containers running rather than cleaning them up. Third, the restart must finish well within the liveness expiry, or the ResourceManager marks the node LOST and the preserved containers are reported lost anyway. Upgrading one node at a time with recovery enabled is how rolling NodeManager upgrades avoid rerunning work.
Auxiliary services: the shuffle lives here
Auxiliary services run inside the NodeManager process for the lifetime of the node, not of any container. The MapReduce ShuffleHandler serves map output to reducers after the map containers have exited; Spark's external shuffle service does the same for Spark executors, which is what lets dynamic allocation release executors without losing their shuffle files.
Because they share the NodeManager's JVM, their failures become node failures: a missing or mismatched Spark shuffle jar stops the NodeManager from starting, and a busy shuffle service competes for its heap. Size the NodeManager heap for it and upgrade shuffle jars with the NodeManager.
Graceful decommission
Removing a node by stopping its NodeManager kills its containers and, for MapReduce, also discards map output that reducers still need, so the maps rerun elsewhere. Graceful decommission avoids that. List the host in the file named by yarn.resourcemanager.nodes.exclude-path, then refresh with a timeout:
# Which nodes are not RUNNING, and why?
yarn node -list -states UNHEALTHY,LOST,DECOMMISSIONING -showDetails
# Everything the RM knows about one node, including the health report string
yarn node -status worker-17.example.com:45454
# The NM's own view, straight from its REST API on port 8042
curl -s http://worker-17.example.com:8042/ws/v1/node/info | python -m json.tool
# Gracefully drain a node: add it to the exclude file, then
yarn rmadmin -refreshNodes -g 3600 -clientThe node moves to DECOMMISSIONING. The scheduler stops allocating there, running containers continue, and the ResourceManager's watcher (polling every 20 seconds by default) decommissions the node once it is drained or the timeout expires. With -client the rmadmin command itself waits and enforces the timeout; with -server the ResourceManager enforces it, defaulting to yarn.resourcemanager.nodemanager-graceful-decommission-timeout-secs = 3,600 s. Choose the timeout from your longest normal container runtime, not from how long you are willing to wait: a long-running service container will simply be killed at the deadline.
Observability
Four views cover most investigations. yarn node -list -all shows every node's state; filter for UNHEALTHY, LOST and DECOMMISSIONING daily. yarn node -status prints the health report string, which contains the script's ERROR line or the disk checker's explanation. The REST API on port 8042 shows the NodeManager's own view, useful when it disagrees with the ResourceManager. Chart containers launched, failed and killed per node from JMX: a node that fails many containers is usually a hardware problem hiding behind application errors.
Failure modes and trade-offs
- Unconfigured capacity. Defaults of 8 GB and 8 vcores on a large machine. Check the Nodes page totals against the hardware after every build.
- Memory-only scheduling on CPU-bound work. The default calculator ignores vcores; match executor shape to node shape or use the dominant-resource calculator.
- Cluster-wide unhealthy. A health script tests a shared service, or every node's log disk crosses 90 per cent together, and running containers everywhere are killed. Test local conditions only and alarm on disk usage well before 90.
- Recovery that does not recover. Ephemeral port, unsupervised shutdown, or a restart longer than the liveness expiry.
- Physical-memory kills. Containers exceeding their size are killed by the monitor; the fix is usually more overhead memory, not disabling the check.
yarn.nodemanager.vmem-pmem-ratio(2.1) causes spurious virtual-memory kills with some JVMs; many clusters disable the virtual check and keep the physical one.
The overall trade-off is responsiveness against stability. Aggressive timeouts and health checks react quickly to real failures and to transient ones, and every false positive reschedules work; lax ones waste capacity on dead nodes. For how the ResourceManager uses these reports, see the YARN ResourceManager; for how applications react to lost containers, see the ApplicationMaster.
What to do next
- Compare each node's registered memory and vcores with its hardware, and set both explicitly after budgeting for co-located daemons.
- Work through the executor arithmetic above for your largest job and pick a resource calculator deliberately.
- Write a health script that checks only local conditions, prints ERROR lines, and runs every one or two minutes with a short timeout.
- Alarm on data-disk utilization at 80 per cent so the disk checker never has to act.
- Enable NodeManager recovery with a fixed port and a supervised restart, then restart one node under load and confirm its containers survive.
- Rehearse a graceful decommission with a timeout based on real container runtimes.