skip to content

What are the consistency and availability tradeoffs of routing all metadata writes through a single active controller and replicated log, compared to the old ZooKeeper model?

level: principalimportance: should knowfreq 30%

answer

  1. Single writer = one total order
  2. commit on quorum majority (Raft)
  3. lose majority -> writes stall, reads OK
  4. fast failover, no ZK reload
  5. scales to millions of partitions

basics

~20 s

The single-writer log gives one ordered history of all metadata changes, so every node converges to the same state and failover is fast. The cost: writes need a controller-quorum majority, and the active controller is a serialization point for metadata throughput.

solid answer

~50 s

Funneling all metadata mutations through one active controller appending to a Raft-replicated log buys strong properties: a single total order of events, deterministic replay (every node converges), commit-on-majority durability, and fast failover because followers already hold the log. Versus ZooKeeper, it removes a separate system, eliminates the controller↔ZK round-trips that bottlenecked large clusters, and scales to far more partitions because metadata propagates as an incremental log replayed locally rather than via ZK watches. Tradeoffs: writes require a *majority* of the controller quorum to be up (lose 2 of 3 controllers and writes stall), and the active controller is a single serialization point — extreme metadata churn is bounded by one node's append/commit rate. You mitigate by sizing the quorum (3 or 5), isolating controllers, using incremental records and snapshots, and accepting that reads stay highly available (local, lock-free) even while writes are blocked.

go deeper

for a junior

Know the single controller writes metadata and a quorum keeps it durable.

for a middle

Contrast with ZooKeeper at a high level: no external system, faster failover.

for a senior

Explain commit-on-majority, the write-availability vs read-availability split, and quorum sizing.

for a principal

Frame the CAP-style tradeoff for metadata, scaling rationale (incremental log vs ZK watches), serialization-point limits, and operational quorum/topology decisions.

**The old ZooKeeper model.** Pre-KRaft, Kafka kept metadata in ZooKeeper (ZK), an external quorum-based coordination service. One broker was the *controller*; it watched ZK znodes and pushed metadata updates to other brokers via `UpdateMetadata`/`LeaderAndIsr` RPCs. Problems at scale: the controller had to read large state from ZK on failover (slow controller failover), ZK watch fan-out and the controller's RPC fan-out limited how many partitions a cluster could support, and you operated *two* distributed systems (Kafka + ZK) with their own tuning and failure modes. There was also a risk of divergence between ZK state and controller in-memory state. **The KRaft single-writer model.** KRaft replaces this with an internal Raft log (`__cluster_metadata`) written by exactly one *active controller*. Properties gained: - **Single total order.** One writer + one log ⇒ a canonical sequence of events. No reconciling concurrent writers. - **Deterministic convergence.** Brokers replay the same ordered, committed records and reach identical state — the event-sourced state machine. - **Commit-on-majority durability (Raft).** A record is committed once a *majority* of the controller quorum has persisted it; committed records survive any minority failure. - **Fast failover.** Followers already replicate the log, so a new active controller doesn't reload everything from an external store; it just continues from its log. Controller failover drops from seconds/minutes to typically sub-second to low seconds. - **Scale.** Metadata is propagated as an incremental, replayable log (plus snapshots) instead of ZK watches and full-state RPCs, enabling clusters with millions of partitions (a stated KIP-500 goal). **The tradeoffs / costs.** - **Write availability needs a quorum majority.** With a 3-node controller quorum you tolerate 1 failure; lose 2 and metadata *writes* halt (you cannot create topics, elect new leaders, etc.) until quorum is restored. A 5-node quorum tolerates 2 failures at the cost of higher commit latency (more nodes must ack). - **Serialization point.** All metadata mutations serialize through one node's append+commit path. For most workloads this is ample, but pathological metadata churn (massive simultaneous reassignments) is bounded by that single pipeline. Incremental records (`PartitionChangeRecord`) and batching mitigate this. - **Reads remain available.** Crucially, brokers serving produce/fetch read metadata *locally* from their in-memory image; a controller-write outage does not stop data-plane reads/writes for existing leaders — it stops *changes* to metadata. This is a deliberate CAP-style choice: favor consistency for metadata writes, keep reads available. **Operational guidance / mitigations:** run a dedicated, odd-sized controller quorum (3 or 5) on isolated, fast-disk nodes; separate `controller` and `broker` roles in large deployments; rely on snapshots to bound recovery; and monitor quorum health and metadata lag. Compared with ZK, you now operate one system, with a clearer failure model (Raft majority) and far better controller-failover and partition-scale characteristics, accepting the inherent single-writer serialization and majority-availability constraints.

  • If 2 of 3 controllers are down, what still works and what stops?
    Metadata writes stop (no quorum majority) — no topic creation, no new leader elections. But existing brokers keep serving produce/fetch for current leaders because they read metadata locally; the data plane is unaffected until a change is needed.
  • Why does KRaft scale to far more partitions than the ZooKeeper model?
    Metadata propagates as an incremental, replayable log with snapshots that brokers apply locally, instead of ZooKeeper watches plus full-state controller RPCs, which were the per-partition bottlenecks.
  • Why choose a 5-node controller quorum over 3?
    It tolerates 2 simultaneous controller failures instead of 1, at the cost of higher write/commit latency since a majority (3) must ack each record.

saying these in an interview costs you the question

  • Saying losing the active controller causes data loss (followers hold the committed log)
  • Claiming metadata reads stop when the controller quorum loses majority
  • Asserting KRaft has no availability tradeoff at all
  • Saying more controllers always means better performance (latency rises with quorum size)

context