What is a broker epoch in KRaft, how is it assigned, and what problem does it solve during broker restarts?
answer
- epoch = offset of registration record in __cluster_metadata
- monotonic, unique per registration
- fencing token / zombie guard
- STALE_BROKER_EPOCH rejects old incarnation
- KRaft analogue of ZK session/zkVersion
basics
~20 sA broker epoch is a monotonically increasing id the controller assigns at registration (the metadata log offset of the registration record). Every heartbeat carries it, so the controller can reject stale requests from a previous broker incarnation.
solid answer
~50 sWhen a broker registers, the controller persists a registration record in the __cluster_metadata log and assigns the broker a broker epoch equal to that record's offset — so it strictly increases over time and is unique per registration. The broker echoes this epoch in every BrokerHeartbeat and in other controller RPCs. This is the KRaft analogue of the ZooKeeper session/zkVersion fencing token. Its purpose is to disambiguate incarnations: if a broker crashes and restarts, it registers afresh (new incarnation id) and gets a higher epoch. Any in-flight or delayed heartbeat from the old incarnation carries the stale epoch and is rejected with a stale-broker-epoch error, preventing a zombie process from interfering with the new one. The controller also uses the epoch to validate that AlterPartition and similar requests come from the current incarnation, avoiding split-brain on ISR updates.
go deeper
Know there is a broker epoch and it is bumped on each registration; deep details optional.
Explain that the epoch travels in heartbeats and lets the controller reject stale requests.
Explain offset-derived monotonicity, STALE_BROKER_EPOCH handling, and the zombie/split-brain problem it prevents.
Compare to ZK session fencing, relate to AlterPartition validation, and distinguish broker vs leader vs partition epochs in a coherent fencing-token model.
## The problem: zombies and split brain Consider a broker that appears to die (it stops heartbeating, gets fenced) but is actually just paused — a long GC, a network blip, or a hung process. Meanwhile it restarts a new process, or the old process resumes. Now there may be **two senders** claiming to be broker 3. If the controller accepted requests from both, the cluster could end up with inconsistent ISR or leadership decisions — a **split-brain** failure. We need a way to tell the *current* incarnation from a *stale* one. That mechanism is the **broker epoch** (a 'fencing token'). ## How the epoch is assigned KRaft stores all cluster state as records in the internal **`__cluster_metadata`** Raft log. When a broker sends `BrokerRegistration`, the active controller appends a **registration record** to that log. The **offset** at which that record is committed becomes the broker's **broker epoch**. Because log offsets only ever increase, every new registration yields a strictly larger epoch than any previous one for that (or any) broker. The registration also includes a per-process **incarnation id** (a random UUID) that lets the controller recognize a genuinely new process versus a duplicate registration of the same one. ## How the epoch is used The broker stores its assigned epoch and includes it in: - every **`BrokerHeartbeat`**, - `AlterPartition` requests (ISR shrink/expand proposals), - controlled-shutdown requests. The controller checks the epoch on each request. If the supplied epoch is **older than the current registered epoch** for that broker id, the request is rejected — typically surfaced as a `STALE_BROKER_EPOCH` error. The stale sender (the zombie) thus cannot mutate cluster state. When the broker sees this error, it knows it has been superseded and must re-register. ## Relationship to the old ZooKeeper world In ZK mode, the equivalent guard was the **ZooKeeper session** plus znode versions (`zkVersion`) used as fencing on controller and ISR writes. KRaft replaces session-based fencing with the **log-offset-based broker epoch**, which is simpler to reason about because it derives directly from the totally-ordered metadata log. ## Edge cases and nuances - **Restart after crash**: new incarnation id -> new registration record -> higher epoch. Old in-flight heartbeats are rejected. - **Controlled shutdown then restart**: still a fresh registration and epoch; the controller does not reuse the old one. - **Epoch vs leader epoch / partition epoch**: do not confuse the *broker* epoch (per broker registration) with the *leader epoch* (per partition leadership term) — they are different counters serving different fencing purposes. - The metadata offset a broker reports in heartbeats (how caught-up it is) is distinct from its broker epoch; the former drives unfencing, the latter drives incarnation validity.
- How does the controller pick the epoch value, and why does that guarantee monotonicity?The epoch is the offset of the broker's registration record in the __cluster_metadata Raft log. Log offsets are strictly increasing, so each new registration necessarily gets a larger epoch.
- What error does a stale broker incarnation receive, and what does it do in response?It receives STALE_BROKER_EPOCH. The broker recognizes it has been superseded by a newer registration and must re-register (effectively restart its registration lifecycle) rather than continue acting on cluster state.
- How is broker epoch different from leader epoch?Broker epoch identifies a broker's registration/incarnation cluster-wide; leader epoch identifies a leadership term for a single partition. Both are fencing tokens but at different scopes.
saying these in an interview costs you the question
- Confusing broker epoch with partition leader epoch
- Saying the epoch is a wall-clock timestamp or a random number rather than the registration record offset
- Claiming epochs can decrease or be reused on restart
- Forgetting that the epoch's purpose is to fence zombie/stale incarnations