Hadoop's master services each needed the same hard thing: agreeing, across machines that can crash or pause, on which instance is in charge. Instead of each project building its own consensus protocol, they all lean on Apache ZooKeeper. HDFS uses it to elect the active NameNode, YARN to elect the active ResourceManager and store its recovery state, and Hive, HBase and HDFS routers use it for discovery, locks and shared state.
That makes a small ensemble of three or five ZooKeeper servers one of the most critical pieces of a cluster, and one of the least understood. This article explains what ZooKeeper guarantees, who in Hadoop depends on it and for what, how sessions and timeouts turn into failover behaviour, and how to size, secure and monitor it. The NameNode failover sequence itself is covered in depth in HDFS high availability.
First principles: what ZooKeeper actually provides
ZooKeeper stores a small tree of nodes called znodes, each holding a few bytes to a few kilobytes of data, replicated on every server in the ensemble. One server is elected leader using the ZAB protocol. Every write goes to the leader, which proposes it to the followers, and commits it once a majority has written it to its transaction log. That is why ensembles have an odd number of servers: three tolerate one failure, five tolerate two, and four tolerate only one while adding write latency.
Reads are served locally by whichever server the client is connected to, so they are fast but can lag slightly behind the latest write. Each client sees updates in order, and a client that needs the latest value can issue a sync first. Writes are totally ordered. For coordination, three more features matter.
- Sessions. A client holds a session by sending heartbeats. If the ensemble hears nothing for the session timeout, the session expires.
- Ephemeral znodes. A znode created as ephemeral is deleted automatically when its owner's session expires. That turns liveness into data: if the node exists, its owner was recently alive.
- Watches. A client can ask to be notified once when a znode changes or is deleted, instead of polling.
Leader election falls out of these directly: every candidate tries to create the same ephemeral znode; exactly one succeeds; the others watch it and try again when it disappears.
Who in Hadoop uses ZooKeeper, and for what
It helps to know exactly which services break, and how, when ZooKeeper misbehaves. Most use it only for coordination, so an outage stops failover and discovery rather than stopping data processing, but the details differ.
| Component | What it keeps in ZooKeeper | If the quorum is lost |
|---|---|---|
| HDFS ZKFC | Ephemeral ActiveStandbyElectorLock and persistent ActiveBreadCrumb under /hadoop-ha/<nameservice> | Active NameNode keeps serving; no automatic failover until the quorum returns |
| YARN ResourceManager | Leader election and, with ZKRMStateStore, application and token state for recovery | Expect YARN to stop accepting new work; running containers are not killed by the outage itself |
| HDFS Router-based federation | Optionally the state store (mount table, membership) and router delegation tokens | Routers serve cached state for a while, then become unsafe to use |
| HiveServer2 | Instance registrations for dynamic service discovery; optionally table locks | New JDBC connections cannot discover servers; existing sessions continue |
| HBase | Active master election, region server liveness, location of hbase:meta | Region servers can abort after session expiry; the cluster degrades quickly |
Two lessons follow. First, ZooKeeper is on the control path, not the data path: HDFS reads and writes never touch it. Second, it is shared, so one misbehaving client can hurt all of them. A YARN state store writing large application records, or HBase with many region servers, generates far more load than a pair of failover controllers.
How election and fencing work together
The HDFS failover controller (ZKFC) runs next to each NameNode, health-checks it, and competes for the ephemeral lock znode on its behalf. The winner also writes a persistent breadcrumb recording which NameNode is active. When a new winner finds a breadcrumb left by someone else, it knows the previous active may still be running, for example paused in a long garbage collection, and fences it before promoting its own NameNode. The JournalNodes' epoch numbers add a second defence, since a stale NameNode's writes are rejected. NameNode HA architecture covers the full sequence.
The principle generalises to every Hadoop user of ZooKeeper: an ephemeral lock tells you who should be in charge, but it cannot stop a paused old owner from waking up and acting. Something else must: fencing, epochs, or ZooKeeper ACLs, which the ResourceManager state store uses so that only the current active can write. The same pattern is described generally in distributed leader election. Here is the core of it in Python with the kazoo client, useful for understanding and for small internal tools:
from kazoo.client import KazooClient
from kazoo.exceptions import NodeExistsError
LOCK = "/myservice/leader"
zk = KazooClient(hosts="zk1:2181,zk2:2181,zk3:2181", timeout=10.0)
zk.start()
zk.ensure_path("/myservice")
def try_lead(me):
try:
# Ephemeral: deleted by ZooKeeper if our session expires
zk.create(LOCK, me.encode(), ephemeral=True)
return True # we are leader; now fence the old one
except NodeExistsError:
# Watch fires once when the current leader's node goes away
zk.exists(LOCK, watch=lambda event: try_lead(me))
return False
Sessions and timeouts: the knob that decides failover behaviour
The session timeout is ZooKeeper's failure detector, and it is a direct trade-off. Short timeouts detect a dead NameNode or ResourceManager quickly, but any pause longer than the timeout looks like death: a stop-the-world garbage collection in the ZKFC or RM process, a stalled virtual machine, or a ZooKeeper server blocked on a slow disk. Long timeouts avoid false failovers but leave the cluster without an active master for longer after a real crash.
| Setting | Where | What it controls |
|---|---|---|
| tickTime | zoo.cfg | Base time unit; the server clamps client session timeouts to 2 to 20 ticks by default |
| minSessionTimeout, maxSessionTimeout | zoo.cfg | Override those clamp bounds explicitly |
| ha.zookeeper.session-timeout.ms | core-site.xml | ZKFC session timeout, which bounds how fast a dead NameNode is noticed |
| initLimit, syncLimit | zoo.cfg | Ticks a follower may take to sync with, or lag behind, the leader |
The clamp is a common surprise: a client that requests a session timeout above twenty ticks silently gets twenty ticks. With a tickTime of 2,000 milliseconds, requests are clamped between 4 and 40 seconds. Read the defaults for the Hadoop timeout from core-default.xml for your release rather than from memory, since they have changed over time.
Worked example. A ZKFC with a ten-second session timeout suffers a twelve-second GC pause. At ten seconds its session expires and ZooKeeper deletes its lock. The standby ZKFC's watch fires, it takes the lock, sees the breadcrumb, fences the old NameNode, and promotes its own. The old NameNode was healthy the whole time; the outage was caused by the failover controller's heap. The fix is a small, fixed heap for ZKFC, GC logging on it, and alerts on session expiries, not a longer timeout.
Sizing and placing the ensemble
- Three or five servers. Three for most clusters; five when you must survive two simultaneous failures or maintenance plus a failure. More voters make writes slower. Add observers, which serve reads but do not vote, if read load grows.
- Failure domains. Put each server on a different host and, ideally, a different rack or zone. Two servers on one hypervisor is one failure away from losing quorum.
- Dedicated disk for the transaction log. Every write waits for an fsync. Point dataLogDir at its own device, away from DataNode and YARN spill disks, or write latency tracks your busiest job.
- Small, stable heap. ZooKeeper keeps the whole tree in memory; a few gigabytes is plenty for typical Hadoop use. Avoid swap entirely.
- Purge old data. Snapshots and logs accumulate forever unless autopurge is enabled.
- Shared or dedicated. A heavy HBase or Kafka deployment deserves its own ensemble so its load cannot starve NameNode failover.
# zoo.cfg for a three-server ensemble
tickTime=2000
initLimit=10
syncLimit=5
dataDir=/data/zookeeper
dataLogDir=/zklog/zookeeper
clientPort=2181
autopurge.snapRetainCount=5
autopurge.purgeInterval=24
4lw.commands.whitelist=mntr,ruok,srvr
server.1=zk1.example.com:2888:3888
server.2=zk2.example.com:2888:3888
server.3=zk3.example.com:2888:3888On the Hadoop side, the quorum string appears in a few places: ha.zookeeper.quorum in core-site.xml for ZKFC, and hadoop.zk.address for the ResourceManager's embedded election and state store, replacing the deprecated yarn.resourcemanager.zk-address. Keep one source of truth for the host list in your configuration management, because a stale list in one file only shows up when that component fails over.
Security
In a Kerberized cluster, clients authenticate to ZooKeeper using SASL with their service principals, and the znodes each service creates should carry ACLs that restrict writes to that principal. ZKFC takes its ACLs and optional digest credentials from ha.zookeeper.acl and ha.zookeeper.auth. Without ACLs, any user who can reach port 2181 can delete the failover lock or corrupt the ResourceManager state store. The Kerberos setup itself is covered in Hadoop Kerberos. Also restrict the four-letter admin commands to a short whitelist, as above, and remember that ZooKeeper 3.5 and later start an embedded admin HTTP server, on port 8080 by default, which often clashes with other daemons; set admin.serverPort or disable it deliberately.
Monitoring and failure modes
The mntr command, or the Prometheus metrics provider in newer releases, exposes what matters. Alert on zk_avg_latency and zk_max_latency for slow disks, zk_outstanding_requests for overload, zk_fsync_threshold_exceed_count for fsync stalls, zk_num_alive_connections for connection storms, and on which server is leader, so leader changes are visible. On the client side, alert on session expiries and failover events in ZKFC and ResourceManager logs.
| Symptom | Likely cause | Fix |
|---|---|---|
| Failovers with no NameNode or RM crash | Session expiry from a long GC pause in ZKFC or RM, or slow ZooKeeper fsync | Fix the pause; move ZooKeeper transaction logs to a dedicated disk; raise the session timeout modestly |
| Every client disconnects at once | Ensemble lost quorum, often two of three servers on the same failed host or rack | Spread servers across failure domains; use five servers where two failures must be survivable |
| ZooKeeper disk fills up | Snapshot and log purging never enabled | Set autopurge.purgeInterval and autopurge.snapRetainCount, and alert on free space |
| Writes rejected, packet length errors | A znode or response larger than jute.maxbuffer, for example a huge RM state entry | Find the oversized entry; raise jute.maxbuffer on servers and clients together only as a last resort |
| Failover controller cannot start | Parent znode missing or ACL mismatch after a rename or security change | Re-run hdfs zkfc -formatZK with NameNodes stopped, and check ha.zookeeper.acl |
When debugging, look at the tree itself. zkCli.sh -server zk1:2181, then ls /hadoop-ha and stat on the lock znode, shows which session owns the lock and since when; comparing that owner with the NameNode that clients think is active resolves most confusion quickly.
What to do next
- Inventory every Hadoop service that points at your ensemble, with its configuration key, and check they all list the same hosts.
- Confirm the servers sit in distinct failure domains and that dataLogDir is on a dedicated disk.
- Enable autopurge and alert on free space in dataDir and dataLogDir.
- Read the effective session timeouts for ZKFC and the ResourceManager, and check them against tickTime clamps.
- Turn on GC logging for ZKFC and the ResourceManager, and alert on pauses approaching the session timeout.
- Scrape mntr or Prometheus metrics for latency, outstanding requests and fsync stalls, and graph leader changes.
- In a Kerberized cluster, verify SASL authentication and ACLs on /hadoop-ha and the RM state store root.
- Run a failover drill: kill the ZooKeeper leader, then a ZKFC, and time each recovery.