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?
answer
- Single writer = one total order
- commit on quorum majority (Raft)
- lose majority -> writes stall, reads OK
- fast failover, no ZK reload
- scales to millions of partitions
basics
~20 sThe 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 sFunneling 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
Know the single controller writes metadata and a quorum keeps it durable.
Contrast with ZooKeeper at a high level: no external system, faster failover.
Explain commit-on-majority, the write-availability vs read-availability split, and quorum sizing.
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)