Every YARN cluster has exactly one active ResourceManager, and everything depends on it. Clients submit applications to it, ApplicationMasters ask it for containers, NodeManagers report to it, and administrators reconfigure queues through it. When the RM is slow, every job is slow; when it is down, nothing new starts. Yet many operators know it only as a web UI on port 8088 and a scheduler configuration file.

This article opens the RM up. It explains the services that face each kind of caller, the event-driven core that ties them together, how a NodeManager heartbeat becomes a container allocation, how high availability and restart recovery actually work, and what to watch and tune. It assumes you know the basic roles of YARN; the YARN overview covers the resource model, node sizing and scheduler choice, and the ApplicationMaster article covers the other side of the allocate protocol.

Advertisement

Four doors, four callers

Inside the ResourceManagerClientsyarn, Spark submitApplicationMastersallocate() callsNodeManagersheartbeatsAdminsyarn rmadmin:8032:8030:8031:8033ClientRMServiceApplicationMasterServiceResourceTrackerServiceAdminServiceAsyncDispatcher: one event queueRMApp, RMAppAttempt, RMContainer and RMNode state machinesNODE_UPDATE, APP_ADDEDSchedulerCapacity (default) or Fair; queues, allocationElectorActiveStandbyElectorRMStateStoreapps, attempts, tokensZooKeeper ensembleleader lock + /rmstoreStandby RMsame code, waiting for the lockwatches
Each kind of caller has its own RPC service and port. Services turn calls into events on one dispatcher, state machines process them, and the scheduler allocates. The active RM holds a ZooKeeper lock and persists application state there.

The RM separates its callers by service, each with its own port and thread pool, so a flood from one class of caller does not starve another. ClientRMService (default port 8032) handles application submission, status and kill requests from clients. ApplicationMasterService (8030) handles AM registration and the allocate heartbeat through which AMs ask for and receive containers. ResourceTrackerService (8031) handles NodeManager registration and heartbeats. AdminService (8033) handles administrative commands such as refreshing queues and HA transitions. The web UI and REST API are on 8088.

The three high-traffic services each default to 50 handler threads (yarn.resourcemanager.client.thread-count, scheduler.client.thread-count and resource-tracker.client.thread-count). On large clusters the resource-tracker and scheduler pools are the first to need more, because their call rate scales with the number of nodes and running applications rather than with human activity.

Around these services sit smaller components worth knowing by name because they appear in logs: RMAppManager admits submitted applications; ApplicationMasterLauncher tells a NodeManager to start an AM container; NMLivenessMonitor and AMLivenessMonitor expire nodes and AMs that stop heartbeating; NodesListManager applies include and exclude lists; and several secret managers issue the tokens that let AMs and containers authenticate.

The core: one dispatcher and four state machines

RPC handlers do very little work themselves. They validate the call, then post an event to the AsyncDispatcher, a queue drained by a dedicated thread that routes each event to its handler. This design keeps the RM's shared state consistent without fine-grained locks, and it means the RM's health is largely the health of that queue. When the RM logs that its event queue has grown into the tens of thousands, events are arriving faster than handlers process them, and every caller will feel the lag.

Most handlers are state machines. RMApp tracks an application through states such as NEW, NEW_SAVING, SUBMITTED, ACCEPTED, RUNNING, FINAL_SAVING, FINISHED, FAILED and KILLED. The SAVING states exist because the RM persists the application to its state store before acknowledging the transition, which is what makes recovery possible. RMAppAttempt tracks each attempt to run the AM; an application gets up to yarn.resourcemanager.am.max-attempts attempts, default 2. RMContainer tracks each container from ALLOCATED to ACQUIRED (the AM picked it up) to RUNNING and COMPLETED. RMNode tracks each NodeManager through RUNNING, UNHEALTHY, DECOMMISSIONING, DECOMMISSIONED and LOST.

Reading an application's diagnostics with these states in mind turns vague symptoms into precise ones. An application sitting in ACCEPTED has been admitted to a queue but its AM container has not been allocated, which is a capacity or queue-limit problem, not an RM fault.

Advertisement

From heartbeat to container

Scheduling in YARN is driven by node heartbeats. Each NodeManager calls the ResourceTrackerService every yarn.resourcemanager.nodemanagers.heartbeat-interval-ms, default 1000 ms, reporting completed containers and its health. The RM records the heartbeat, processes completions, and raises a node-update event for the scheduler. The scheduler then looks at the free resources on that node and walks its queues to decide which outstanding requests to satisfy there.

# Simplified: what the RM does with one NodeManager heartbeat
def on_node_heartbeat(node_id, status):
    node = rm_nodes[node_id]
    node.last_seen = now()                       # NMLivenessMonitor uses this
    for c in status.completed_containers:
        dispatch(ContainerFinishedEvent(c))      # RMContainer -> COMPLETED, AM told on next allocate
    dispatch(NodeUpdateSchedulerEvent(node))     # scheduler reacts asynchronously
    return HeartbeatResponse(
        containers_to_clean=node.pending_kills,  # e.g. preempted or from finished apps
        next_interval_ms=1000)                   # yarn.resourcemanager.nodemanagers.heartbeat-interval-ms

