Learn Labs
7. Reliable Data Delivery

7.2 Replication recap — and what "in sync" precisely means

A replica is in sync if it is the leader, or if it is a follower that:

"Kafka's replication mechanism, with its multiple replicas per partition, is at the core of all of Kafka's reliability guarantees."

Foundations: a partition is stored on a single disk; ordering is guaranteed within it; a partition is either online (available) or offline (unavailable); all events are produced to the leader and usually consumed from it; if the leader becomes unavailable, one of the in-sync replicas becomes the new leader (with one exception — unclean election, §3.2).

The three conditions for being in-sync

A replica is in sync if it is the leader, or if it is a follower that:

  1. Has an active session with ZooKeeper — sent a heartbeat in the last 6 seconds (configurable). → zookeeper.session.timeout.ms, raised from 6 s to 18 s in release 2.5.0.
  2. Fetched messages from the leader in the last 10 seconds (configurable).
  3. Fetched the most recent messages from the leader in the last 10 seconds. “That is, it isn’t enough that the follower is still getting messages from the leader; it must have had no lag at least once in the last 10 seconds.” → replica.lag.time.max.ms, raised from 10 s to 30 s in release 2.5.0.

Condition ③ is the subtle one. A follower that steadily fetches but is permanently 5 seconds behind is out of sync — it never achieves zero lag. Continuous progress is not sufficient; momentary full catch-up is required.

Getting back in: "An out-of-sync replica gets back into sync when it connects to ZooKeeper again and catches up to the most recent message written to the leader. This usually happens quickly after a temporary network glitch is healed but can take a while if the broker the replica is stored on was down for a longer period."

💡 The counterintuitive performance effect of falling out of sync

"An in-sync replica that is slightly behind can slow down producers and consumers — since they wait for all the in-sync replicas to get the message before it is committed. Once a replica falls out of sync, we no longer wait for it to get messages. It is still behind, but now THERE IS NO PERFORMANCE IMPACT. The catch is that with fewer in-sync replicas, the effective replication factor of the partition is lower, and therefore there is a higher risk for downtime or data loss."

we stop waiting for itslightly behind, still in ISRSLOWS EVERYONE — acks=all waitsdrops out of ISRIMPACT VANISHES — looks “fixed”durability: still RF = 3durability: effective RF = 2 — SILENTLY LESS SAFE

The trap: latency recovering is not the incident resolving. It is the symptom moving from latency to risk, where nobody’s dashboard is looking. This is precisely why UnderReplicatedPartitions is a top-tier alert.

Figure 7.2.2The counterintuitive performance effect of falling out of sync

OUT-OF-SYNC REPLICA FLAPPING — historical context

*"In older versions of Kafka, it was not uncommon to see one or more replicas rapidly flip between in-sync and out-of-sync status. This was a sure sign that something was wrong with the cluster. A relatively common cause was a large maximum request size and large JVM heap that required tuning to prevent long garbage collection pauses that would cause the broker to temporarily disconnect from ZooKeeper.

These days the problem is very rare, especially with Kafka 2.5.0+ and its default ZooKeeper connection timeout and maximum replica lag. JVM 8+ (now the minimum supported) with G1 helped curb this — "although tuning may still be required for large messages."

"Generally speaking, Kafka's replication protocol became significantly more reliable in the years since the first edition." References: Jason Gustafson, "Hardening Apache Kafka Replication"; Gwen Shapira, "Please Upgrade Apache Kafka Now."


On this page