Kafka: replication, the ISR and the high watermark
On the live cluster, the gold L marks a leader, green F an in-sync follower, red × a replica that fell out of sync. Press Slow down or Stop on a broker to watch the ISR change.
Leader and followers
Each partition has replication.factor replicas on different brokers. One is the leader: every write goes to it, and by default every read too. The others are followers: they send fetch requests to the leader, exactly like a consumer, and append what they get. The first broker in the replica list is the preferred leader.
The ISR
The in-sync replicas are the leader plus the followers that are keeping up. A follower drops out of the ISR when it hasn't caught up to the leader's log end for replica.lag.time.max.ms (30 s by default). It is about time, not record count, so a short burst doesn't kick anyone out but a stuck disk does. Once it catches up again it rejoins.
LEO and the high watermark
- Log end offset (LEO): the offset the next record will get, per replica.
- High watermark (HW): the smallest LEO among the ISR. Every record below it is on every in-sync replica, so it is committed.
Consumers only see records below the HW. A record the leader has but the followers don't yet is invisible to consumers, because it could still be lost.
acks and min.insync.replicas
Producer acks | Leader answers when… | Lost if the leader dies right after? |
|---|---|---|
0 | Never; the producer doesn't wait | Possibly, and nobody knows |
1 | The leader wrote it to its own log | Yes, if followers hadn't fetched it |
all (-1) | Every ISR member has it (the HW passed it) | No, as long as one ISR replica survives |
acks=all alone isn't enough: if the ISR shrinks to just the leader, "all" means one copy. Set the topic's min.insync.replicas (commonly 2 with RF 3): with fewer in-sync replicas than that, acks=all writes fail with NOT_ENOUGH_REPLICAS instead of being accepted with weak durability. acks=all has been the producer default since Kafka 3.0.
When the leader dies
The controller picks the next replica in the list that is alive and in the ISR, and bumps the leader epoch. The new leader has every committed record by definition.
The old leader may have had records beyond the HW that never got replicated (an acks=1 write, say). When it comes back as a follower it asks the new leader where its last epoch ended and truncates everything after that point. Those writes are gone; the producer was only told "written to one replica".
If no ISR replica is alive, the partition goes offline (Leader: none) and waits for one to come back. With unclean.leader.election.enable=true it instead elects any live replica, staying available at the cost of losing committed data. The default is false.
Why the HW lags, and what leader epochs fixed
Followers pull. The only way the leader learns how far a follower got is the offset in its next fetch, so the HW moves one fetch after a copy arrives, and the follower hears the new HW one fetch later still. Before Kafka 0.11, a follower that restarted cut its log back to its own (stale) HW. KIP-101 showed two ways that loses data: a quick restart followed by a leader failure drops an acked record, and a power cut on both replicas leaves them holding different records at the same offset while both count as in sync.
Since 0.11 every record carries its leader epoch. A returning follower asks the leader "where did my last epoch end?" (OffsetsForLeaderEpoch) and truncates only past that point.
Kafka doesn't fsync each write either: records sit in the OS page cache, and durability comes from copies on separate machines. acks=all with min.insync.replicas=2 and replication factor 3 is what makes one failure survivable.
Try it
- acks=1 loses acked writes
- acks=all is only as strong as the ISR
- Unclean leader election
- KIP-101: a restart loses an acked record, then the same with leader epochs
- KIP-101: logs diverge after a power cut
Commands
kafka-topics --bootstrap-server localhost:9092 --describe --under-replicated-partitions
kafka-topics --bootstrap-server localhost:9092 --describe --unavailable-partitions
kafka-leader-election --bootstrap-server localhost:9092 --election-type PREFERRED --all-topic-partitions