skip to content

How does ZooKeeper's Zab protocol combine with ephemeral sequential znodes to implement leader election, how does this design avoid a 'thundering herd' of watch notifications when the leader changes, and how does quorum-based membership prevent split-brain under a network partition?

level: principalimportance: should knowfreq 35%

answer

  1. Zab elects ZooKeeper's own internal leader via majority quorum
  2. minority partition can't commit writes -> split-brain avoided
  3. ephemeral znode dies with the session
  4. sequential znode = lowest number is leader
  5. watch only your immediate predecessor to avoid herd effect

basics

~20 s

ZooKeeper clients each create a temporary numbered file to compete for leadership. Lowest number wins; others watch only the entry just ahead of them, so a change wakes one client, not everyone. It needs more than half its servers reachable.

solid answer

~50 s

Zab keeps ZooKeeper's own servers in a consistent, ordered replicated log with a single active internal leader, elected only when a majority (quorum) of servers can communicate; a minority partition simply cannot process writes, which is ZooKeeper's core split-brain defense. Application-level leader election is a pattern built on top, not a built-in primitive: each candidate creates an ephemeral, sequential znode under a shared parent path; ephemeral means it's auto-deleted if that client's session dies, sequential means ZooKeeper appends a monotonically increasing suffix. The client holding the lowest-numbered znode is leader. Every other client, instead of watching the whole list, sets a watch only on the znode immediately preceding its own; when that one gets deleted, exactly one waiter is notified and re-checks, avoiding a 'thundering herd' where every waiting client wakes on every single leadership change.

go deeper

for a junior

Should grasp the basic idea: clients create numbered temporary entries, and the lowest number is leader; entries disappear if a client dies.

for a middle

Should distinguish ephemeral (dies with session) from sequential (ordered numbering), and describe checking the children list to find the lowest.

for a senior

Should explain the predecessor-only watch pattern and why it avoids a herd effect, and connect quorum loss in the ZooKeeper ensemble itself to write unavailability rather than incorrect results.

for a principal

Should distinguish Zab's internal ensemble leadership from the application-level election recipe, discuss the session-timeout-vs-pause zombie-leader risk and how znode versioning can serve as a fencing token, and cite real systems (Kafka's pre-KRaft controller election, Curator LeaderLatch) that use this exact pattern.

