What problem did having ZooKeeper AND the Kafka controller both hold metadata create, and how could they diverge?
answer
- ZK + controller cache + broker caches = ≥2 truths
- ZK write and RPC fan-out not atomic
- write-to-ZK-then-crash before pushing RPC
- no single global metadata offset
- KRaft: one ordered Raft log, brokers are consumers
basics
~20 sMetadata lived in ZooKeeper but the controller cached it in memory and also pushed it to brokers. With two copies that update separately, they could drift, causing brokers to act on stale or inconsistent metadata.
solid answer
~50 sLegacy Kafka had effectively two places metadata lived: ZooKeeper (the durable store) and the controller's in-memory cache, which it then propagated to brokers via direct RPCs (LeaderAndIsr, UpdateMetadata). These updates were not a single ordered, replicated log — ZK writes, the controller's cache, and the broker-side caches could get out of sync. A controller could write to ZK then fail before pushing the RPC, or brokers could receive updates out of order, leaving them with stale leadership info. There was no single monotonic version that every component agreed on. KIP-500's answer is a single source of truth: an ordered, Raft-replicated metadata log with a metadata offset. Every broker consumes the same log in the same order, so 'who is the leader for partition X at offset N' is unambiguous, and there's no second store to reconcile.
go deeper
Know metadata existed in ZooKeeper and also in the controller/brokers in memory, and two copies can get out of sync.
Explain that ZK writes and broker RPC pushes weren't atomic, so brokers could hold stale leadership info.
Articulate the dual-source-of-truth smell, concrete divergence scenarios, and why one ordered replicated log resolves them.
Discuss consistency guarantees: bounded ordered lag vs divergence, fencing vs total order, and the general anti-pattern of two stores that must agree.
## The two stores In the ZooKeeper era, cluster metadata physically existed in more than one place: 1. **ZooKeeper** — the durable, authoritative store (znodes for topics, partitions, ISR, configs). 2. **The controller's in-memory cache** — what the controller used to make fast decisions. 3. **Each broker's local metadata cache** — pushed to brokers by the controller so they know which partitions they lead/follow and where to route requests. That's effectively **three representations of the same truth**, and crucially they were **updated by different mechanisms**: the controller wrote to ZK *and* separately sent RPCs (`LeaderAndIsr`, `UpdateMetadata`, `StopReplica`) to brokers. ZK changes and RPC fan-out were not one atomic, ordered operation. ## How they diverged Because the updates weren't a single ordered, replicated log, several skew scenarios were possible: - **Write-then-crash:** the controller persists a change to ZK, then dies before pushing the corresponding RPC to brokers. Brokers now lag the durable state. - **Out-of-order / lost RPCs:** metadata pushes were individual RPCs; a slow or dropped one could leave a broker with stale leader/ISR info. - **Cache vs ZK skew on failover:** a new controller reloads from ZK (the snapshot model), but in-flight changes that hadn't been durably reconciled could be momentarily inconsistent. - **No global version:** there was no single monotonically increasing offset that *every* component agreed defined 'the current metadata'. The controller epoch helped fence stale controllers but didn't give a unified, replayable order for all metadata. The symptom for users: brokers occasionally acting on stale leadership (e.g., `NOT_LEADER_FOR_PARTITION` errors), and operational complexity reconciling 'what ZK says' vs 'what the cluster is doing'. ## Why a single source of truth fixes it KRaft (KIP-500) collapses this to **one ordered, Raft-replicated metadata log** (the internal `__cluster_metadata` topic). Metadata changes are **records appended at increasing offsets**. The active controller is the only writer; every other controller and broker is a **consumer** of that same log, applying records **in offset order**. Consequences: - There is exactly **one authoritative ordering**. 'Leader of partition X as of metadata offset N' is well-defined and identical everywhere once N is applied. - Brokers learn their own lag: a broker knows which metadata offset it has caught up to. - No second store (ZK) to keep consistent — eliminating the entire class of reconciliation bugs. - Failover is safe because the log is already replicated to standbys. ## Nuance / edge cases - KRaft still has *propagation lag* — a broker may be a few offsets behind — but it's **bounded and ordered**, not a divergence between two independent stores. Staleness is now 'old but consistent', not 'inconsistent'. - The **controller epoch** in the ZK design provided fencing (rejecting requests from a deposed controller) but did not give the unified, replayable metadata ordering KRaft provides. - 'Dual source of truth' is a design smell generally: two stores that must agree but update independently will eventually disagree.
- Does KRaft eliminate metadata propagation lag entirely?No — brokers can still be a few offsets behind the log. But the staleness is bounded and consistently ordered ('old but correct'), not a divergence between two independent stores that can contradict each other.
- What was the controller epoch for, if not ordering metadata?It fenced off a stale/deposed controller so brokers would reject its requests (zombie-controller protection). It did not provide a single replayable ordering for all metadata across the cluster.
saying these in an interview costs you the question
- Saying KRaft removes all staleness — it removes divergence between stores, but ordered propagation lag still exists.
- Treating ZooKeeper and the controller cache as always perfectly in sync — the whole point is they could drift because updates weren't a single ordered log.
- Confusing controller epoch (fencing) with a unified metadata ordering.