skip to content

How does replica.lag.time.max.ms control ISR shrink and expansion, and why did Kafka switch from a message-count-based lag check?

level: middleimportance: must knowfreq 70%

answer

  1. Default 30000 ms (30s)
  2. Leader tracks last caught-up-to-LEO time per follower
  3. Catches stalled AND slow followers uniformly
  4. Old replica.lag.max.messages broke on bursts
  5. Check runs every lag.time.max.ms / 2

basics

~20 s

A follower stays in the ISR if it fetched up to the leader's latest offset within replica.lag.time.max.ms (default 30s). Miss that window and the leader removes it; catch back up and it rejoins. The old count-based check (lag in messages) misjudged bursty traffic, so it was replaced by time.

solid answer

~50 s

`replica.lag.time.max.ms` (default 30000 ms) is the single threshold governing ISR membership for followers. The leader tracks, per follower, the last time that follower's fetch request reached the leader's log end offset (i.e., was fully caught up). If a follower either stops sending fetch requests or keeps fetching but never catches up to the leader's LEO for longer than this window, the leader shrinks the ISR by removing it. Once the follower catches up again, the leader expands the ISR. Before KIP-stabilized behavior, Kafka also had `replica.lag.max.messages` — a follower lagging by more than N messages was ejected. That broke under bursty producers: a sudden spike would push every follower past the message threshold and falsely shrink the ISR even though followers were keeping pace. Time-based lag handles both slow followers and stalled followers uniformly without a magic message count.

go deeper

for a junior

Know the config name, the 30s default, and that missing the window drops a follower from the ISR.

for a middle

Explain both the stalled and slow-follower cases, and shrink/expand mechanics.

for a senior

Discuss why time replaced message-count lag and the tuning tradeoffs of the window.

for a principal

Reason about ISR-churn vs. durability tradeoffs at scale and the interaction with min.insync.replicas and alerting.

## The single knob: replica.lag.time.max.ms **`replica.lag.time.max.ms`** (default **30000 ms** = 30 seconds) is a broker-level config that defines how long a follower may go without being fully caught up before the leader evicts it from the ISR. ### How the leader tracks each follower The leader maintains, for every follower, a timestamp: the last time that follower **caught up to the leader's log end offset (LEO)** — the offset of the next record to be appended. Each time a follower's `FetchRequest` arrives requesting an offset equal to the leader's current LEO, the leader knows that follower is caught up *as of now* and updates the timestamp. ### Two ways a follower falls behind 1. **Stalled fetcher:** the follower stops sending fetch requests entirely (crashed, GC pause, network partition, disk stall). No fetch arrives, the caught-up timestamp ages, and after `replica.lag.time.max.ms` the leader shrinks the ISR. 2. **Slow follower:** the follower keeps fetching but the leader is being written faster than the follower can replicate. Its fetches never reach the leader's LEO. The caught-up timestamp stops advancing, and after the window it is removed. Both cases are caught by the *same* time-based rule — that's the elegance of the design. ### Shrink and expand - **Shrink:** when a follower exceeds the window, the leader removes it from the ISR and propagates the new ISR (historically via ZooKeeper; in KRaft mode via an `AlterPartition` request to the controller). - **Expand:** when a removed follower resumes fetching and catches up to the leader's LEO, the leader adds it back to the ISR. The leader checks for lagging replicas periodically (every `replica.lag.time.max.ms / 2` by default). ## Why message-count lag was abandoned The old config **`replica.lag.max.messages`** ejected a follower lagging by more than N messages. The fatal flaw: with a **bursty producer**, a single spike of, say, 10,000 messages would instantly push *every* follower past a threshold of (e.g.) 4000 messages — even followers replicating perfectly well. The ISR would collapse and then re-expand once the burst drained, causing needless churn, false under-replication alerts, and durability scares. There was no single N that worked for both steady and bursty workloads. The time-based approach asks the right question — "has this follower made progress recently?" — independent of throughput. ## Practical tuning notes - Lowering `replica.lag.time.max.ms` makes the cluster eject slow followers faster (tighter durability semantics, but more ISR churn). - Raising it tolerates longer follower stalls before shrinking (fewer false positives, but a slow follower stays counted as in-sync longer, weakening the durability guarantee of acks=all). - It interacts with `replica.fetch.wait.max.ms` and the follower's ability to keep up under load.

  • What exactly does the leader compare to decide a follower is lagging?
    It compares now() against the last time that follower's fetch reached the leader's current log end offset. If that gap exceeds replica.lag.time.max.ms, the follower is removed from the ISR.
  • What was the concrete problem with replica.lag.max.messages?
    Under bursty produce traffic, a single spike could push all followers past the message threshold simultaneously — even healthy ones — causing the ISR to collapse and re-expand spuriously. No fixed message count fit both steady and bursty loads.

saying these in an interview costs you the question

  • Claiming Kafka still uses a message-count lag threshold for ISR (replica.lag.max.messages was removed).
  • Saying the default is something other than 30s without qualification.
  • Thinking only crashed followers fall out — slow-but-alive followers fall out too.
  • Confusing replica.lag.time.max.ms with consumer lag or session timeouts.

context