Almost every highly available Hadoop cluster depends on ZooKeeper, and few have measured what happens when it is needed. The first failover test passes on an idle cluster; the next real failover happens two years later at peak load, on a different version, with an ensemble now shared with HBase and Hive.

This article is about proving the design rather than describing it. The companion piece on ZooKeeper in Hadoop lists who stores what in the tree, and HDFS high availability explains the NameNode pair, the JournalNodes and fencing. Here we take those as given and work out the failover time budget term by term, look closely at the part most teams misconfigure (the YARN ResourceManager state store), and then run a game day: seven fault injections, what each one should do, how to measure it, and what a failed result tells you to fix.

Advertisement

What ZooKeeper decides, and what it does not

ZooKeeper's job in Hadoop is narrow. It answers one question per HA service: which of the candidates is allowed to be active right now. It does that with an ephemeral znode that exists only while the session that created it is alive, and with watches that tell the other candidate the moment that znode disappears. It stores no blocks, no namespace, no edit log and no job output.

That narrowness has two practical consequences that people forget. First, ZooKeeper going down does not take HDFS down. The Hadoop HA guide is explicit: if the ZooKeeper cluster crashes, no automatic failovers are triggered, but HDFS continues to run, and it reconnects when ZooKeeper returns. Second, ZooKeeper is not what keeps a deposed NameNode from corrupting the namespace. That job belongs to the JournalNodes, which reject writes from any writer with an older epoch. ZooKeeper chooses; the JournalNodes enforce. A design review that confuses the two will put effort in the wrong place.

ZooKeeper decides who is active; it never holds HDFS data or YARN jobs' outputZooKeeper ensemble3 or 5 servers, quorumZKFC 1health monitor + electorRM 1 (active)embedded electorRM 2 (standby)embedded electorZKFC 2health monitor + electorlock, breadcrumbwatch lock/yarn-leader-election and /rmstoreNameNode 1activeJournalNodes x3edit log quorum, epochsNameNode 2standbyhealth RPChealth RPCwrite editstail editsZooKeeper answers 'who should be active'; the JournalNodes' epoch check is what stops a deposed NameNode corrupting the logIf the ensemble loses quorum, both services keep running on their current active; only automatic failover stops
The HA control plane. Each ZKFC watches its own NameNode and competes for the lock; the ResourceManagers carry their own embedded electors and also keep application state in ZooKeeper.

The failover time budget

Failover time is a sum, and every term has its own owner. Writing it down before the game day tells you which number to look at when the result is slow.

T_failover = T_detect + T_fence + T_promote + T_client

T_detect   how long until the standby side learns the active is gone
T_fence    how long the winner spends making sure the old active cannot act
T_promote  the standby NameNode finishing the edit-log tail and leaving safe states
T_client   clients' failover proxies noticing and retrying against the new active

The detection term depends entirely on what failed, and this is the part most runbooks get wrong.

What failedWho noticesDetection time is roughly
NameNode process dies, host and ZKFC aliveLocal ZKFC's health check fails; it quits the election and its lock znode is removedA few seconds: one health-check interval plus the RPC failure
NameNode hangs (long GC pause, stuck disk)Local ZKFC's health RPC times outha.health-monitor.rpc-timeout.ms, which is much longer than a crash
Whole host dies or is cut off from ZooKeeperNobody locally; ZooKeeper expires the ZKFC's sessionha.zookeeper.session-timeout.ms, 10,000 ms by default in current releases
ZKFC process dies, NameNode healthyZooKeeper expires the ZKFC sessionSession timeout, then a failover of a healthy NameNode

Note the session timeout default. Older Hadoop releases defaulted to 5 seconds, and the HA guide still says so; HADOOP-15449 raised it to 10 seconds because short timeouts caused needless failovers whenever a ZooKeeper disk or the network was briefly slow. Check the value your cluster actually runs with rather than trusting either document.

Fencing is cheap when it works and expensive when it does not. The winner reads the persistent ActiveBreadCrumb znode to learn who was active last and, if it was the other node, runs the configured fencing methods in order. sshfence against a powered-off host waits for the SSH connection timeout before falling through to the next method. If you rely on shell(/bin/true) as the last method (common with Quorum Journal Manager, because the JournalNode epoch already protects the log), fencing costs nothing, but you have accepted that the old NameNode may keep answering stale reads until it notices it has lost its epoch.

Advertisement

YARN: the embedded elector and the state store

The ResourceManager does not use a separate failover controller. With yarn.resourcemanager.ha.automatic-failover.embedded (true by default) each RM runs an elector in-process and competes for a lock under yarn.resourcemanager.ha.automatic-failover.zk-base-path, which defaults to /yarn-leader-election, in a child named after yarn.resourcemanager.cluster-id. Two clusters sharing an ensemble with the same cluster id will elect one RM between them, which is a memorable way to learn that the id must be unique.

