Explain how brokers consume the metadata log as an event-sourced state machine, and what guarantees this provides.
answer
- State = fold over ordered records
- MetadataImage / MetadataDelta / publishers
- metadata offset = how caught up
- snapshot then replay tail
- committed (quorum) before applied
basics
~20 sBrokers fetch the replicated metadata log and apply each record in order to build their in-memory view of the cluster. Because everyone replays the same ordered log, all nodes converge to the same state (eventual consistency), with the offset marking how far each has caught up.
solid answer
~40 sEach broker runs a metadata loader that continuously fetches the `__cluster_metadata` log from the controller quorum and applies records to an in-memory `MetadataImage` — an event-sourced state machine where state is the fold over the ordered record stream. The broker tracks its *metadata offset* (how far it has applied); the `MetadataApplier`/publishers react to deltas (e.g. start hosting a new partition leader). Because the active controller is the single writer and the log is totally ordered, any node replaying the same prefix derives identical state — strong per-node consistency, eventual convergence cluster-wide. Brokers periodically load *snapshots* (compacted point-in-time images) instead of replaying from offset 0, then continue from the snapshot's offset. A lagging broker is simply behind on the log; it catches up by fetching more. Fencing prevents a too-far-behind broker from serving as leader.
go deeper
Know brokers replay the log to learn cluster state and all end up consistent.
Describe applying records to an in-memory image and the metadata offset as a catch-up measure.
Explain MetadataImage/Delta/publishers, commit-before-apply, eventual consistency, and snapshots.
Reason about determinism guarantees, fencing of laggards, atomic image swaps, and how event sourcing simplifies recovery vs ZooKeeper watches.
**Event sourcing recap.** In an event-sourced system you do not store the current state; you store an ordered list of *events* (changes). Current state is derived by starting from an empty state and applying every event in order — a *fold* or *reduce*. KRaft applies this to cluster metadata. **The flow.** The active controller appends records (TopicRecord, PartitionRecord, RegisterBrokerRecord, etc.) to `__cluster_metadata`. Every broker acts as a Raft *observer*: it issues `Fetch` requests against the controller quorum to pull new records. Internally, a component (historically `BrokerMetadataListener`, now organized around the `MetadataLoader` + `MetadataImage`/`MetadataDelta` types) applies each record to an immutable in-memory image. Applying records produces a new `MetadataImage`; the difference between consecutive images is a `MetadataDelta`, which *publishers* turn into actions — for instance, when a `PartitionRecord`/`PartitionChangeRecord` makes this broker the leader of a partition, the broker starts accepting produce traffic for it. **Metadata offset.** Each broker knows the *offset* of the last record it applied. This is its position in the event stream and a direct measure of how current its view is. Tools and metrics (e.g. `current-metadata-offset`) expose this; the gap between a broker's offset and the log end is its metadata *lag*. **Guarantees.** 1. *Total order*: single writer + single-partition log ⇒ one canonical event order. 2. *Determinism*: applying the same prefix yields the same state on every node — so two brokers at the same offset are byte-for-byte consistent in their metadata view. 3. *Eventual consistency*: brokers may momentarily be at different offsets, but each is a consistent snapshot of some committed prefix, and all converge as they catch up. 4. *Durability*: only *committed* records (replicated to a quorum majority) are exposed, so replayed state never reflects an uncommitted change that could be lost. **Snapshots.** Replaying from the beginning forever is impractical, so the metadata log is periodically **snapshotted**: a compacted image of state at offset N is written, and the log before N can be truncated. A starting or recovering broker loads the latest snapshot, then replays only records after it. This bounds startup time and disk use. **Edge cases.** A broker that lags badly (e.g. slow disk, GC) keeps a stale view; KRaft *fences* brokers that miss heartbeats so they don't serve as partition leaders with stale metadata. The metadata view is read-locally and lock-free because images are immutable and swapped atomically. And because this is a state machine, recovery is just 'load snapshot, replay tail' — there is no ad-hoc reconciliation logic as there was with ZooKeeper watches.
- Why are metadata snapshots necessary?Without them a node would replay the entire log from offset 0 on startup, which is slow and grows unbounded. Snapshots provide a compacted point-in-time image so only the tail after the snapshot must be replayed.
- How does a broker know how up-to-date its metadata is?It tracks the offset of the last applied record (e.g. current-metadata-offset). The gap to the log end offset is its metadata lag.
- What stops a lagging broker from serving stale data as a leader?Fencing: a broker that misses heartbeats to the controller is fenced and removed from leadership eligibility until it re-registers and catches up.
saying these in an interview costs you the question
- Saying brokers query the controller synchronously per request for metadata
- Claiming all brokers always have identical offsets in real time (it is eventual, not instantaneous)
- Forgetting snapshots and asserting replay always starts at offset 0
- Saying uncommitted records are applied