skip to content

How does a follower replica actually stay in sync with the leader? Describe the ReplicaFetcherThread and its fetch loop.

level: middleimportance: must knowfreq 60%

answer

  1. Pull, not push — follower asks
  2. ReplicaFetcherThread per leader broker, grouped
  3. fetchOffset = follower LEO = implicit ack
  4. Long-poll up to wait.max.ms
  5. Leader advances HW from min ISR LEO

basics

~20 s

Each follower broker runs ReplicaFetcherThreads. A thread sends Fetch requests to the leader asking for records starting at the follower's current log-end offset, appends what it gets to its local log, and advances its offset — repeating continuously.

solid answer

~50 s

Replication is pull-based: the follower fetches, the leader does not push. On a follower broker, the `ReplicaFetcherManager` creates `ReplicaFetcherThread` instances, each responsible for a set of partitions whose leaders live on a given remote broker. A thread loops: it builds a Fetch request listing, per partition, the offset to start from (the follower's current log-end offset / LEO), sends it to the leader, and waits up to `replica.fetch.wait.max.ms` for data (long-poll). The leader returns records from that offset, up to `replica.fetch.max.bytes` per partition and `replica.fetch.response.max.bytes` overall. The follower appends them to its local log, advances its LEO, and the *next* fetch's offset tells the leader how far the follower has progressed — that's how the leader updates the follower's position and the high watermark. If the follower is fully caught up, the request blocks until new data arrives or the wait timeout elapses, then returns empty.

go deeper

for a junior

Know that followers pull data from the leader using fetch requests in a loop.

for a middle

Describe the ReplicaFetcherThread loop, fetchOffset-as-LEO, and long-poll wait behavior.

for a senior

Explain HW advancement from min ISR LEO, fetcher-to-leader grouping, and log truncation via leader epochs.

for a principal

Reason about fetcher contention, throttling during reassignment, and the implications of the pull model for leader statelessness and scalability.

**The pull model.** A crucial design choice: Kafka replication is **pull-based**, not push-based. The leader never proactively sends data to followers. Instead each follower *asks* the leader for data using the very same **Fetch** API that consumers use. This means the leader is largely stateless about replication progress — it learns each follower's position from the offsets in their fetch requests. **Who runs the loop.** On a broker acting as a follower, the `ReplicaManager` owns a `ReplicaFetcherManager`. The manager creates a pool of **`ReplicaFetcherThread`** objects. Partitions are assigned to threads by a hash of the partition over the configured fetcher count, *grouped by the leader broker* — i.e. one thread handles partitions whose leader is broker X, another handles partitions led by broker Y. So the number of fetcher threads a follower runs toward one leader is bounded by `num.replica.fetchers` (default 1). **The loop, step by step.** 1. The thread assembles a Fetch request. For each assigned partition it includes `fetchOffset = follower's current log-end offset (LEO)` — the offset of the next record it needs. 2. It sends the request to the leader broker and waits. The leader uses *long polling*: if there is no new data, it holds the request open up to **`replica.fetch.wait.max.ms`** (default 500 ms) or until at least `replica.fetch.min.bytes` (default 1) of data accumulates, then responds. 3. The leader returns records starting at `fetchOffset`, capped at **`replica.fetch.max.bytes`** (default ~1 MB) per partition and `replica.fetch.response.max.bytes` for the whole response. 4. The follower **appends** the returned records to its local log segment, advancing its LEO. 5. The arrival of the *next* fetch request — now carrying the higher `fetchOffset` — tells the leader the follower has persisted up to that point. The leader records this as the follower's LEO and uses the minimum LEO across the ISR to advance the **high watermark** (the highest offset considered committed/readable by consumers). **Why offsets-as-acknowledgment is elegant.** The follower's progress is implicit in its next request's start offset. No separate ack protocol is needed; the fetch request *is* the acknowledgment of everything below its `fetchOffset`. **Edge cases & nuances.** - **Catch-up vs steady state:** A lagging follower's fetches return full batches every time (bounded by the byte limits); a caught-up follower's fetches mostly block and return empty. - **Offset out of range / truncation:** If a follower's offset doesn't match the leader's log (e.g. after an unclean failover), the follower truncates its log to the leader's, using leader epochs (KIP-101/KIP-279) to find the correct divergence point, then resumes fetching. - **Throttling:** Replication can be rate-limited (`leader.replication.throttled.rate` / `follower...`), used during reassignment so fetchers don't saturate the network. - **One thread, many partitions:** Because a single thread multiplexes many partitions, a slow/large partition can delay others on the same thread — a reason to raise `num.replica.fetchers`.

  • How does the leader learn that a follower has persisted a given offset?
    From the follower's next Fetch request: its fetchOffset is the follower's new log-end offset, which implicitly acknowledges everything below it. The leader records that as the follower's LEO and advances the high watermark to the minimum LEO across the ISR.
  • What happens to the fetch request when the follower is fully caught up?
    The leader long-polls — it holds the request open up to replica.fetch.wait.max.ms (or until replica.fetch.min.bytes accumulates), then returns an empty response. This avoids tight busy-looping while keeping latency low.
  • What happens if a follower's fetch offset is out of range relative to the leader?
    The follower truncates its log to the leader's, using leader-epoch information to find the correct point of divergence (KIP-101/KIP-279), then resumes fetching from there.

saying these in an interview costs you the question

  • Saying the leader pushes data to followers — replication is pull-based
  • Claiming each partition has its own dedicated fetcher thread by default
  • Forgetting that the fetch offset itself is the acknowledgment
  • Confusing high watermark advance (min ISR LEO) with simple leader LEO

context