ZooKeeper is a small, replicated, strongly ordered tree that distributed systems use to agree on things: who is the leader, which servers are alive, what the current configuration is, and who holds a lock. It is not a data store and it is not fast at writes. It offers a narrow set of guarantees that are hard to build yourself, through an API simple enough that coordination patterns are assembled from it on the client side.
This page explains ZooKeeper as a system: the data model, how an ensemble orders writes with the Zab protocol, sessions and ephemeral nodes, watches, what the consistency guarantees really promise (reads can be stale), and how to operate an ensemble. The Hadoop-specific view of which components store what is covered separately in ZooKeeper in Hadoop.
What ZooKeeper is for, and what it is not
ZooKeeper solves coordination: a group of processes needs one agreed answer to a small question, and the answer must survive the crash of any minority of servers. Leader election, membership, service discovery, configuration and locks all need the same primitives: a history of changes every server applies in the same order, a way to tie state to a client's liveness, and a way to be told when something changes.
It is not a database. The whole tree lives in memory on every server, each znode is limited by jute.maxbuffer (just under 1 MB by default; stay far below it), and every write goes through one leader and a disk flush on a majority. That is plenty for coordination and far too little for a data path. Store the pointer, not the payload: "partition 7 is owned by broker 3, epoch 12" goes in ZooKeeper, the partition goes elsewhere.
The data model: znodes and their types
The namespace is a tree of znodes addressed by paths such as /services/billing/leader. Every znode holds a small byte array and may have children. Its Stat carries what coordination code relies on: czxid and mzxid (the transactions that created and last modified it), version, cversion (children version) and ephemeralOwner (the owning session, or zero).
Versions make optimistic concurrency cheap: setData(path, data, expectedVersion) fails with BadVersion if someone changed the node first, so a read-modify-write loop needs no lock. multi() applies a list of create, delete, setData and check operations atomically.
| Node type | Lifetime | Typical use |
|---|---|---|
| Persistent | until deleted | configuration, directory nodes |
| Ephemeral | until the creating session ends | liveness, membership, lock ownership |
| Sequential (either kind) | as above; the server appends a 10-digit, monotonically increasing counter | queues, fair locks, election order |
| Container | deleted by the server once its last child is removed | parent directories of lock and election nodes |
| TTL (persistent only) | deleted if unmodified and childless for the TTL; must be enabled with zookeeper.extendedTypesEnabled | leases that outlive a session |
Ephemeral nodes cannot have children, and sequence counters are per parent.
The ensemble and the request pipeline
Every server holds a full replica of the tree. One is elected leader; the other voters are followers. Optional observers receive commits but do not vote, adding read and connection capacity without slowing writes; they are the usual way to serve a second data center.
A client keeps a TCP session to one server, which answers reads from memory. Writes go to the leader, which turns each into an idempotent transaction (a conditional write becomes a concrete value and version, so replay cannot diverge), assigns it a zxid and sends a PROPOSAL. Followers log and fsync it before acknowledging; once a majority of voters, leader included, has acknowledged, the leader sends COMMIT to followers and INFORM to observers. Servers apply commits in zxid order and then reply to waiting clients.
Zab: how the order survives a leader crash
Zab (ZooKeeper Atomic Broadcast) is the protocol behind that pipeline. Its key object is the zxid, a 64-bit number: the high 32 bits are the epoch, incremented each time a new leader takes over, and the low 32 bits are a counter within the epoch. Comparing zxids therefore compares both leadership era and position in history.
Zab runs in phases. In election, servers vote for the candidate with the most up-to-date history (highest zxid, ties broken by server id). In synchronization, the new leader establishes a new epoch with a quorum and brings followers up to date with missing transactions, a truncation of uncommitted proposals, or a full snapshot. Only then does it broadcast new writes. The result: anything committed by one leader is delivered by every later leader, an uncommitted proposal from a crashed leader is committed everywhere or nowhere, and every server delivers the same order.
Zab resembles Raft but is built for primary-backup replication: the primary sends state changes, not client commands. See Raft consensus and total order broadcast.
Sessions, heartbeats and ephemeral nodes
A session is ZooKeeper's notion of a live client. The server clamps the requested timeout between minSessionTimeout and maxSessionTimeout, by default 2 and 20 times tickTime. If the client's server dies, the library reconnects elsewhere with the same session id; within the timeout, the session, its ephemeral nodes and its watches survive.
Expiry is decided by the ensemble. When the leader hears nothing for the full timeout, it expires the session as a transaction, deleting its ephemeral nodes and firing their watches. That is liveness detection: a dead service's ephemeral node disappears within about one timeout.
The subtle part is on the client: while disconnected it cannot know whether its session has expired, and a leader that keeps acting may already have been replaced. Stop leader-only work as soon as the state becomes Disconnected (Curator calls it SUSPENDED), resume only on reconnect with the same session, and treat Expired as final. Even then a paused process can wake and act on stale beliefs, so resources the leader writes to should check a fencing token such as the zxid that created the leader's znode.
Watches: one-shot, persistent, and the herd effect
A read can leave a watch: getData and exists watch data, creation and deletion; getChildren watches the child list. The classic watch is one-shot and carries only the event type and path, so the client must re-read and re-register; intermediate changes are not delivered, only the latest state. A client always sees the watch event before the new data that triggered it.
ZooKeeper 3.6 added addWatch with persistent and persistent-recursive modes, which stay registered and can cover a subtree, removing the re-registration round trip for caches at the cost of more events.
The herd effect is the classic mistake: a thousand clients watch one lock node, it is released, and all wake at once. The recipes make each waiter watch only its predecessor.
What is guaranteed, and when reads are stale
Writes are linearizable: one global order that respects real time. Each client's requests run in the order sent, and a client never sees older state after newer state, even after reconnecting, because a server behind the client's last seen zxid refuses the connection.
Reads are not guaranteed current. They are served locally, and a follower may lag, so if A writes and tells B out of band, B can still read the old value. When a read must reflect everything committed before it, call sync(path) first; it returns once the connected server has caught up with the leader.
Recipe: leader election that scales
Every candidate creates an ephemeral sequential node; the lowest is the leader; every other candidate watches the node just before its own, so a leader crash wakes one client. The code uses the plain Java client; in production prefer Apache Curator's LeaderLatch or LeaderSelector, which implement the same recipe with careful connection-state handling.
// Leader election: one EPHEMERAL_SEQUENTIAL znode per candidate, each watching its predecessor.
public final class Candidate implements Watcher {
private final ZooKeeper zk;
private final String dir = "/services/billing/election";
private String me; // e.g. /services/billing/election/n_0000000042
public Candidate(ZooKeeper zk) { this.zk = zk; }
public void enter() throws KeeperException, InterruptedException {
me = zk.create(dir + "/n_", hostId(), ZooDefs.Ids.OPEN_ACL_UNSAFE,
CreateMode.EPHEMERAL_SEQUENTIAL);
check();
}
private void check() throws KeeperException, InterruptedException {
while (true) {
List<String> kids = zk.getChildren(dir, false);
Collections.sort(kids); // fixed-width suffix sorts correctly
String mine = me.substring(dir.length() + 1);
int idx = kids.indexOf(mine);
if (idx < 0) throw new IllegalStateException("my node vanished: session expired");
if (idx == 0) { becomeLeader(zk.exists(me, false).getCzxid()); return; }
String prev = dir + "/" + kids.get(idx - 1);
if (zk.exists(prev, this) != null) return; // wait for NodeDeleted on predecessor
// predecessor disappeared between getChildren and exists: loop and re-check
}
}
@Override public void process(WatchedEvent e) {
if (e.getType() == Event.EventType.NodeDeleted) {
try { check(); } catch (Exception ex) { stepDown(); }
} else if (e.getState() == Event.KeeperState.Disconnected) {
stepDown(); // cannot prove we still hold leadership
}
}
}Three details matter: re-check after setting the watch, because the predecessor may vanish in between; pass the node's czxid into leader work as a fencing token, since it increases with each new leader; and step down on Disconnected, not only Expired. The same pattern gives fair locks; see leader election for the problem in general.
Worked example: sizing a five-node ensemble
A platform runs 400 service instances with ephemeral registrations, 50 leader-elected controllers each updating a small status node once a second, and a configuration tree every instance reads: about 60 writes per second, bursts of 400 when the fleet restarts, and a few thousand reads per second.
Quorum first. A write needs a majority of voters: three tolerate one failure, five tolerate two, and four still tolerate only one, so even counts buy nothing. Five lets one server be down for maintenance while surviving an unplanned failure, which is why it is the usual size. Seven makes every write wait for four fsyncs; add observers for read scale instead.
Latency next. A write costs the leader's fsync, one round trip and the fastest majority of follower fsyncs. With a dedicated SSD log at about 1 ms per fsync and 0.5 ms round trips, a write commits in roughly 2 to 3 ms. Share the log disk with snapshots or application I/O and 50 ms fsyncs appear; followers exceed syncLimit, get dropped and resynchronize. The write rate here is easy; what hurts is 400 new sessions at once, each a transaction, and the watch storm when 400 ephemeral nodes vanish together. Finally, size the JVM heap for the tree plus snapshot serialization and never let the process swap: a swapping server misses heartbeats and expires sessions fleet-wide.
Operating an ensemble
# zoo.cfg for one voter of a five-voter ensemble
tickTime=2000
initLimit=10 # ticks a follower may take to connect and sync with the leader
syncLimit=5 # ticks a follower may lag before the leader drops it
dataDir=/var/lib/zookeeper/data # snapshots + myid
dataLogDir=/zk-txlog # transaction log on its own low-latency device
clientPort=2181
autopurge.snapRetainCount=5
autopurge.purgeInterval=24 # hours; 0 disables purging and the disk fills
4lw.commands.whitelist=ruok,mntr,srvr,stat
server.1=zk1:2888:3888
server.2=zk2:2888:3888
server.3=zk3:2888:3888
server.4=zk4:2888:3888
server.5=zk5:2888:3888
server.6=zk6:2888:3888:observer # on zk6 itself also set peerType=observer- Disks and purge. Put
dataLogDiron a dedicated low-latency device, and setautopurge.purgeIntervalor snapshots and logs accumulate forever. - Monitoring.
mntr(four-letter commands must be whitelisted since 3.5) and the AdminServer expose latency, outstanding requests, znode, watch and ephemeral counts. Alert on leader changes, rising outstanding requests and fsync warnings. - Changes. Use dynamic reconfiguration (3.5+) or rolling restarts that keep a majority; upgrade followers first, the leader last.
- Security. SASL and TLS on client and quorum ports, ACLs on sensitive paths;
OPEN_ACL_UNSAFEabove is illustration only.
Failure modes
- Split leadership in the client. A process keeps acting as leader after Disconnected because it only listens for Expired. Fence every side effect.
- Slow fsync. Shared or throttled disks stall the commit path, followers drop, elections churn. Watch fsync warnings.
- Session storms. A network blip longer than the session timeout expires thousands of sessions; every client reconnects and re-registers at once. Use jittered backoff and sensible timeouts.
- Oversized znodes and huge child lists. Large payloads or directories with hundreds of thousands of children slow snapshots and
getChildren, and can exceedjute.maxbufferon the response. - Stale reads mistaken for bugs. A reader on a lagging follower misses a write; use
syncwhen it matters. - Even-sized or cross-region voting ensembles. Four voters, or voters spread across two sites, lose quorum when one site fails. Use three sites or observers.
Trade-offs: ZooKeeper versus Raft-based stores
etcd and Consul offer the same core service with Raft, a flat key space, leases instead of sessions, and linearizable reads by default. Systems are also absorbing coordination: Kafka's KRaft mode replaces ZooKeeper with an internal Raft quorum. ZooKeeper remains sound where it is already embedded (HBase, Hadoop HA, Solr) and where Curator's recipes are wanted. For a new system, prefer the store your platform already runs, and build on leases and fencing either way; see leases.
What to do next
- Draw your znode tree, mark each node's type, and move any payload that belongs elsewhere.
- Make every leader action stop on Disconnected and carry a fencing token such as the czxid.
- Replace hand-written recipes with Apache Curator.
- Add
syncbefore reads that must observe another client's write. - Run five voters in three failure domains, the transaction log on a dedicated device, autopurge on, and alerts on fsync warnings and leader changes.
- Test a leader kill, a follower disk stall and a network partition in staging, and measure session expiry time against your timeout.