02 / 05

What is the high watermark used for conceptually?

Difficulty: 4/10
Offsets, Watermarks, Coordinators

The High Watermark: Replication and Visibility Boundary

The high watermark (HW) is the offset up to which all in-sync replicas (ISR) of a partition have replicated. Conceptually, it is the boundary between records that are fully replicated and records that are not. Records below the HW are considered committed and durable: even if the leader fails, a new leader elected from the ISR will have those records. Records above the HW are not yet fully replicated; if the leader fails and a new leader is elected, those records may be lost. The HW is therefore the visibility boundary for consumers: a consumer using the default read_uncommitted isolation can read up to the LEO, but it may read records that are not yet replicated and could be lost. A consumer using read_committed can only read up to the last stable offset (LSO), which is at or below the HW. The HW is maintained by the leader, which computes it as the minimum LEO of the ISR. It is replicated to followers as part of the fetch protocol, so all replicas know the HW. The HW is a key part of Kafka's replication and durability model.

The mechanism that updates the HW is the replication protocol. When a follower fetches from the leader, it includes its own LEO in the fetch request. The leader uses the follower's LEO to update its view of the follower's progress. The HW is the minimum of the LEOs of all ISR members. When a follower catches up and its LEO advances, the leader may advance the HW. The HW is then included in the fetch response to the followers, so they can update their own HW. This is how the HW propagates through the replica set. The HW is used for two purposes: it is the boundary for committed data, and it is the boundary for consumer visibility in read_committed mode. The trade-off is between freshness and safety. Reading up to the LEO gives the freshest data but may expose records that are lost on leader failover. Reading up to the HW gives only replicated records, which is safer but slightly stale. For most use cases, the difference is milliseconds, but for latency-sensitive applications, it matters. Version note: the concept of the HW has been stable since early Kafka. With transactions (Kafka 0.11+), the LSO was introduced as the visibility boundary for read_committed consumers; the LSO is the minimum of the HW and the first open transaction's offset. For consumers not using transactions, the LSO equals the HW. In KRaft mode, the data-plane HW is unchanged; the metadata log has its own HW managed by the controller quorum.

A common mistake is to confuse the HW with the LEO. The LEO is the log's extent; the HW is the replicated boundary. On a healthy cluster with no lagging followers, the HW and LEO are close, but they are not the same. Another mistake is to assume that read_uncommitted consumers can read up to the LEO. In practice, they can, but they may see records that are not committed. A third mistake is to use the HW as a measure of consumer progress; the HW is a property of the partition, not the consumer. The trade-off is between durability and latency. Waiting for the HW to advance before consuming adds a small delay but ensures that the consumer never sees data that could be lost. For a financial pipeline, this is the right choice. For a metrics pipeline where occasional loss is acceptable, reading up to the LEO is fine. Version note: the HW is exposed via JMX as kafka.log:type=Log,name=HighWatermark,topic=,partition=. Monitoring the HW and the LEO together helps detect replication lag and understand consumer visibility. In KRaft, the controller quorum's metadata log also has a HW, which is used for metadata consistency.

javascript
  1. 1

    HW is the offset up to which all ISR members have replicated.

  2. 2

    Records below the HW are committed and durable; records above may be lost on leader failover.

  3. 3

    HW is the minimum LEO of the ISR members.

  4. 4

    HW is propagated to followers via the fetch protocol.

  5. 5

    read_uncommitted consumers can read up to LEO; read_committed consumers read up to LSO.

  6. 6

    LSO equals HW when there are no open transactions; otherwise it is the first open transaction's offset.

  7. 7

    HW and LEO are close on a healthy cluster but are not the same.

  8. 8

    Monitoring HW vs LEO helps detect replication lag.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.