After a failover, how does a former-leader replica that rejoins reconcile its log with the new leader, and what role does the high watermark play?
answer
- LEO = log end; HW = committed-to-all-ISR
- orphans live above old HW
- OffsetsForLeaderEpoch -> divergence offset
- truncate then re-fetch, rejoin ISR
- consumers never saw above-HW records
basics
~20 sWhen the old leader comes back as a follower, it may hold uncommitted records past the high watermark that the new leader never had. Using leader-epoch lookups it finds the divergence point, truncates the extra records, then re-fetches from the new leader to catch up and rejoin the ISR.
solid answer
~40 sThe high watermark (HW) is the highest offset confirmed replicated to all ISR members — everything below it is committed and exposed to consumers; records above it are uncommitted and may be discarded. A crashed leader can have appended records beyond its HW that were never replicated. When it rejoins as a follower of the new leader, those extra records are invalid. It uses OffsetsForLeaderEpoch to ask the new leader for the end offset of each epoch it shares, finds the exact divergence offset, truncates everything past it, and resumes fetching. This epoch-based truncation (KIP-101/279) replaced the older HW-only approach that could diverge. Once the follower catches up within replica.lag.time.max.ms it re-enters the ISR. Consumers were never exposed to the truncated records because they sat above the HW.
go deeper
Know there's a 'safe read line' (high watermark) and that uncommitted records can be dropped on failover.
Distinguish LEO vs HW and state that a rejoining replica truncates records above the old HW.
Explain epoch-based divergence-offset lookup and why it replaced HW truncation; tie to acks levels.
Reason end-to-end about HW propagation lag, KIP-101/279 correctness, and how acks/min.insync choices determine which records are truncatable vs durable.
## Two offsets you must distinguish - **Log End Offset (LEO):** the offset of the next record to be appended to a replica's log — i.e. how far that replica's log goes. - **High Watermark (HW):** the highest offset that is known to be replicated to **all** ISR members. The leader advances the HW only once the slowest in-sync follower has the record. **Consumers can only read up to the HW** — records above it are 'uncommitted' and may still vanish. So on the leader, LEO >= HW. The gap (HW..LEO) is records appended but not yet fully replicated — durable enough to keep, but **not yet promised** to consumers. ## How divergence arises The leader may append records (advancing its LEO) and **crash before** those records reach the followers and before the HW advances past them. The old leader's on-disk log now contains records that: - are **above the old HW**, and - the **new leader never received**. The new leader was elected from the ISR, so its log is the authoritative truth from the cluster's standpoint. The old leader's extra records are orphans. ## Reconciliation when the old leader rejoins 1. It comes back and becomes a **follower** of the new leader. 2. It must not blindly append from the new leader's current end — its own log diverges *earlier*, at the orphaned records. It needs the **divergence offset**. 3. It sends **OffsetsForLeaderEpoch** for the most recent leader epoch it shares with the new leader: *"what's the end offset of epoch E?"* If the follower's own data for E extends past the leader's epoch boundary, it truncates to that boundary. KIP-279 generalized this to walk back through multiple epochs if needed so two followers can't diverge. 4. It **truncates** its log to the divergence offset, deleting the orphaned records. 5. It **fetches** from there forward, replicating the new leader's records, until its LEO catches the leader's. 6. Once it has been caught up (within `replica.lag.time.max.ms`, default 30000ms), the leader re-adds it to the **ISR**. ## Why this is safe for consumers The discarded records were **above the old HW**, so consumers were **never allowed to read them** — there is no consumer-visible rollback for committed data. This is the whole point of the HW: it draws the line between 'durable and visible' and 'speculative, may be truncated'. ## Why epoch-based (not HW-based) truncation The legacy approach truncated the follower to its own HW on restart. Because the HW propagates asynchronously and lags, specific crash orderings let two replicas keep *different* records below their respective HWs, causing **permanent log divergence or silent data loss** (KIP-101's motivating bug). Epoch boundaries are exact and agreed, making truncation deterministic. KIP-279 closed a remaining follower-vs-follower gap. ## Edge interactions - With **acks=all + min.insync.replicas>=2**, the HW only advances when enough replicas hold a record, so the records that get truncated are precisely those never acknowledged to the producer — no acknowledged data is lost. - With **acks=1**, a record could be acknowledged by the leader yet sit above the HW (not yet on a follower); a crash can then truncate it — acknowledged-but-lost. That's why acks=1 is weaker.
- Why can records above the high watermark be safely truncated without any consumer noticing?Consumers can only read up to the HW; anything above it is uncommitted and was never delivered. Truncating those orphaned records removes data no consumer was ever allowed to see, so there's no visible rollback of committed data.
- With acks=1, how can a record be acknowledged to the producer yet still be lost on failover?acks=1 acknowledges as soon as the leader writes it, before any follower replicates it — so it can sit above the HW. If the leader crashes before replication, the new leader never had it, and on rejoin the old leader truncates it: acknowledged but lost. acks=all avoids this.
saying these in an interview costs you the question
- Saying consumers can read past the high watermark (they cannot).
- Claiming the rejoining old leader keeps its extra records and the new leader catches up to them (it's the reverse — the old leader truncates).
- Using HW-based truncation as the current mechanism (epoch-based replaced it).
- Equating LEO and HW — the gap between them is exactly the truncatable region.