skip to content

What is the high watermark, how does it relate to LEO, and why can a consumer not read up to the LEO?

level: seniorimportance: must knowfreq 60%

answer

  1. HW = min LEO across ISR
  2. consumers read offset < HW only
  3. tail between HW and LEO = not yet replicated
  4. acks=all ack == reaches HW
  5. read_committed bound = LSO <= HW

basics

~20 s

The high watermark is the offset up to which records are replicated to all in-sync replicas, so they are 'committed' and safe. Consumers can only read records below the high watermark, never up to the LEO, because un-replicated records could be lost on leader failure.

solid answer

~50 s

The **high watermark (HW)** of a partition is the offset of the first *uncommitted* record — i.e. records below it are fully replicated to every in-sync replica (ISR) and won't be lost if the leader fails. The leader computes it as the minimum LEO across the ISR. Consumers are only allowed to fetch records with offset **< HW**; everything between HW and the leader's LEO exists on the leader but isn't yet acknowledged by all ISR members, so exposing it would risk reading data that vanishes after an unclean truncation or leader change. As followers fetch and advance their LEOs, the leader raises the HW and the new records become visible. With `acks=all` a producer's send is acknowledged once the record reaches HW. (Transactional/read-committed consumers have an even tighter bound — the last stable offset — but the HW is the baseline visibility rule.)

go deeper

for a junior

Know that consumers can only read 'committed' (fully replicated) records, not the very newest unreplicated ones.

for a middle

Define HW as min ISR LEO and explain the read-below-HW rule.

for a senior

Explain the replication-latency gap, acks=all tie-in, and the durability rationale for stopping at HW.

for a principal

Reason about unclean leader election moving HW backward, LSO/read_committed, and the consistency-vs-availability tradeoffs of the watermark design.

## Setup: replicas and the ISR Each partition has one **leader** replica and zero or more **follower** replicas on other brokers. Followers continuously **fetch** from the leader to copy its log. The set of replicas that are sufficiently caught up is the **in-sync replica set (ISR)**. A replica falls out of the ISR if it lags beyond `replica.lag.time.max.ms`. ## Each replica has a LEO Recall the **log-end-offset (LEO)** is one past the last record a replica has. The leader's LEO is the furthest; each follower's LEO trails it depending on how much it has fetched. ## Definition of the high watermark The **high watermark (HW)** is the **minimum LEO across all replicas currently in the ISR**. Equivalently, it is the offset of the first record that has *not* yet been replicated to every in-sync replica. Records with offset `< HW` are **committed**: they are present on all ISR members and are durable against a single leader failure (within the ISR durability guarantee). ## Why consumers stop at HW Kafka only lets consumers read up to **HW − 1** (offsets strictly below HW). Records between the HW and the leader's LEO are physically on the leader but **not yet acknowledged by all ISR followers**. If the leader crashed and a follower that hadn't replicated those records became the new leader, those records would be **truncated** (lost). Exposing them to consumers would mean a consumer could process a record that later disappears — violating consistency. So the HW is the **consumer-visible boundary**: read-uncommitted exposure of un-replicated tail data is disallowed. ## How the HW advances 1. Leader appends a record -> leader LEO rises (HW unchanged yet). 2. Followers fetch the record -> their LEOs rise. 3. On the next fetch, the leader sees all ISR LEOs have advanced and raises the HW to the new minimum. 4. The leader propagates the HW to followers in fetch responses so they know what's committed. This is why there is a small replication-latency gap between a record being written and it becoming consumer-visible. ## Producer view (acks) - `acks=0/1`: ack before full ISR replication — faster, weaker durability. - `acks=all` (with `min.insync.replicas`): the producer is acknowledged only once the record is committed (reaches the HW), tying producer durability directly to the HW. ## Transactions: the last stable offset (LSO) For `isolation.level=read_committed` consumers, the effective ceiling is the **last stable offset (LSO)** — the offset of the first record belonging to an open (not-yet-committed/aborted) transaction. The LSO ≤ HW. Records from in-flight transactions sit below HW but above LSO and remain invisible until the transaction commits. For non-transactional reads (`read_uncommitted`, the default), the bound is simply the HW. ## Edge cases / gotchas - HW ≤ leader LEO **always**; equality means the ISR is fully caught up. - **Unclean leader election** (`unclean.leader.election.enable=true`) can elect an out-of-sync replica and move the HW *backward*, intentionally trading consistency for availability. - A consumer's `position()` can never exceed the HW even though the partition's LEO is higher. - The HW is per-partition and per-leader; it is not a global concept.

  • Why not just let consumers read up to the leader's LEO?
    Records between HW and LEO are only on the leader and not yet replicated to all ISR members. If the leader failed and a less-caught-up follower took over, those records would be truncated and lost — so a consumer could process records that later vanish, breaking consistency.
  • How does acks=all relate to the high watermark?
    With acks=all and min.insync.replicas satisfied, the producer's send is acknowledged only once the record has been replicated to all in-sync replicas — i.e. once it falls below the high watermark and is committed and consumer-visible.
  • What is the last stable offset and how does it differ from the HW?
    The last stable offset (LSO) is the first offset belonging to an open transaction. read_committed consumers can only read below the LSO, which is ≤ HW, so records from in-flight transactions stay hidden until commit/abort even though they're below the HW.

saying these in an interview costs you the question

  • Saying consumers can read up to the LEO.
  • Defining the HW as the leader's LEO rather than the minimum ISR LEO.
  • Claiming the HW can never move backward (unclean leader election can move it back).
  • Ignoring that read_committed adds the tighter LSO bound.

context