skip to content

Walk through how the Kafka controller was elected via ZooKeeper, and what happened on controller failover.

level: seniorimportance: should knowfreq 50%

answer

  1. race to create ephemeral /controller
  2. losers watch /controller
  3. /controller_epoch = fencing token
  4. failover = reload all metadata from ZK
  5. O(partitions) recovery → KRaft fix

basics

~20 s

Brokers raced to create a single ephemeral znode at /controller; whoever created it first became the controller. If that broker died, its session expired, /controller was deleted, every other broker got a watch notification, and they raced again to elect a new controller.

solid answer

~40 s

Controller election was a classic ZooKeeper leader-election: all brokers attempted to create the **ephemeral** znode `/controller`. ZK guarantees only one create succeeds, so exactly one broker became the controller; the rest set a **watch** on `/controller`. The winner also bumped `/controller_epoch` (a fencing token). The controller is the broker responsible for partition leader assignment, ISR management, and reacting to broker membership changes. On failover, the dead controller's ZK session expired, ZK deleted the ephemeral `/controller`, and the watch fired on every other broker; they raced to recreate it. The new controller incremented `controller_epoch`, then **reloaded the full cluster metadata from ZK** (topics, partitions, ISR, broker list) to rebuild its in-memory state before resuming. This metadata reload made failover slow at high partition counts — a key pain point KIP-500/KRaft addressed.

go deeper

for a junior

Know one broker is the controller and it's chosen by creating the /controller znode.

for a middle

Explain the race-to-create + watch pattern and that failover re-runs the election.

for a senior

Add controller_epoch fencing and the metadata-reload cost as the real failover bottleneck.

for a principal

Contrast with KRaft's Raft-log controller where metadata is already in memory, and reason about fencing/split-brain guarantees end to end.

## What the controller is In a ZK-based Kafka cluster, exactly one broker acts as the **controller**. It's a normal broker that additionally owns cluster-wide control decisions: assigning partition leaders, managing the **ISR** (in-sync replica) sets, handling partition reassignment, and reacting to brokers joining/leaving. Other brokers are followers of these decisions. ## The election mechanism 1. **Race to create `/controller`** — every broker tries to create the ephemeral znode `/controller`. ZooKeeper's create is atomic and rejects a second create on an existing node, so exactly one broker wins. The winner writes its broker id into the znode. 2. **Losers set a watch** — the brokers that lost register a one-shot **watch** on `/controller` so they'll be notified when it disappears. 3. **Epoch bump** — the new controller increments the persistent `/controller_epoch` znode. This epoch is a **fencing token**: requests carry the epoch, and brokers reject any command stamped with a stale (older) epoch. This prevents a 'zombie' old controller from issuing conflicting commands after a partition. ## Controller failover 1. The controller crashes or is partitioned. Its ZK **session** eventually expires (bounded by the session timeout). 2. ZK auto-deletes the ephemeral `/controller` (it was owned by that session). 3. The deletion fires the watch on all other brokers. 4. They race to recreate `/controller`; one wins and becomes the new controller, incrementing `controller_epoch` again. 5. **State rebuild** — the new controller must reconstruct its in-memory view by reading metadata from ZK: the full list of topics, partitions, replica assignments, leaders, and ISR. With tens of thousands of partitions, this read storm took a long time, during which leadership changes stalled. ## Edge cases / gotchas - **Split-brain via stale controller**: a slow/partitioned old controller might still think it's in charge. The `controller_epoch` fencing token is what makes its commands get rejected once a newer controller exists. - **Watch is one-shot**: after firing, brokers re-register the watch. - **Failover latency**: the dominant cost was the metadata reload from ZK, not the election itself. This O(partitions) recovery time was a core scalability limit. - **Single-writer**: only the controller mutated leadership; this serialization simplified correctness but funneled all metadata changes through one node plus ZK. ## Why this motivated KRaft Under KIP-500/KRaft, metadata is itself a replicated log managed by a Raft quorum of controllers, so the active controller already has metadata in memory (it's the tail of the log it replicates). Failover no longer requires re-reading all metadata from an external store, collapsing failover time from seconds/minutes to roughly the time to elect a new Raft leader.

  • What is the role of /controller_epoch?
    It's a monotonically increasing fencing token bumped on each new controller. Brokers reject commands stamped with an older epoch, preventing a partitioned/zombie old controller from issuing conflicting leadership commands (split-brain protection).
  • Why was ZK-based controller failover slow at large partition counts?
    The newly elected controller had to reload the entire cluster metadata (topics, partitions, replicas, leaders, ISR) from ZooKeeper to rebuild its in-memory state. That read cost scaled with the number of partitions, so recovery time grew with cluster size.

saying these in an interview costs you the question

  • Claiming multiple controllers can be active at once — exactly one wins /controller; epoch fencing handles stragglers.
  • Saying election is slow — the election is cheap; the metadata reload is what was slow.
  • Forgetting controller_epoch / fencing, leaving no answer for split-brain.

context