skip to content

Request Purgatory and Delayed Operations

How the broker parks requests it cannot answer yet, such as acks=all produces and long-polling fetches, in the purgatory with its timing wheel. A deeper internals question that shows whether you know how long polling and durable acks are implemented.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

How do fetch.min.bytes and fetch.max.wait.ms drive DelayedFetch, and when does the broker decide to park a fetch versus answer it immediately?

level: middleimportance: must knowfreq 50%

answer

  1. min.bytes met now? answer immediately
  2. else park DelayedFetch, timer = max.wait.ms
  3. HW advance ⇒ recheck bytes
  4. timeout ⇒ return what's there (maybe empty)
  5. long-poll: latency vs batching

basics

~20 s

If a fetch can already return at least fetch.min.bytes of data, the broker replies immediately. If not, it parks a DelayedFetch and waits up to fetch.max.wait.ms for more data to accumulate; if that timer fires first, it returns whatever it has (possibly empty).

solid answer

~50 s

A consumer fetch carries fetch.min.bytes (default 1) and fetch.max.wait.ms (default 500). When the request arrives, the broker checks how many bytes of new data are available across the requested partitions. If that already meets fetch.min.bytes — or there's an error/partition state that means it should respond now — it answers immediately, no purgatory. Otherwise it builds a DelayedFetch, parks it in the fetch purgatory under each requested TopicPartition key, and arms a timer with fetch.max.wait.ms. As producers append and the partition's high watermark advances (for consumers, which read up to the HW), the broker calls checkAndComplete on those keys; tryComplete re-checks the accumulated bytes and completes once fetch.min.bytes is satisfied. If max.wait elapses first, onExpiration returns whatever is available — which may be an empty response. This is Kafka's long-poll: it trades a little latency for far better batching and lower request overhead.

go deeper

for a junior

Know the broker either answers a fetch right away or waits a bit (up to fetch.max.wait.ms) for enough data.

for a middle

Explain the min.bytes-met-now vs park decision, the timer, completion on HW advance, and that empty responses on idle topics are normal.

for a senior

Distinguish consumer fetch (reads to HW) vs replica fetch (reads to LEO), and how one HW-advance event completes both fetch and produce waiters.

for a principal

Reason about the latency/throughput trade-off of these knobs and how server-side long-poll beats client-side busy polling at scale.

