What is request purgatory in Kafka, and how does it relate to RemoteTimeMs for Fetch and Produce requests?
answer
- purgatory = park delayed requests, free the thread
- wait time → RemoteTimeMs
- Produce acks=all waits for ISR; Fetch waits for min.bytes/max.wait
- PurgatorySize gauge per delayedOperation
- hierarchical timing wheel, watcher lists
basics
~10 sPurgatory is where the broker parks requests that can't complete immediately and must wait for a condition (replication acks, or enough fetch data). That waiting time shows up as RemoteTimeMs.
solid answer
~40 sPurgatory is the broker's mechanism for handling delayed (asynchronous) operations — requests that can't be answered right away and must wait for a condition or a timeout. A Produce with acks=all parks in the produce purgatory until enough ISR followers replicate the records (or it times out). A Fetch parks in the fetch purgatory until fetch.min.bytes of data accumulates or fetch.max.wait.ms elapses (long-polling). The time spent waiting in purgatory is reported as RemoteTimeMs. You can observe purgatory depth via kafka.server:type=DelayedOperationPurgatory,name=PurgatorySize,delayedOperation=Produce|Fetch. Key insight: high RemoteTimeMs on Fetch is usually normal idling consumers, while high RemoteTimeMs on Produce indicates slow replication. Purgatory uses a hierarchical timing wheel so that watching huge numbers of delayed requests is cheap.
go deeper
Know purgatory is a waiting area for requests that can't finish immediately.
Connect produce acks=all and fetch min.bytes/max.wait to purgatory and to RemoteTimeMs.
Interpret produce vs fetch RemoteTimeMs differently and read PurgatorySize gauges.
Explain the timing-wheel/watcher-list design and set replication-health vs benign-longpoll alerting policy.
## What problem purgatory solves Many Kafka requests cannot be answered synchronously. Two big examples: - A **Produce with `acks=all`** must wait until all in-sync replicas (ISR) have replicated the records before the broker can acknowledge. - A **Fetch** with `fetch.min.bytes > 1` should wait (long-poll) until at least that many bytes are available, or until `fetch.max.wait.ms` elapses — so consumers aren't hammered with empty responses. If the broker held an I/O thread blocked for each such wait, it would run out of threads instantly. **Purgatory** decouples the wait from the thread: the request handler completes its synchronous work, then *parks* the request in a purgatory (a `DelayedOperationPurgatory`) and frees the thread. The request is completed later — either when its condition is satisfied (a follower's fetch advances the high-water mark; new data arrives) or when its timeout fires. ## Where the time shows up The wall-clock time a request spends parked in purgatory is recorded as **RemoteTimeMs** in RequestMetrics. "Remote" because the completion depends on something *other than* this broker's local processing — typically other brokers (replication) or future client activity (more produces filling a fetch). ## Two purgatories, two interpretations - **Produce purgatory** — high RemoteTimeMs here means ISR followers are slow to replicate: network issues, lagging brokers, or under-replicated partitions. This is a real latency problem. - **Fetch purgatory** — high RemoteTimeMs here is frequently *benign*: consumers configured with a fetch wait are simply long-polling an idle topic. Don't alert on it blindly. ## Observability - `kafka.server:type=DelayedOperationPurgatory,name=PurgatorySize,delayedOperation=Produce` (and `Fetch`) — number of requests currently parked. - `name=NumDelayedOperations` — similar gauge. There are also purgatories for other delayed operations (DeleteRecords, Heartbeat, Rebalance, etc.). ## Implementation detail Purgatory is built on a **hierarchical timing wheel**, an O(1) data structure for scheduling huge numbers of timeouts efficiently, plus a watcher-list keyed by the partitions a request waits on, so completion checks are triggered only by relevant events rather than scanning all parked requests. ## Edge cases - A request can complete *before* its timeout (condition met) or *at* timeout (forced completion returning whatever is available). - Large fetch purgatory size with low RemoteTimeMs simply means many idle long-polls — normal. - Producer with `acks=1` or `acks=0` skips produce-purgatory replication wait, so RemoteTimeMs stays near zero.
- Why does Kafka use purgatory instead of just blocking the request-handler thread until the condition is met?Blocking a thread per delayed request would exhaust the small I/O thread pool under load. Purgatory parks the request and frees the thread immediately, so a broker can have thousands of pending fetches/produces with only a handful of handler threads.
- A monitoring alert fires on high RemoteTimeMs for FetchConsumer. Is this necessarily a problem?No — it's usually normal long-polling: consumers waiting up to fetch.max.wait.ms for fetch.min.bytes of data on idle partitions. Only investigate if it correlates with consumer lag or throughput drops. For Produce, by contrast, high RemoteTimeMs does indicate slow replication.
saying these in an interview costs you the question
- Claiming purgatory blocks an I/O thread while waiting (it frees the thread).
- Treating high fetch RemoteTimeMs as always a fault rather than benign long-polling.
- Saying RemoteTimeMs is network send time — that's ResponseSendTimeMs.