Yarn ResourceManager High Availability: Minimal Design and Operational Checks
A concise architecture note covering Yarn ResourceManager HA requirements, minimal two‑node design, trust boundaries, operational verification steps, typical failure modes, and when the design would need to change.
22 Oct 2025, 06:48 UTC

Requirements
The core requirement for Yarn ResourceManager (RM) HA is continuous scheduler availability so that running applications are not disrupted when the active RM fails. This demands a mechanism for automatic failover and a shared, highly available state store that both RM instances can read and write.
Smallest Suitable Design
The minimal HA deployment consists of two RM processes (one active, one standby) that coordinate leadership through a highly available store such as Apache ZooKeeper or an HA‑enabled HDFS namenode. Both RMs run the same version of Hadoop/Yarn and share the same configuration directory.
Example yarn-site.xml snippet
<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>rm1.example.com</value>
</property>
<property>
<name>yarn.resourcemanager.hostname.rm2</name>
<value>rm2.example.com</value>
</property>
<property>
<name>yarn.resourcemanager.ha.automatic-failover.enabled</name>
<value>true</value>
</property>
<property>
<name>yarn.resourcemanager.ha.automatic-failover.embedded</name>
<value>true</value>
</property>
<property>
<name>yarn.resourcemanager.ha.automatic-failover.zk-base-path</name>
<value>/yarn-leader-manager</value>
</property>
<property>
<name>yarn.resourcemanager.ha.automatic-failover.zk-address</name>
<value>zk1.example.com:2181,zk2.example.com:2181,zk3.example.com:2181</value>
</property>
Only the RM processes and the ZooKeeper Failover Controller (ZKFC) interact with the ZooKeeper ensemble; NodeManagers and application containers never access ZooKeeper directly, preserving isolation.
Trust and Data Boundaries
- RM processes – read/write leadership ephemeral nodes in ZooKeeper.
- ZKFC – monitors RM health and initiates failover by creating/deleting the leadership node.
- Application containers & NodeManagers – communicate with the active RM via RPC; they have no direct ZooKeeper access.
Operational Checks
- Verify each RM’s service state from any host with the Yarn client installed (typically run as the
yarnuser or with sudo):
Expected result: one RM returnsyarn rmadmin -getServiceState <rm-id>active, the otherstandby. - Check ZooKeeper session health using the ZooKeeper CLI (
zkCli.sh) on a ZooKeeper node:
You should see ephemeral nodes named after the RM IDs (e.g.,ls /yarn-leader-managerrm1,rm2). Only the node belonging to the current active RM will be present; the standby’s node is absent. - Inspect RM logs for state transition events (look for lines containing
Transitioning to activeorTransitioning to standby) to confirm that failover is being processed correctly.
Failure Modes
- Split‑brain – If both RMs lose their ZooKeeper session simultaneously (e.g., due to a network partition affecting the ZooKeeper ensemble), each may assume the active role, leading to duplicate scheduling.
- Failover delay – Determined by the ZooKeeper session timeout (
zookeeper.session.timeoutin the RM config). A longer timeout increases the window during which no RM is active. - Loss of shared state store – If the ZooKeeper ensemble becomes unavailable, neither RM can acquire leadership, resulting in a total loss of scheduling capability until store connectivity is restored.
- Network partition isolating one RM – The partitioned RM may remain in standby while the other stays active; however, if the partition also cuts off ZooKeeper access for the active RM, both may enter standby.
Conditions That Would Prompt a Redesign
- Need for sub‑second failover: the current design relies on ZooKeeper session timeouts (typically seconds). If faster recovery is required, a different consensus protocol (e.g., etcd with leader lease) might be considered.
- Scale‑out of RM instances beyond two active standbys: the current HA model assumes a single active RM. Supporting multiple concurrent active RMs would require a sharded scheduler or a different resource‑management architecture.
- Restrictions on using ZooKeeper due to operational overhead: if an organization cannot manage a ZooKeeper ensemble, migrating to an HA‑enabled HDFS namenode as the state store (or a cloud‑based lock service) would necessitate changes to the failover controller.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.