def on_scheduler_node_update(node):
    while node.unallocated_resources > 0:
        req = pick_request(queues, node)         # queue order, user limits, locality
        if req is None:
            break
        container = allocate(req, node)          # RMContainer ALLOCATED
        pending_for_am[req.app_attempt].append(container)
    # The AM receives these containers on ITS next allocate() call, then asks the NM to launch.

Allocation and launch are deliberately decoupled. The scheduler does not start anything; it marks a container as allocated and holds it until the owning AM's next allocate call, which returns it together with a token. The AM then asks the NodeManager to launch it. If the AM never picks the container up and launches it, the RM eventually reclaims it through a container-allocation expiry. This is why an AM that stops heartbeating wastes cluster capacity until it is expired.

Heartbeat-driven scheduling has a cost on large clusters: each allocation decision considers one node at a time, and a busy scheduler lock serialises them. The Capacity Scheduler can therefore also schedule asynchronously from background threads, controlled by yarn.scheduler.capacity.schedule-asynchronously.enable; the current documentation lists it as enabled by default, but older releases differ, so check the value your version actually runs with. Asynchronous scheduling raises allocation throughput on busy clusters at the price of more behaviour to reason about. Queue behaviour itself is covered in the Capacity Scheduler article and reclaiming capacity in the preemption article.

Nodes that stop heartbeating are declared LOST after yarn.nm.liveness-monitor.expiry-interval-ms, default 600,000 ms (ten minutes). Their containers are reported to AMs as failed, and the AMs decide whether to retry the work elsewhere. Ten minutes is conservative on purpose: declaring a node lost during a brief network blip would kill healthy work.

High availability: election and fencing

The RM is a single point of failure unless you run two or more with HA enabled. Each RM runs the same code; exactly one is active. Election uses ZooKeeper through an embedded ActiveStandbyElector: the RMs race to create an ephemeral lock znode, the winner becomes active, and the others watch the lock. If the active RM's ZooKeeper session expires, because the process died, the host lost the network or a long garbage-collection pause stopped heartbeats to ZooKeeper, the lock disappears and a standby takes over. Automatic failover and the embedded elector are both on by default once HA is enabled.

Election alone does not prevent two RMs from acting at once: an old active RM that was only paused may wake up believing it is still in charge. The ZooKeeper state store handles this by fencing. The active RM claims exclusive write access to the store's root node, so when a new RM takes over, the old one's writes fail and it transitions itself to standby. Clients, AMs and NodeManagers find the active RM with a failover proxy that tries the configured RM ids in turn.

<!-- yarn-site.xml: two ResourceManagers with automatic failover and recovery -->
<property><name>yarn.resourcemanager.ha.enabled</name><value>true</value></property>
<property><name>yarn.resourcemanager.cluster-id</name><value>prod-yarn</value></property>
<property><name>yarn.resourcemanager.ha.rm-ids</name><value>rm1,rm2</value></property>
<property><name>yarn.resourcemanager.hostname.rm1</name><value>rm1.example.com</value></property>
<property><name>yarn.resourcemanager.hostname.rm2</name><value>rm2.example.com</value></property>
<property><name>yarn.resourcemanager.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.resourcemanager.work-preserving-recovery.enabled</name><value>true</value></property>
<property><name>yarn.resourcemanager.max-completed-applications</name><value>1000</value></property>

The ensemble is named with yarn.resourcemanager.zk-address, as in the Hadoop HA documentation; some distributions set a shared ZooKeeper address in core-site instead, so follow your version's documentation. Note that recovery is off by default (yarn.resourcemanager.recovery.enabled is false) and the default store class is a filesystem store, so an HA setup that forgets these two lines fails over to an RM that has lost every running application.

Restart recovery, from the RM&amp;#x27;s side

Recovery has two layers. The state store persists what cannot be rebuilt from anyone else: each application's submission context, its attempts, final states, and the secret keys and delegation tokens the RM issued. Under ZooKeeper this lives beneath /rmstore by default. On becoming active, an RM loads this state and recreates every application.

Work-preserving recovery, on by default when recovery is enabled, rebuilds the rest from the cluster itself. NodeManagers keep their containers running while the RM is away, then resynchronise and report every live container. From those reports the new RM reconstructs the scheduler's view of which containers run where and charges them back to their queues. Running AMs re-register and carry on. To avoid handing out resources before that picture is complete, the RM waits yarn.resourcemanager.work-preserving-recovery.scheduling-wait-ms, default 10,000 ms, before scheduling new containers.

