Why does a record produced to the leader not become visible to consumers immediately, and how do followers and the high watermark control commit/visibility?
answer
- LEO = log tail; HW = commit point
- HW = min ISR LEO
- Consumers read only up to HW
- Committed = on all ISR replicas
- acks=all waits for HW to pass the record
basics
~20 sA record appended on the leader is only 'committed' once all in-sync followers have replicated it. The high watermark marks the highest committed offset; consumers can only read up to it, so un-replicated tail records stay invisible.
solid answer
~50 sWhen a producer write lands on the leader, the leader appends it to its log immediately, advancing its log-end offset (LEO). But that record is not yet *committed* and not yet readable. A record is committed only when every replica in the ISR has fetched and appended it. The leader tracks each follower's LEO (learned from their fetch offsets) and sets the partition's **high watermark (HW)** to the minimum LEO across the current ISR. Consumers are only allowed to read up to the HW — never the leader's un-replicated tail — so that any record a consumer can see is guaranteed to survive a leader failure (it exists on all in-sync replicas). With `acks=all`, the producer's send isn't acknowledged until the HW reaches its record, tying producer durability to the same boundary. This is why replication lag directly delays both consumer visibility and producer acks.
go deeper
Know consumers only see records after they're replicated, marked by the high watermark.
Define LEO vs HW and that committed means replicated across the ISR.
Connect HW = min ISR LEO to acks=all, min.insync.replicas, and consumer-visibility guarantees, including the round-trip lag.
Reason about the durability/availability trade across ISR shrink, unclean election, and epoch-based truncation on failover.
**Two offsets per replica.** Each replica tracks its **log-end offset (LEO)** — the offset of the next record to be written, i.e. the tail of its log. The leader additionally maintains the **high watermark (HW)** for the partition. **Append != commit.** When a producer sends to the leader, the leader writes the record to its local log right away, bumping its own LEO. At this instant the record exists on exactly one broker. If the leader crashed now and a follower took over, that record could vanish. So Kafka draws a line between *appended* and *committed*. **Defining committed.** A record is **committed** once **all replicas in the ISR** (in-sync replica set) have replicated it. The leader knows each follower's LEO because every follower's Fetch request carries a `fetchOffset` equal to its LEO (the follower has persisted everything below it). The leader computes: `HW = min(LEO of every replica currently in the ISR)` The HW is the highest offset that is fully replicated across the ISR — the commit point. **Consumer visibility.** Consumers are served records only **up to (HW - 1)**; they can never read the leader's un-replicated tail between HW and the leader's LEO. This guarantees the *read-your-survivable-data* property: anything a consumer has seen is already on every in-sync replica, so a leader failover (which promotes an ISR member) won't make previously-read records disappear. Without this, consumers could read a record that later evaporates — breaking consistency. **Producer acks tie-in.** - `acks=0`: producer doesn't wait at all. - `acks=1`: leader acks after its own append (before followers replicate) — fast, but a record can be lost if the leader dies before followers catch up. - `acks=all` (a.k.a. -1): the leader acks only once the record is committed — i.e., the HW has advanced past it, meaning all in-sync replicas have it. Combined with `min.insync.replicas`, this is the strong-durability setting. **`min.insync.replicas`.** This caps how small the ISR may shrink while still accepting `acks=all` writes. If the ISR drops below it (e.g. only the leader left), the leader rejects `acks=all` produces with `NotEnoughReplicas` rather than commit data that lives on too few brokers. **How HW advances over time.** Leader appends record at offset N (leader LEO = N+1). Followers fetch, append, and their next fetch reports LEO = N+1. Once the slowest ISR follower reports N+1, the leader raises HW to N+1, committing offset N. Note HW advancement is therefore one *fetch round-trip behind* the leader's append — the origin of inherent replication latency. **Edge cases.** - A follower lagging beyond `replica.lag.time.max.ms` is ejected from the ISR; HW can then advance using the remaining ISR (availability vs durability trade governed by `min.insync.replicas` and `unclean.leader.election.enable`). - On leader failover, the new leader's HW may be behind some followers' LEOs; followers truncate to the new leader's HW/epoch to stay consistent (KIP-101).
- With acks=all and min.insync.replicas=2 in a 3-replica partition, what happens if two of the three brokers fail?Only the leader remains in the ISR, which is below min.insync.replicas=2. The leader rejects acks=all produce requests with NotEnoughReplicas(AfterAppend) so it won't commit data that lives on too few replicas. Existing committed data stays readable up to the HW. Producers must wait/retry until a follower rejoins the ISR.
- Why is the high watermark always at least one fetch round-trip behind the leader's log-end offset under continuous load?The leader appends immediately (advancing its LEO), but it can only raise the HW after the slowest ISR follower fetches that offset and its next fetch reports the new LEO. That round-trip — fetch the data, then report progress on the subsequent fetch — is the inherent lag between append and commit.
saying these in an interview costs you the question
- Saying consumers can read the leader's latest appended record before it's replicated
- Claiming a record is committed as soon as the leader writes it (true only for acks=1, not durability)
- Confusing LEO (log tail) with HW (commit point)
- Thinking acks controls replication itself rather than when the producer is acknowledged