## Zab and the majority quorum ZooKeeper's own internal high availability rests on **Zab** (ZooKeeper Atomic Broadcast), a crash-recovery, primary-backup replication protocol conceptually similar in spirit to Raft: the ZooKeeper ensemble (typically 3, 5, or 7 servers, always an odd count) elects one of its own members as the internal leader responsible for ordering and broadcasting all state-mutating operations to followers, and that internal leader is only valid, and can only commit writes, while it has active, healthy connections to a majority (a quorum) of the ensemble's servers. If the ensemble network partitions such that no side has a majority, or the minority side is cut off from the majority, the minority side cannot elect or sustain a leader capable of committing writes at all - it simply stops serving writes (and, depending on configuration, can also stop serving reads to avoid returning stale data), while the majority side continues operating normally. This majority-quorum requirement is precisely what prevents ZooKeeper itself from split-brain: by construction, any two possible majorities out of the same fixed ensemble size must overlap in at least one server, so at most one side of any partition can ever hold a quorum at a time, and that overlap is what makes 'exactly one active leader ensemble-wide' provable rather than merely likely. ## The two primitives the recipe is built from Application-level leader election, the pattern most engineers actually build on top of ZooKeeper (and its close relative etcd, though etcd typically uses lease-plus-Raft directly), is not a separate protocol inside ZooKeeper - it's a recipe built from two primitives ZooKeeper exposes: ephemeral znodes and sequential znode naming. A **znode** is roughly a file-like node in ZooKeeper's hierarchical namespace. - **Ephemeral** means the znode is automatically deleted the moment the client session that created it ends, whether by explicit close, crash, or session timeout from a lost heartbeat - this is the mechanism that ties 'am I still alive' directly to the coordination data itself, no separate heartbeat-and-cleanup logic required. - **Sequential** means ZooKeeper appends a monotonically increasing, ensemble-wide-unique numeric suffix to the requested znode name at creation time, guaranteeing a total, gap-free creation order even under concurrent creation by many clients. ## Running the election, one watch at a time To run an election, each candidate process creates an ephemeral sequential znode under a shared parent path, for example `/election/candidate-`, receiving back something like `/election/candidate-0000000042`. 1. The candidate then lists the parent path's children and checks whether its own znode has the lowest sequence number among them; if so, it is the leader. 2. If not, rather than watching every znode ahead of it (which would mean every waiting candidate gets notified on every single change to the list, an `O(n)` notification fan-out on every leadership transition - the 'thundering herd' or 'herd effect'), each candidate sets a ZooKeeper watch on exactly one znode: the one with the sequence number immediately preceding its own. 3. When that specific znode is deleted (its owner either resigned or its session died), only the one candidate watching it receives a notification; that candidate then re-lists the children (or simply checks if it's now the lowest) and, if not yet lowest, sets a new watch on whatever now immediately precedes it. This chained, one-predecessor-only watch design means a single leadership change triggers exactly one notification and one re-check, not a stampede of every waiting node re-reading the full candidate list simultaneously, which matters a great deal at scale since ZooKeeper watches are one-shot and re-arming many of them under load is itself expensive. ## The subtlety that still bites The combination of these two layers - Zab's majority-quorum internal leadership for ZooKeeper's own consistency, and the ephemeral-sequential-znode recipe for application leader election - is what makes ZooKeeper-based leader election meaningfully safer against split-brain than a naive Bully or lease-only scheme, provided the application respects ZooKeeper's session-timeout semantics correctly. The subtlety production teams get bitten by is **session timeout versus GC pauses**: if an application-level leader experiences a long stop-the-world pause longer than its ZooKeeper session timeout, the session expires, ZooKeeper deletes its ephemeral znode, and a new leader is legitimately elected - but the original process, once it resumes, may not immediately realize its session died and could still attempt leader-only work before noticing the session-expired event, which is the same zombie-leader class of bug seen with lease-based leadership; the standard mitigation is, again, to pair election with fencing tokens (ZooKeeper's own znode version/czxid can serve this role) enforced at the point where the leader actually touches shared state, rather than trusting 'I still think I'm the leader' locally. ## Real deployments Real deployments that popularized this exact pattern include: - **Apache Kafka's original controller election** (pre-KRaft, Kafka used ZooKeeper ephemeral znodes specifically for controller leader election); - **HBase's master election**; - **Apache Curator's LeaderLatch/LeaderSelector recipes**, which implement precisely the ephemeral-sequential-znode-plus-predecessor-watch pattern described above as a reusable library rather than something every team hand-rolls.

  • Why watch only the immediate predecessor znode instead of watching the current leader's znode directly?
    Watching the leader directly would mean every waiting candidate gets notified the instant the leader changes, causing every one of them to simultaneously re-list and re-evaluate the full candidate set on every single transition - an O(n) notification storm on each change. Watching only your immediate predecessor means each transition wakes exactly one candidate, who then re-checks and potentially sets a new watch, keeping notification cost constant per transition regardless of how many candidates are waiting.
  • What happens to application-level leadership if the ZooKeeper ensemble itself loses quorum, separately from any single client's session?
    If the ensemble can't maintain a quorum, it stops accepting writes (and typically also serving reads, depending on configuration) entirely, so no ephemeral znode can be created, deleted, or reliably watched during that window - application-level leader election effectively pauses along with all other ZooKeeper-dependent coordination until quorum is restored, rather than producing an incorrect result.
  • How does a znode's version number (czxid) serve as a fencing token in this pattern?
    Each znode carries a globally unique, monotonically increasing transaction ID assigned by Zab at creation, so an application can have the leader attach its own leader-znode's czxid to every write against shared storage; the storage layer then rejects any write bearing an older czxid than one it has already accepted, closing the same zombie-leader gap that lease-based fencing tokens address, without needing a separate token-issuing system.

It's like a numbered deli-counter ticket system with a twist: your ticket vanishes automatically the moment you leave the shop (ephemeral), the lowest ticket number currently in the machine gets served next (sequential), and instead of everyone in line staring at the display waiting for their number, each person only watches the one ticket directly ahead of theirs - so when that one person gets served and steps away, only the very next person in line is tapped on the shoulder, not the whole room.

saying these in an interview costs you the question

  • Conflates Zab (ZooKeeper's internal replication/leader protocol) with the application-level ephemeral-znode election recipe as if they're the same mechanism
  • Thinks watching every candidate's znode directly is how the herd effect is avoided, rather than watching only the immediate predecessor
  • Believes ZooKeeper can keep serving writes during a minority-side partition
  • Doesn't mention that ephemeral znode deletion is tied to session expiry, not just process crash
  • Assumes ZooKeeper-based election is immune to zombie-leader/pause-related split-brain without fencing

context