The number of completed applications kept in the store is bounded by yarn.resourcemanager.state-store.max-completed-applications, which defaults to the in-memory limit of 1,000. Raising either number makes history easier to browse and recovery slower, because every retained application is read back on failover.

Worked example: a failover, second by second

A cluster runs rm1 active and rm2 standby, with recovery and work preservation enabled. rm1's host suffers a kernel hang. Over the next few seconds, rm1 stops sending ZooKeeper heartbeats; once its session timeout passes, ZooKeeper deletes its ephemeral lock. rm2's watch fires, it wins the lock, fences the state store and loads roughly a thousand applications from /rmstore, which takes a few seconds on a healthy ensemble.

rm2 now starts its services. NodeManagers, whose heartbeats to rm1 have been failing, fail over to rm2, resynchronise and report their running containers; no container was killed. AMs' allocate calls fail over too, and each AM re-registers. For ten seconds rm2 accepts reports but does not schedule. Then it resumes. From a user's point of view, running Spark and MapReduce jobs paused briefly in acquiring new containers, and submissions during the gap were retried by the client. The postmortem question is how long the ZooKeeper session timeout was, because that dominates detection time.

Contrast the same event with recovery disabled: rm2 becomes active with no applications, NodeManagers are told to kill everything, and every job in the cluster restarts from scratch or fails.

Operating the RM: commands and signals

# Which RM is active?
yarn rmadmin -getServiceState rm1
yarn rmadmin -getServiceState rm2

# Reload queue definitions and include/exclude node lists without a restart
yarn rmadmin -refreshQueues
yarn rmadmin -refreshNodes

# Cluster-level health from the REST API (served by the active RM on :8088)
curl -s http://rm1.example.com:8088/ws/v1/cluster/metrics | jq '.clusterMetrics |
  {appsPending, appsRunning, activeNodes, lostNodes, unhealthyNodes,
   availableMB, allocatedMB, containersPending}'

# Applications stuck before their AM started
curl -s 'http://rm1.example.com:8088/ws/v1/cluster/apps?states=ACCEPTED' | jq '.apps.app[] |
  {id, queue, user, diagnostics}'

Watch a small set of signals. From the cluster metrics, track pending applications and pending containers against available memory: pending work with free capacity points to queue limits or placement constraints, not shortage. Track lostNodes and unhealthyNodes, which should be near zero. From the RM's JVM, watch heap use and garbage-collection pauses, since a long pause can expire the ZooKeeper session and cause a needless failover. From logs, watch the dispatcher's event-queue size and state-store operation times. The HDFS and YARN architecture article shows how these pieces sit in the wider cluster.

Failure modes

SymptomLikely causeFix
Failovers every few hours with no host failureRM garbage-collection pauses exceed the ZooKeeper session timeoutRight-size the heap, cap retained applications, tune GC; check the session timeout
All jobs lost after failoverRecovery disabled or filesystem store not sharedEnable recovery with ZKRMStateStore
Slow failover, minutes to become activeToo many applications retained in the state storeLower retained completed applications
State-store writes fail for some applicationsVery large application data exceeds ZooKeeper's znode size limitTrim oversized diagnostics or tokens; review the ZooKeeper buffer limit with care
All scheduling slow, RPC latency highDispatcher queue backlog or scheduler lock contentionConfirm async scheduling is on, add handler threads, reduce heartbeat load
Applications stuck in ACCEPTEDQueue AM-resource limit or no node fits the AM containerInspect queue limits and the application's diagnostics

Trade-offs

A single active scheduler gives YARN a global, consistent view that makes queue fairness and preemption straightforward, at the cost of a ceiling on cluster size and allocation rate. Federation lifts that ceiling by running several sub-clusters, each with its own RM, behind a router, but adds a routing layer and cross-cluster policy to operate. Shortening heartbeat intervals speeds allocation and loads the RM; lengthening liveness timeouts avoids false alarms and delays real recovery. Choose deliberately and write the choice down.

What to do next

  1. Check whether HA, recovery and ZKRMStateStore are all enabled; if any is missing, fix it before the next incident.
  2. Run yarn rmadmin -getServiceState for each RM id and confirm clients and NodeManagers list every RM id.
  3. Perform a planned failover in a maintenance window and time each phase against the worked example above.
  4. Dashboard pending containers, lost and unhealthy nodes, RM heap and GC pause time, and alert on repeated failovers.
  5. Review retained completed applications and the RM heap together, since both drive failover time.
  6. On large or busy clusters, confirm whether asynchronous scheduling is on and test larger RPC handler pools in staging before production.
Key takeaway: The ResourceManager is four RPC services feeding one event dispatcher, a set of state machines, and a scheduler driven by NodeManager heartbeats. Allocations are handed to AMs on their next allocate call, not launched by the RM. High availability needs both ZooKeeper election and recovery with ZKRMStateStore, because recovery is off by default. Work-preserving restart rebuilds container state from NodeManager reports, and failover time is dominated by ZooKeeper session timeout and the number of applications retained.