## The two configs - **`fetch.min.bytes`** (consumer-side, default `1`): the minimum amount of data the broker should accumulate before answering this fetch. Raising it improves batching/throughput at the cost of latency. - **`fetch.max.wait.ms`** (consumer-side, default `500`): the maximum time the broker will hold the fetch waiting to reach `fetch.min.bytes`. It caps the added latency. Together they implement **long polling**: 'give me at least N bytes, but wait no longer than T ms'. ## Immediate vs parked decision When a fetch arrives, the broker reads, per requested partition, the available data between the consumer's requested offset and the offset it's allowed to read up to: - For ordinary consumers that's the **high watermark (HW)** — the highest fully-replicated offset. - For followers fetching (replica fetch) and read-committed/read-uncommitted semantics the readable boundary differs (HW vs log-end-offset vs last stable offset), but the purgatory machinery is the same. The broker sums available bytes across partitions. It responds **immediately** (no `DelayedFetch`) if any of these hold: 1. Accumulated bytes ≥ `fetch.min.bytes`. 2. `fetch.max.wait.ms <= 0` (caller wants no waiting). 3. A partition has an **error or changed state** (e.g. this broker is no longer the leader, offset out of range) that must be reported now. 4. The fetch reads from a partition whose readable data already satisfies the request. Otherwise it parks a **`DelayedFetch`**. ## The parked path 1. Build a `DelayedFetch` capturing the fetch metadata, the required minimum bytes, and the per-partition fetch positions. 2. Insert it into the **fetch `DelayedOperationPurgatory`**, watched under each requested `TopicPartition` key. 3. Arm the timer with `fetch.max.wait.ms`. Free the handler thread. 4. **Trigger:** when producers append new records and the partition's **HW advances**, the broker calls `checkAndComplete(topicPartition)` on the fetch purgatory. `tryComplete()` recomputes available bytes; once they meet `fetch.min.bytes`, the op completes and `onComplete()` reads the data and sends the response. 5. **Timeout:** if `fetch.max.wait.ms` elapses first, `onExpiration()` completes the op and returns whatever is available — possibly an **empty** fetch response (this is normal and expected when an idle topic has no new data). ## Edge cases and nuances - **Empty responses are normal.** On an idle partition, the consumer's fetch parks for `fetch.max.wait.ms` and returns empty; the consumer immediately re-fetches. This steady long-poll is by design, not a bug. - **fetch.max.bytes / max.partition.fetch.bytes** cap the *upper* size of a response; they don't affect the park/no-park decision (which is about the *minimum*). A single huge message larger than the limit is still returned to avoid stalls. - **HW advance, not raw append, is what matters for consumers.** Records appended but not yet above the HW (not yet committed) are not readable by ordinary consumers, so they don't satisfy a parked consumer fetch until the HW moves. - **Replica fetches** from followers also use `DelayedFetch` (with `replica.fetch.min.bytes` / `replica.fetch.wait.max.ms`), enabling efficient batched replication, but they read up to the leader's log-end-offset rather than the HW. - **Interaction with DelayedProduce.** The same HW-advance event can complete both a waiting consumer `DelayedFetch` (new readable data) and a `DelayedProduce` (offsets now committed) on the same partition key — one event, two watcher lists checked. ## Why long-poll instead of busy polling Without this, a consumer that wanted batching would have to poll repeatedly and sleep client-side, wasting requests and adding jitter. Server-side long-poll via purgatory lets the consumer issue one fetch and get woken precisely when enough data exists, with a bounded worst-case latency.

  • A consumer is fetching from a completely idle topic. What happens, and is an empty response a problem?
    The fetch parks as a DelayedFetch for fetch.max.wait.ms (default 500 ms) since no data accumulates, then returns an empty response on timeout. This is normal long-poll behavior — the consumer simply re-fetches; it is not an error.
  • Does raising fetch.min.bytes affect tail latency, and how is that bounded?
    Yes — a higher fetch.min.bytes makes the broker wait longer to accumulate data, increasing latency on low-traffic partitions. fetch.max.wait.ms bounds that wait, so the worst-case added latency for a fetch is fetch.max.wait.ms.

saying these in an interview costs you the question

  • Saying fetch.max.wait.ms is a minimum or a poll interval — it's the maximum the broker holds the request.
  • Thinking an empty fetch response means an error.
  • Confusing fetch.min.bytes (park decision) with fetch.max.bytes (response size cap).
  • Claiming consumers can read newly appended records before the HW advances.

context

open as a page

Walk through how a produce request with acks=all flows through DelayedProduce and what triggers its completion.

level: middleimportance: must knowfreq 48%

basics

~20 s

The 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.

open as a page

What is the request purgatory in a Kafka broker, and why does it exist?

level: juniorimportance: should knowfreq 35%

basics

~20 s

Purgatory is a broker-side holding area for requests that can't be answered immediately. Instead of blocking a thread, the broker parks the request there and completes it later when a condition is met or a timeout expires.

open as a page

Explain watcher keys and the checkAndComplete mechanism: how does a parked operation get woken, and why can an operation register under multiple keys?

level: seniorimportance: should knowfreq 32%

basics

~20 s

Each delayed operation is filed in the purgatory under one or more watcher keys (typically TopicPartition). When an event affecting a key occurs — usually an ISR/high-watermark advance — the broker calls checkAndComplete on that key, which re-runs tryComplete on every operation watching it. Multiple keys let one request span several partitions.

open as a page

Why does Kafka use a hierarchical timing wheel for delayed-operation timeouts instead of a java.util.concurrent DelayQueue or a priority queue?

level: principalimportance: nice to knowfreq 22%

basics

~20 s

A timing wheel adds and cancels a timeout in O(1), versus O(log n) for a heap/DelayQueue. Since brokers create and (usually) cancel a huge number of short-lived timeouts that complete early via events, O(1) insert/cancel matters far more than precise ordering.

open as a page