Walk through how a produce request with acks=all flows through DelayedProduce and what triggers its completion.
answer
- append to leader log → maybe park
- watch each TopicPartition key
- HW advance ⇒ checkAndComplete
- all partitions' HW ≥ required offset ⇒ ack
- timeout ⇒ NOT_ENOUGH_REPLICAS / timeout
basics
~20 sThe leader writes the records to its log, then parks a DelayedProduce in purgatory keyed by each target partition. When followers replicate and the high watermark advances past the produced offsets for all partitions, the operation completes and the broker acknowledges the producer.
solid answer
~50 sWith acks=all, the leader first appends the batch to its own log, getting an offset. It can't ack yet — it must wait for all in-sync replicas (ISR) to replicate. So it builds a DelayedProduce holding the required offset per partition and the response callback, and registers it in the produce purgatory under each TopicPartition watcher key. The handler thread returns to the pool. Followers issue fetch requests; as they replicate and acknowledge, the leader advances each partition's high watermark (HW). On every HW advance the leader calls checkAndComplete on the produce purgatory for that partition, which re-runs tryComplete on parked DelayedProduce ops. Once the HW has reached the required offset for every partition in the request, tryComplete succeeds, onComplete sends a successful ack. If replication stalls and the operation's timer (request.timeout.ms-derived deadline) fires first, onExpiration completes it — typically with a NOT_ENOUGH_REPLICAS or timeout-style error, since not all ISR members caught up.
go deeper
Know acks=all makes the producer wait until replicas copy the data, and the broker parks the request meanwhile.
Trace append → park in purgatory → HW advances via follower fetches → complete; know the timeout path and min.insync.replicas.
Explain watcher keys per partition, exactly-once completion, ISR-shrink behavior, and leadership-change error handling.
Connect this to the durability contract and latency model; reason about the interaction of ISR dynamics, min.insync.replicas, and HW advancement under failure.
## Setup: what acks=all means `acks` is a producer config controlling durability: - `acks=0` — fire-and-forget, never enters purgatory. - `acks=1` — leader acks as soon as it writes to its own log; no waiting, no purgatory. - `acks=all` (a.k.a. `acks=-1`) — the leader must wait until **all in-sync replicas (ISR)** have replicated the records before acknowledging. This is the only produce mode that uses **`DelayedProduce`**. **ISR** = the set of replicas (including the leader) that are sufficiently caught up to the leader. The broker config `min.insync.replicas` sets the floor: if `acks=all` and the ISR is smaller than `min.insync.replicas`, the broker rejects the write with `NOT_ENOUGH_REPLICAS` rather than parking it. ## Step-by-step flow 1. **Append.** The leader appends the producer's batch to its local log, obtaining a base offset and a last offset for each partition in the request. 2. **Can we ack now?** The leader checks whether the **high watermark (HW)** for every target partition already covers the last offset written. The HW is the highest offset known to be replicated to all ISR members; consumers can only read up to it. If the HW already covers everything (e.g. ISR is just the leader, or followers were already caught up), the produce completes immediately — no purgatory. 3. **Park.** Otherwise the leader constructs a `DelayedProduce` carrying, per partition, the offset that must become 'committed' (i.e. HW must reach it), plus the response callback and an error map. It is inserted into the **produce `DelayedOperationPurgatory`**, watched under each partition's `TopicPartition` key, and also armed in the timer with a deadline derived from the client's `request.timeout.ms`. The handler thread is freed. 4. **Replication.** Follower brokers continuously send **fetch requests** to the leader to pull new records. As each follower advances its fetch position, the leader updates that follower's log-end-offset. When *all* ISR members have replicated past a given offset, the leader **advances the HW**. 5. **Trigger completion.** Advancing the HW for a partition causes the leader to call `replicaManager.tryCompleteDelayedProduce` / `purgatory.checkAndComplete(topicPartition)`. This walks the watcher list for that key and calls `tryComplete()` on each parked `DelayedProduce`. 6. **tryComplete logic.** The operation checks: for *every* partition in the request, has the HW reached the required offset (or has that partition errored, e.g. leadership moved)? If all partitions are satisfied or errored, it returns true and is completed exactly once via `onComplete()`, which sends the produce response (success, or partial errors per partition). 7. **Timeout path.** If replication stalls — a follower is slow or drops out of ISR — the timer deadline fires first. `onExpiration()` force-completes the operation. The producer sees a timeout/`NOT_ENOUGH_REPLICAS`-style outcome because not all required offsets were committed. ## Key edge cases - **Multi-partition produce.** One request can target many partitions; the single `DelayedProduce` watches *all* of them and only completes when the slowest partition's HW catches up (or it times out). It still completes exactly once. - **Leadership change.** If the leader loses a partition mid-wait, that partition's slot is marked with an error (e.g. `NOT_LEADER_OR_FOLLOWER`) and treated as 'resolved' so the op can complete and report the error rather than hang. - **ISR shrink.** If a follower falls out of the ISR, the HW can advance based on the *remaining* ISR, which may unblock the produce — acks=all means 'all currently in-sync replicas', not 'all replicas'. This is why `min.insync.replicas` exists: to forbid acking when the ISR has shrunk too far. - **Immediate completion.** Many acks=all produces never actually park — if followers are already caught up, tryComplete succeeds on the first check inline. ## Why this matters This is the mechanism behind Kafka's durability guarantee. The producer's latency under acks=all is essentially 'time until the slowest ISR follower replicates this batch', and the purgatory is what lets the broker wait for that without burning a thread per outstanding produce.
- Does acks=all wait for all replicas or all in-sync replicas, and how does min.insync.replicas interact?It waits for all replicas currently in the ISR, not the full replica set. min.insync.replicas sets a floor: if the ISR is below it, the broker rejects the write with NOT_ENOUGH_REPLICAS instead of parking the produce, preventing a silent durability downgrade when followers drop out.
- What advances the high watermark and thus triggers DelayedProduce completion?Follower fetch requests. As followers replicate and report their log-end-offsets, the leader advances the HW to the minimum offset replicated across all ISR members; each advance triggers checkAndComplete on the produce purgatory for that partition.
saying these in an interview costs you the question
- Saying acks=all waits for all replicas — it waits for the current ISR.
- Claiming the leader acks before any follower replicates (that's acks=1).
- Thinking the producer's own request.timeout.ms is unrelated — it derives the DelayedProduce deadline.
- Saying completion is purely timer-driven; the normal path is event-driven on HW advance.