Election alone only gives you a new RM that knows nothing. For running applications to survive, the RM must persist state, and two defaults work against you: yarn.resourcemanager.recovery.enabled is false, and yarn.resourcemanager.store.class defaults to the FileSystemRMStateStore. For HA you want recovery on and the ZooKeeper store, which writes beneath yarn.resourcemanager.zk-state-store.parent-path (default /rmstore).

<!-- yarn-site.xml: the minimum for RM HA that preserves running applications -->
<property><name>yarn.resourcemanager.ha.enabled</name><value>true</value></property>
<property><name>yarn.resourcemanager.cluster-id</name><value>prod-east-1</value></property>
<property><name>yarn.resourcemanager.ha.rm-ids</name><value>rm1,rm2</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>
<!-- core-site.xml in Hadoop 3: the ensemble address used by the RM -->
<property><name>hadoop.zk.address</name><value>zk1:2181,zk2:2181,zk3:2181</value></property>

The state store fences itself. When an RM becomes active it claims the root of the store through a ZooKeeper ACL (yarn.resourcemanager.zk-state-store.root-node.acl lets you set it explicitly) so that only it can create and delete children. A deposed RM that wakes up from a pause and tries to write gets an authorisation failure and steps down instead of overwriting its successor's state.

Two limits bite busy clusters. yarn.resourcemanager.zk-max-znode-size.bytes defaults to 1 MB, matching ZooKeeper's own jute.maxbuffer ceiling, so an application with a very large submission context or many attempts can fail to persist. And a flat directory of tens of thousands of application znodes makes the child listing returned on recovery enormous; yarn.resourcemanager.zk-appid-node.split-index (0, meaning no split, by default) spreads applications across intermediate nodes. Raise it before you need it, because recovery is exactly when a huge listing hurts.

The other tenants of the ensemble

The ensemble that elects your NameNode is rarely used by HDFS alone, and every tenant adds load and risk to the shared session timeout budget.

  • HBase keeps the active master's lock, region server liveness and the location of the meta table in ZooKeeper. It is by far the chattiest tenant, and a region server storm reconnecting after a network blip can slow every other client of the ensemble.
  • HiveServer2 can register instances for client-side discovery when hive.server2.support.dynamic.service.discovery is true, under the namespace in hive.server2.zookeeper.namespace. JDBC clients then connect with serviceDiscoveryMode=zooKeeper in the URL. Light load, but a dependency many teams forget exists until discovery breaks.
  • Kafka clusters often shared the ensemble; Kafka 4.0 removed ZooKeeper in favour of KRaft, so migrating takes that load off.

Election traffic is tiny and latency-critical; HBase traffic is large and bursty. With HBase present, give the HA services their own ensemble, or at least treat ensemble latency as an HDFS availability signal.

The game day: seven experiments

Run these on a staging cluster that matches production's versions and configuration, under a realistic load (a TeraSort or a replay of real jobs), with a stopwatch script running. Each experiment has an expected outcome; anything else is a finding.

#InjectionExpected outcomePass if
1kill -9 the active NameNode processLocal ZKFC notices, quits; standby wins and becomes activeWrites resume within your target, no job fails
2Power off the active NameNode hostSession expires after the timeout; standby fences via breadcrumb, then promotesTotal time close to session timeout plus fencing plus promotion
3kill -STOP the active NameNode for 60 s, then kill -CONTHealth RPC times out, failover happens; on resume the old NN is rejected by the JournalNodes and goes standby or abortsNever two actives writing; no corruption on fsck
4Stop the ZooKeeper leaderEnsemble re-elects in seconds; ZKFC and RM sessions surviveNo HDFS or YARN failover at all
5Stop two of three ZooKeeper servers for 5 minutesQuorum lost; actives keep serving; no failover possibleNo client-visible errors; reconnection clean afterwards
6kill -9 the active ResourceManagerStandby RM wins, loads /rmstore, running apps continueRunning jobs finish without resubmission
7Firewall the active NameNode host from ZooKeeper onlyIts ZKFC session expires; the other side wins and must fence a still-running NameNodeFencing succeeds or the epoch check stops the old active

Experiment 7 finds real bugs: the old NameNode is healthy and reachable by clients, which is split brain on purpose. Watch which fencing method fires and whether mid-write clients get a clean error.

Measuring it

Two small tools are enough. The first polls service state from outside and prints transitions with timestamps; hdfs haadmin -getAllServiceState reports every NameNode in the nameservice.

#!/usr/bin/env bash
# failover_clock.sh: print every HA state change with a millisecond timestamp
prev=""
while true; do
  now=$(hdfs haadmin -getAllServiceState 2>/dev/null | tr '\n' ' ')
  if [ "$now" != "$prev" ]; then
    echo "$(date +%H:%M:%S.%3N) $now"
    prev="$now"
  fi
  sleep 0.5
done

