skip to content

What is the leader epoch, and why is it bumped on every leadership change?

level: seniorimportance: should knowfreq 52%

answer

  1. monotonic term number per partition
  2. OffsetsForLeaderEpoch -> truncation boundary
  3. KIP-101 / KIP-279 fix HW truncation
  4. fences zombie old leader
  5. FENCED_LEADER_EPOCH to clients

basics

~20 s

The leader epoch is a number that increases by one each time a partition gets a new leader. It lets brokers and clients tell stale leadership apart from current, so followers truncate correctly and old leaders can't corrupt the log.

solid answer

~40 s

Each partition has a leader epoch: a monotonically increasing integer stamped into the log, incremented every time the controller elects a new leader. It solves two problems. First, log divergence: after a crash, a follower asks the new leader 'what was the end offset of epoch N?' via an OffsetsForLeaderEpoch request, and truncates anything beyond that point — replacing the old, unreliable high-watermark-based truncation that could lose or duplicate data (fixed by KIP-101 and KIP-279). Second, fencing: a zombie former leader that comes back with a stale epoch is rejected, so it can't append or mislead clients. Clients also cache the epoch and use it to detect stale metadata. The epoch is the version token that makes truncation deterministic and leadership unambiguous.

go deeper

for a junior

Just know there's a version number that increases each time leadership changes.

for a middle

Explain that followers use the epoch to truncate to the right offset and that it tells new leadership from old.

for a senior

Explain the OffsetsForLeaderEpoch lookup, why HW-based truncation (pre-KIP-101) was unsafe, and zombie fencing.

for a principal

Discuss the Raft-term analogy, KIP-101/279 correctness guarantees, and how epoch fencing interacts with client metadata consistency.

## What it is The **leader epoch** is a per-partition integer that the controller increments by exactly one on every leadership change. The leader writes its current epoch into the partition log alongside the records (this is the *leader epoch cache* / `leader-epoch-checkpoint` file). So the log is effectively a sequence of segments tagged with which leadership term produced them. Think of it as a **monotonic version number for who is in charge** of the partition — analogous to a term number in Raft. ## Problem 1 it solves: deterministic follower truncation Followers replicate by fetching from the leader. After a crash and failover, a follower may have a log that diverges from the new leader's — it might hold records the new leader never had, or be missing some. The old (pre-KIP-101) recovery used the **high watermark** (HW = highest offset known replicated to all ISR members) to decide where to truncate. Because the HW propagates asynchronously, this could cause **data loss or log divergence** in certain crash/restart orderings. With leader epochs, a recovering follower sends an **OffsetsForLeaderEpoch** request: *"For epoch N, what was your last offset?"* The leader replies with the first offset of epoch N+1 (the boundary). The follower truncates its log to that boundary, discarding anything it wrongly kept, then resumes fetching. This is **deterministic** and matches the authoritative leader exactly. KIP-101 introduced this; **KIP-279** fixed a follower-to-follower divergence gap by making the epoch lookup chase down the correct boundary across multiple epochs. ## Problem 2 it solves: fencing zombie leaders Suppose a leader is network-partitioned, a new leader is elected (epoch bumped), then the old leader reconnects still thinking it's the leader (a 'zombie'). Because its epoch is now stale, brokers and the controller reject its actions. The epoch is the token that **fences** stale leadership so it cannot append records or hand stale answers to clients. ## Client side Producers and consumers also track the leader epoch in their metadata. A consumer includes the epoch in fetch/offset requests, so if it has stale metadata pointing at an old epoch, the broker can signal `FENCED_LEADER_EPOCH` and the client refreshes. This prevents a client from, e.g., reading from a deposed leader. ## Where you see it - `kafka-dump-log.sh` shows epoch entries. - The `leader-epoch-checkpoint` file in each partition directory. - Metrics/admin output expose the current epoch via DescribeTopics / the partition state. ## Why 'bump on every change' Monotonic increment is what makes 'newer beats older' unambiguous across the whole cluster without coordination. If epochs ever repeated or moved backward, truncation and fencing decisions would be ambiguous — exactly the bug class KIP-101 removed.

  • Before leader epochs, what mechanism decided where a follower truncated, and why was it unsafe?
    It used the high watermark (highest offset replicated to all ISR). Because the HW propagates asynchronously, certain crash/restart sequences caused followers to truncate to the wrong point, leading to data loss or log divergence between replicas. KIP-101 replaced it with the epoch-boundary lookup.
  • How does the leader epoch help fence a 'zombie' former leader?
    On every leadership change the epoch increments. A returning old leader carries a stale (lower) epoch, so its append/fetch attempts are rejected and clients targeting it get FENCED_LEADER_EPOCH and refresh — the stale leader can't corrupt the log or serve stale reads.

saying these in an interview costs you the question

  • Saying followers still truncate using only the high watermark in modern Kafka (KIP-101 replaced that).
  • Confusing leader epoch with the producer epoch (idempotent-producer fencing) — different mechanisms.
  • Claiming the epoch is global per broker; it is per-partition.
  • Saying the epoch can reset or decrease — it is monotonic.

context