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.
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.
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 activeThe detection term depends entirely on what failed, and this is the part most runbooks get wrong.
| What failed | Who notices | Detection time is roughly |
|---|---|---|
| NameNode process dies, host and ZKFC alive | Local ZKFC's health check fails; it quits the election and its lock znode is removed | A few seconds: one health-check interval plus the RPC failure |
| NameNode hangs (long GC pause, stuck disk) | Local ZKFC's health RPC times out | ha.health-monitor.rpc-timeout.ms, which is much longer than a crash |
| Whole host dies or is cut off from ZooKeeper | Nobody locally; ZooKeeper expires the ZKFC's session | ha.zookeeper.session-timeout.ms, 10,000 ms by default in current releases |
| ZKFC process dies, NameNode healthy | ZooKeeper expires the ZKFC session | Session 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.
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.discoveryis true, under the namespace inhive.server2.zookeeper.namespace. JDBC clients then connect withserviceDiscoveryMode=zooKeeperin 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.
| # | Injection | Expected outcome | Pass if |
|---|---|---|---|
| 1 | kill -9 the active NameNode process | Local ZKFC notices, quits; standby wins and becomes active | Writes resume within your target, no job fails |
| 2 | Power off the active NameNode host | Session expires after the timeout; standby fences via breadcrumb, then promotes | Total time close to session timeout plus fencing plus promotion |
| 3 | kill -STOP the active NameNode for 60 s, then kill -CONT | Health RPC times out, failover happens; on resume the old NN is rejected by the JournalNodes and goes standby or aborts | Never two actives writing; no corruption on fsck |
| 4 | Stop the ZooKeeper leader | Ensemble re-elects in seconds; ZKFC and RM sessions survive | No HDFS or YARN failover at all |
| 5 | Stop two of three ZooKeeper servers for 5 minutes | Quorum lost; actives keep serving; no failover possible | No client-visible errors; reconnection clean afterwards |
| 6 | kill -9 the active ResourceManager | Standby RM wins, loads /rmstore, running apps continue | Running jobs finish without resubmission |
| 7 | Firewall the active NameNode host from ZooKeeper only | Its ZKFC session expires; the other side wins and must fence a still-running NameNode | Fencing 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
doneThe 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) | Event | Budget term |
|---|---|---|
| 0.0 | Host nn1 loses power. Clients' in-flight RPCs start timing out. | - |
| 10.2 | ZooKeeper expires nn1's ZKFC session; the lock znode disappears; nn2's ZKFC watch fires. | T_detect about 10 s |
| 10.3 | nn2's ZKFC creates the lock, reads the breadcrumb, finds nn1, starts fencing. | - |
| 40.4 | sshfence gives up after its 30-second default connect timeout; shell(/bin/true) succeeds. | T_fence about 30 s |
| 41.9 | nn2 finishes tailing edits and transitions to active. | T_promote 1.5 s |
| 43 to 46 | Clients' 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
maxClientCnxnslimits 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
- 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.
- Confirm
ha.zookeeper.session-timeout.msas actually configured, not as documented, on every ZKFC host. - Check the YARN state store: recovery enabled, ZKRMStateStore selected, a unique cluster id, and a split index sized for your application count.
- List every tenant of the ensemble and decide whether HDFS and YARN election should move to a dedicated one.
- Run the seven experiments on staging under load with the two measuring scripts, and record the timeline for each.
- Fix the largest budget term first (usually fencing), rerun the affected experiment, and repeat the game day after every Hadoop or ZooKeeper upgrade.