What happens to running YARN applications when the active ResourceManager fails over in an HA cluster?
answer
- the work is not on the master
- who keeps the containers alive?
- state store holds apps, not containers
- everyone re-registers with the new boss
- you lose scheduling, not progress
basics
~20 sWith work-preserving restart, containers already running on NodeManagers keep running. The standby ResourceManager becomes active, rebuilds application state from its ZooKeeper state store, and NodeManagers and ApplicationMasters resync with it. Clients retry transparently; new scheduling stalls for the failover window.
solid answer
~40 sYARN HA runs two ResourceManagers, `rm1` and `rm2`, listed in `yarn.resourcemanager.ha.rm-ids`, with an embedded ZooKeeper-based elector choosing the active. Application state — submission contexts, attempts, tokens — is persisted to a `ZKRMStateStore` when `yarn.resourcemanager.recovery.enabled=true`. On failover the standby becomes active, reloads that state, and then **work-preserving restart** does the rest: it does not kill anything. NodeManagers re-register and report the containers they are still running, and ApplicationMasters resync, so a long MapReduce or Spark job survives the failover instead of restarting from zero. What you do lose is the gap: nothing is scheduled while the new active is rebuilding, the RM web UI and REST API blip, and submissions during the window rely on client retries through `ConfiguredRMFailoverProxyProvider`. Verify with `yarn rmadmin -getServiceState rm1`.
code
xml · 20 lines<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.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>go deeper
Know that production YARN runs two ResourceManagers, active and standby, and that a failover does not by itself kill running jobs.
Explain the pieces: ZooKeeper-based election, the ZKRMStateStore holding application state, and the resync by NodeManagers and ApplicationMasters that rebuilds the scheduler's view.
Show the operational judgment — what actually degrades during the window, why fine-grained scheduler state is deliberately not persisted, when repeated failovers mean a GC problem, and how NodeManager recovery enables rolling upgrades.
Decide how much this availability is worth: RM and ZooKeeper placement across failure domains, acceptable scheduling-gap targets, and whether long-running services on YARN need stronger guarantees than batch jobs do.
## Why ResourceManager HA is different from NameNode HA HDFS HA protects data availability: without a NameNode, no file can be read. YARN HA protects *scheduling*. A dead ResourceManager does not, by itself, destroy work in progress — the containers are running on NodeManagers, which do not need the RM moment to moment. The design goal is therefore not "take over in seconds so reads keep working" but "take over without killing what is already running". ## The HA setup ``` yarn.resourcemanager.ha.enabled = true yarn.resourcemanager.cluster-id = prod yarn.resourcemanager.ha.rm-ids = rm1,rm2 yarn.resourcemanager.hostname.rm1 = rm-a.example.com yarn.resourcemanager.hostname.rm2 = rm-b.example.com yarn.resourcemanager.zk-address = zk1:2181,zk2:2181,zk3:2181 ``` Unlike HDFS, YARN needs no separate failover-controller daemon: each ResourceManager embeds an elector that contests a ZooKeeper lock, and the winner transitions itself to active while the other stays in standby, redirecting web requests to the active. ## The state store Recovery only works if the new active knows what applications exist. With `yarn.resourcemanager.recovery.enabled=true` and `yarn.resourcemanager.store.class` set to `ZKRMStateStore`, the RM persists each application's submission context, its attempt records, and its security tokens to ZooKeeper as they change. (A filesystem-backed store exists, but ZooKeeper is the usual choice in HA because it is already there and gives strongly consistent small writes.) What is *not* in the store is just as important: the fine-grained scheduler state — which container is on which node, how much of each queue is in use — is not persisted. That is rebuilt from the cluster itself when NodeManagers re-register. ## Work-preserving restart This is the part interviewers are testing. Modern YARN performs a **work-preserving** restart: when the RM comes back (whether after a crash or a failover), it does **not** kill running containers. 1. The new active loads application state from the store and re-enters a recovery phase. 2. Every NodeManager notices the RM change, re-registers with the new active, and reports the containers it is currently running plus their resource usage. 3. From those reports the RM rebuilds its picture of cluster and queue utilisation. 4. ApplicationMasters resync with the new RM on their next heartbeat and continue negotiating containers. The practical effect is that a job that has been running for four hours keeps its completed work. Before work-preserving restart existed, an RM restart killed everything and applications were re-run from the start, which is the behaviour candidates often still describe. Note the AM's obligation: an ApplicationMaster must handle a resync instruction from the RM. The MapReduce and Spark AMs do. A custom AM that ignores it will misbehave across failover. ## What clients see Clients and AMs use a failover proxy (`ConfiguredRMFailoverProxyProvider`) built from the `rm-ids` list: on a connection failure they retry the other ResourceManager, bounded by `yarn.resourcemanager.connect.max-wait.ms` and the retry interval. A `yarn application -submit` issued exactly during the window will block and retry rather than fail outright, provided the wait budget covers the recovery time. ## What you actually lose - **A scheduling gap.** Nothing new is allocated until recovery completes, so short jobs queue and container churn stalls. Job *latency* suffers even though nothing is lost. - **UI and REST discontinuity.** Dashboards and monitoring that poll the RM see errors or redirects during the transition. - **Nothing from the RM's memory that was not in the store.** Anything an application kept only in the RM's in-memory view is reconstructed, not restored. ## The NodeManager equivalent A parallel feature exists one level down: NodeManager work-preserving restart, enabled with `yarn.nodemanager.recovery.enabled=true` and a `yarn.nodemanager.recovery.dir` for its state. It lets you restart a NodeManager — for a patch or a config change — without killing the containers on that node, provided `yarn.nodemanager.address` is pinned to a fixed port so containers can reconnect. This is what makes a rolling YARN upgrade tolerable on a busy cluster. ## Operating it `yarn rmadmin -getServiceState rm1` reports active or standby; `yarn rmadmin -transitionToStandby` exists for manual control when automatic failover is disabled. The things worth alerting on: both RMs standby (a ZooKeeper problem), state-store write errors, and repeated failovers, which usually mean the active is GC-pausing long enough to lose its ZooKeeper session rather than genuinely failing.
- If the ResourceManager's state store holds application state, how does the new active learn which containers are running?From the NodeManagers. On failover each NodeManager re-registers with the new active and reports the containers it currently hosts and their resource usage, letting the RM rebuild cluster and queue utilisation. That is why fine-grained scheduler state is deliberately not persisted — the cluster is the authoritative source, and reconstructing from it avoids the store going stale.
- How does a NodeManager restart differ from a ResourceManager failover for a running container?With `yarn.nodemanager.recovery.enabled=true` and a fixed `yarn.nodemanager.address`, a restarted NodeManager recovers its own state from `yarn.nodemanager.recovery.dir` and re-adopts the containers still running on that host. Without recovery enabled, restarting a NodeManager kills its containers outright — the application then loses those tasks and must re-run them.
- Repeated ResourceManager failovers with no host failure — what is the usual cause?A long stop-the-world GC pause on the active RM that outlives its ZooKeeper session timeout. The elector concludes the RM is gone and promotes the peer, and the cycle can repeat. Look at RM heap sizing and GC logs before touching the failover configuration, and check ZooKeeper health and session timeouts as the second suspect.
saying these in an interview costs you the question
- Says all running containers are killed on failover
- Thinks the standby ResourceManager schedules containers too
- Claims ZooKeeper holds the full scheduler state
- Believes a YARN failover needs a separate ZKFC daemon
- Says restarting a NodeManager is always safe for its containers