The second watches the lock znode itself, which tells you when ZooKeeper's view changed, separately from when the NameNode finished promoting. The difference between the two timestamps is your T_promote. This uses the kazoo Python client; the path follows the default /hadoop-ha parent and your nameservice id.

import time
from kazoo.client import KazooClient

NS = "mycluster"
LOCK = f"/hadoop-ha/{NS}/ActiveStandbyElectorLock"

zk = KazooClient(hosts="zk1:2181,zk2:2181,zk3:2181")
zk.start()

@zk.DataWatch(LOCK)
def on_change(data, stat):
    ts = time.strftime("%H:%M:%S")
    if stat is None:
        print(ts, "lock released: no active")
    else:
        # The payload is a serialized record; the hostname shows up as printable bytes.
        owner = "".join(chr(b) if 32 <= b < 127 else "." for b in data)
        print(ts, "lock held, session", hex(stat.ephemeralOwner), owner)

while True:
    time.sleep(1)

Worked example: experiment 2 on a loaded cluster

An illustrative timeline for powering off the active NameNode's host, with the default 10-second session timeout, sshfence followed by shell(/bin/true), and a write-heavy workload. Your numbers will differ; the shape is what matters.

t (s)EventBudget term
0.0Host nn1 loses power. Clients' in-flight RPCs start timing out.-
10.2ZooKeeper expires nn1's ZKFC session; the lock znode disappears; nn2's ZKFC watch fires.T_detect about 10 s
10.3nn2's ZKFC creates the lock, reads the breadcrumb, finds nn1, starts fencing.-
40.4sshfence gives up after its 30-second default connect timeout; shell(/bin/true) succeeds.T_fence about 30 s
41.9nn2 finishes tailing edits and transitions to active.T_promote 1.5 s
43 to 46Clients' failover proxies retry and reach nn2; writes resume.T_client 1 to 4 s

The finding writes itself: two thirds of the outage was sshfence waiting on a dead host. The fixes are a short dfs.ha.fencing.ssh.connect-timeout, or a power-fencing method that talks to the host's management controller and succeeds immediately for a powered-off host. Shortening the session timeout would cut the other half but trades it for spurious failovers, which is the subject of the trade-offs below.

Failure modes seen in real clusters

  • Failover never happens, or the RM cannot persist state. A missing or mis-parented HA znode and oversized state-store entries are covered in the failure table of ZooKeeper in Hadoop; experiments 2 and 6 surface both.
  • Failover flaps every few minutes. ZooKeeper's transaction log shares a disk with something busy, fsync latency spikes, and sessions expire. Put the ZooKeeper data log on its own device and alert on fsync time.
  • RM failover loses all running applications. Recovery is off or the store class is still the filesystem default. Experiment 6 catches this in minutes.
  • Kerberized cluster, sudden authentication failures. ZooKeeper SASL configuration or keytabs expired, and the ZKFCs cannot reach their own znodes. See Kerberos in Hadoop for the principal layout.
  • Session storm after a network blip. Thousands of HBase clients reconnect at once and starve the ZKFC sessions. Separate ensembles, or raise maxClientCnxns limits only after measuring.

Trade-offs

Short versus long session timeout. A short timeout cuts detection time for host failures but turns every ZooKeeper hiccup and every long JVM pause in the ZKFC into a failover, and an unnecessary failover is itself a short outage. Most clusters are better off keeping the 10-second default and shrinking the fencing term instead.

Shared versus dedicated ensemble. A shared ensemble is one less thing to run; a dedicated one isolates the latency-critical election traffic from bulk tenants. The cost of three small extra servers is low compared with an HBase reconnect storm causing a NameNode failover.

Automatic versus manual failover. Some operators disable automatic failover and use hdfs haadmin -failover by hand, accepting minutes of outage in exchange for never failing over by mistake. The underlying primitive, a lock tied to a session plus a separate fencing check, is the same one described generally in fencing tokens and ZooKeeper's architecture.

What to do next

  1. Write down your failover target, then fill in the four budget terms from your configuration: session timeout, fencing methods and their timeouts, and client retry settings.
  2. Confirm ha.zookeeper.session-timeout.ms as actually configured, not as documented, on every ZKFC host.
  3. Check the YARN state store: recovery enabled, ZKRMStateStore selected, a unique cluster id, and a split index sized for your application count.
  4. List every tenant of the ensemble and decide whether HDFS and YARN election should move to a dedicated one.
  5. Run the seven experiments on staging under load with the two measuring scripts, and record the timeline for each.
  6. Fix the largest budget term first (usually fencing), rerun the affected experiment, and repeat the game day after every Hadoop or ZooKeeper upgrade.
Key takeaway: ZooKeeper in Hadoop decides who is active and nothing more; the JournalNodes' epochs protect the namespace and the RM state store's ACL protects application state. Treat failover time as detection plus fencing plus promotion plus client retry, configure the YARN store so recovery actually works, keep election traffic away from noisy tenants, and prove all of it with a measured game day rather than a configuration review.