skip to content

Walk through exactly what data loss and log divergence happens when an unclean leader election promotes an out-of-sync replica.

level: seniorimportance: must knowfreq 60%

answer

  1. old HW 1000, lagging at 800
  2. 801-1000 committed but lost
  3. new leader reuses offsets -> divergence
  4. returning replica truncates via leader epoch (KIP-101/279)
  5. HW moves backward for consumers

basics

~20 s

The lagging replica becomes leader at its own (shorter) log end. Records the old leader had committed but this replica never fetched are gone. When the old replica returns, it must truncate its longer log down to the new leader's offset to re-join, throwing away those records.

solid answer

~50 s

Say the old leader committed records up to offset 1000, but an out-of-sync replica only fetched up to offset 800. If all ISR replicas die and the out-of-sync one is elected (unclean), it becomes leader with a log ending at 800. New writes append from 801 onward on the NEW leader's epoch. The committed records 801-1000 from the old leader are permanently lost — consumers that already read them now see different data written at those offsets, and consumers that hadn't read them never will. When an old ISR replica recovers, it discovers its log diverges from the new leader (same offsets, different records). Using leader-epoch information it truncates its log back to the divergence point (800) and re-replicates from the new leader. So offsets are reused for different data: that is the log divergence. The High Watermark effectively moves backward from the consumer's perspective.

go deeper

for a junior

Understand the gist: the new leader is behind, so the records it never received are gone.

for a middle

Be able to narrate the offset example and that returning replicas truncate to converge.

for a senior

Explain leader-epoch truncation (KIP-101/279), offset reuse, and the consumer-visible HW-moves-backward effect.

for a principal

Reason about end-to-end data-integrity guarantees, downstream phantom-read handling, and why pre-epoch HW truncation was itself unsafe.

## Setup terms - **Offset**: the monotonically increasing position of a record within a partition log. - **Log End Offset (LEO)**: the offset just past the last record a replica has. - **High Watermark (HW)**: the highest offset that is **committed** (replicated to all ISR members) and therefore readable by consumers. - **Leader epoch**: a number that increments each time leadership changes; stamped on records so replicas can detect and reconcile divergence. ## The scenario 1. Leader L and follower F1 are in the ISR. Both have records up to offset **1000**; HW = 1000, so consumers have read up to 1000. 2. Follower F2 had network trouble, fell behind at offset **800**, and was **removed from the ISR**. It is alive but out of sync. 3. L and F1 both crash. The ISR is now empty of live members. ## With unclean.leader.election.enable = false The partition goes **offline**. No leader is elected. Producers/consumers for it block until L or F1 returns. No data is lost. ## With unclean.leader.election.enable = true Kafka elects **F2** (the out-of-sync survivor) as leader. - F2's log ends at **800**. It is now leader at a **new leader epoch**. - Records **801–1000**, which were committed and already consumed by some clients, exist on no live broker. They are **permanently lost**. - New producer writes append starting at offset **801** again — but these are *different records* than the original 801–1000. ## Log divergence and truncation When L or F1 recovers, it has records at offsets 801–1000 that **disagree** with the new leader's records at those same offsets. This is **log divergence**: the same offset range holds different data on different replicas. Kafka resolves it with **leader-epoch–based truncation** (KIP-101, refined by KIP-279). The returning replica issues an `OffsetsForLeaderEpoch` request, finds that its epoch diverges at offset 800, **truncates** everything from 801 onward, and re-fetches the new leader's records from 801. The old committed records are discarded so the cluster converges on one log. ## Consumer-visible effect - Consumers that read 801–1000 before the failure have processed records that no longer exist — effectively **phantom reads**. - The High Watermark, from a consumer's standpoint, **moved backward**: offsets it had already passed now contain new, different data. - Consumer offsets stored for that partition may now point past the new LEO or into rewritten territory, causing `OffsetOutOfRange` handling (reset to earliest/latest) depending on `auto.offset.reset`. ## Why epochs make this safe-ish Before leader epochs existed (pre-0.11), HW-based truncation could itself cause silent divergence even without unclean election. Leader epochs give every replica a deterministic divergence point so recovery is correct; they do **not**, however, prevent the loss inherent to choosing a behind replica.

  • What mechanism does a recovered replica use to reconcile its diverged log with the new leader?
    Leader-epoch–based truncation (KIP-101, KIP-279). It sends an OffsetsForLeaderEpoch request, finds the offset where epochs diverge, truncates everything after it, and re-replicates from the new leader.
  • Can consumers detect that records they already processed have been overwritten?
    Not directly — there is no built-in 'this offset's data changed' signal. They may observe OffsetOutOfRange if their committed offset is now beyond the new LEO, but already-processed records that got rewritten are effectively silent phantom reads handled at the application level.

saying these in an interview costs you the question

  • Saying no data is lost because the replica re-replicates — re-replication discards the old records; the loss is real and permanent.
  • Claiming offsets are never reused — under unclean election the new leader DOES write different data at previously committed offsets.
  • Ignoring leader epochs and asserting truncation uses the High Watermark (that's the pre-0.11 buggy behavior KIP-101 fixed).

context