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.

Advertisement

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.

One ZooKeeper ensemble, many Hadoop clients, each holding a sessionZooKeeper ensemble (odd number of servers)leaderorders writesfollowervotes, servesfollowervotes, serveswrite = leader proposal + quorum ack + fsyncHDFS ZKFCNameNode electionYARN RMelection + state storeHDFS Routersstate store, tokensHiveServer2, HBasediscovery, master, metaEphemeral znode = livenessexists only while its session lives; lock, leader, memberSession timeout = failure detectortoo short: false failovers; too long: slow failoversZooKeeper decides who is in charge. It does not stop the old owner acting: that is fencing, done by each component.
Hadoop services each hold a ZooKeeper session and use ephemeral znodes for locks and membership. The session timeout is the failure detector, and fencing is the component's own job.

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.

ComponentWhat it keeps in ZooKeeperIf the quorum is lost
HDFS ZKFCEphemeral ActiveStandbyElectorLock and persistent ActiveBreadCrumb under /hadoop-ha/<nameservice>Active NameNode keeps serving; no automatic failover until the quorum returns
YARN ResourceManagerLeader election and, with ZKRMStateStore, application and token state for recoveryExpect YARN to stop accepting new work; running containers are not killed by the outage itself
HDFS Router-based federationOptionally the state store (mount table, membership) and router delegation tokensRouters serve cached state for a while, then become unsafe to use
HiveServer2Instance registrations for dynamic service discovery; optionally table locksNew JDBC connections cannot discover servers; existing sessions continue
HBaseActive master election, region server liveness, location of hbase:metaRegion 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.

Advertisement

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.

SettingWhereWhat it controls
tickTimezoo.cfgBase time unit; the server clamps client session timeouts to 2 to 20 ticks by default
minSessionTimeout, maxSessionTimeoutzoo.cfgOverride those clamp bounds explicitly
ha.zookeeper.session-timeout.mscore-site.xmlZKFC session timeout, which bounds how fast a dead NameNode is noticed
initLimit, syncLimitzoo.cfgTicks 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:3888

On 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.

SymptomLikely causeFix
Failovers with no NameNode or RM crashSession expiry from a long GC pause in ZKFC or RM, or slow ZooKeeper fsyncFix the pause; move ZooKeeper transaction logs to a dedicated disk; raise the session timeout modestly
Every client disconnects at onceEnsemble lost quorum, often two of three servers on the same failed host or rackSpread servers across failure domains; use five servers where two failures must be survivable
ZooKeeper disk fills upSnapshot and log purging never enabledSet autopurge.purgeInterval and autopurge.snapRetainCount, and alert on free space
Writes rejected, packet length errorsA znode or response larger than jute.maxbuffer, for example a huge RM state entryFind the oversized entry; raise jute.maxbuffer on servers and clients together only as a last resort
Failover controller cannot startParent znode missing or ACL mismatch after a rename or security changeRe-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

  1. Inventory every Hadoop service that points at your ensemble, with its configuration key, and check they all list the same hosts.
  2. Confirm the servers sit in distinct failure domains and that dataLogDir is on a dedicated disk.
  3. Enable autopurge and alert on free space in dataDir and dataLogDir.
  4. Read the effective session timeouts for ZKFC and the ResourceManager, and check them against tickTime clamps.
  5. Turn on GC logging for ZKFC and the ResourceManager, and alert on pauses approaching the session timeout.
  6. Scrape mntr or Prometheus metrics for latency, outstanding requests and fsync stalls, and graph leader changes.
  7. In a Kerberized cluster, verify SASL authentication and ACLs on /hadoop-ha and the RM state store root.
  8. Run a failover drill: kill the ZooKeeper leader, then a ZKFC, and time each recovery.
Key takeaway: ZooKeeper gives Hadoop one shared, correct answer to who is in charge: writes go through a quorum, sessions detect failure, and ephemeral znodes turn liveness into data. HDFS, YARN, routers, Hive and HBase build election, discovery and recovery state on top, while each still owns its own fencing. Run three or five servers across failure domains on dedicated log disks, tune session timeouts against real GC pauses, lock the tree down with SASL and ACLs, and watch latency and session expiries before they become false failovers.