Explain watcher keys and the checkAndComplete mechanism: how does a parked operation get woken, and why can an operation register under multiple keys?
answer
- purgatory = Map<key, watcher list> + timer
- tryCompleteElseWatch: inline try, else file under all keys
- HW/ISR advance ⇒ checkAndComplete(key)
- AtomicBoolean completed = exactly once
- reaper/purge.interval.requests sweeps stale entries
basics
~20 sEach 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.
solid answer
~50 sThe purgatory is essentially a map from a watcher key to a list of parked DelayedOperations. A watcher key is whatever the operation depends on — for produce and fetch it's a TopicPartition; for group operations it's a member/group key. When the broker does something that might satisfy waiters — the most important being advancing a partition's high watermark or changing its ISR after follower replication — it calls purgatory.checkAndComplete(key). That looks up the watcher list for that key and invokes tryComplete on each operation; those whose conditions are now met are force-completed exactly once and unlinked. An operation registers under multiple keys when it depends on several partitions: a multi-partition produce or fetch watches every TopicPartition it touches, so a triggering event on any of them re-checks it, but it only completes when all its partitions are satisfied. A separate timer guarantees completion on timeout regardless of events, and a reaper periodically purges completed-but-still-linked entries.
go deeper
Know operations are filed under a key and re-checked when something happens on that key.
Explain watcher key = TopicPartition, checkAndComplete re-runs tryComplete, and a request can span multiple partition keys.
Detail tryCompleteElseWatch, the exactly-once CAS guard against timer/event races, and the reaper/purge.interval.requests hygiene.
Reason about the data-structure design (map + timing wheel), speculative checkAndComplete cost, and observability via per-purgatory JMX metrics.
## What a watcher key is The **`DelayedOperationPurgatory`** maintains, conceptually, a `Map<Key, Watchers>` where `Watchers` is a linked list of parked `DelayedOperation`s plus the same operations are also held in a **timer** keyed by their deadlines. A **watcher key** identifies the external condition an operation is waiting on. For the operations in this topic: - `DelayedProduce` and `DelayedFetch` use a **`TopicPartition`** as the key — they wait on that partition's data/replication state. - Other op types use other keys (e.g. group/member identifiers for `DelayedHeartbeat`/`DelayedJoin`). ## Registration When the broker parks an operation it calls `purgatory.tryCompleteElseWatch(op, keys)`: 1. It first runs `tryComplete()` once inline — if the condition is already met, it completes immediately and never actually parks (a common fast path). 2. Otherwise it adds the operation to the **watcher list of every key** in `keys`, and arms it in the timer with its deadline. Because the same operation object lives in several watcher lists, a multi-partition request is woken by activity on *any* of its partitions, yet its `tryComplete()` still requires *all* partitions to be satisfied before returning true. ## Waking via checkAndComplete The broker calls **`purgatory.checkAndComplete(key)`** whenever an event might satisfy operations watching that key. The dominant trigger for produce/fetch is the **ReplicaManager advancing the high watermark (HW)** of a partition after follower fetches / ISR changes. `checkAndComplete`: 1. Looks up the watcher list for `key`. 2. Calls `tryComplete()` on each parked op. 3. For each that returns true, it **force-completes** the op (`forceComplete()` flips an atomic `completed` flag with compare-and-set, so it runs exactly once even if a timeout races), runs `onComplete()`, and removes it from *all* its watcher lists and the timer. Operations that aren't ready stay parked; `tryComplete()` may be invoked many times, so it must be cheap and idempotent. ## Exactly-once completion under races An operation can be raced by (a) an event-driven `checkAndComplete` and (b) its timer expiring. `DelayedOperation` guards this with an `AtomicBoolean completed`; whichever path wins the CAS runs `onComplete()`/`onExpiration()`, the other becomes a no-op. This prevents double responses to the client. ## Why multiple keys (concrete) A producer sends one request writing to partitions P0, P1, P2 with `acks=all`. The single `DelayedProduce` is filed under keys P0, P1, P2. Suppose P0's HW advances first → `checkAndComplete(P0)` runs `tryComplete`, which sees P1 and P2 still uncommitted and returns false. Later P1 and then P2 advance; the final `checkAndComplete(P2)` makes `tryComplete` return true and the op completes once. The same applies to a multi-partition consumer `DelayedFetch` accumulating `fetch.min.bytes` across partitions. ## Memory hygiene: the reaper / purge When an op completes via one key's `checkAndComplete`, it's removed from the lists it can reach, but stale references can linger in *other* watcher lists until they're walked. A background **reaper (expiration) thread** plus periodic purging (triggered by `purge.interval.requests`, default 1000 requests) sweeps already-completed operations out of watcher lists and the timer so memory doesn't grow unbounded under high churn. The purgatory also exposes JMX metrics like `PurgatorySize` and `NumDelayedOperations` per purgatory (Produce, Fetch, etc.) for observability. ## Edge cases - A partition leadership change marks that partition's slot resolved-with-error inside `tryComplete`, letting multi-partition ops complete and report the error rather than hang. - If no event ever fires (e.g. a fetch on a truly idle topic), only the timer completes the op — the watcher-key path is never exercised for it. - `checkAndComplete` is called speculatively on many events; most calls find nothing to complete, so it's designed to be cheap when watcher lists are short.
- If an operation watches three partition keys and partition P0's HW advances, does it complete?Only if tryComplete now finds all three partitions satisfied (or errored). A checkAndComplete on P0 re-runs tryComplete, but the op stays parked until P1 and P2 are also satisfied; it completes exactly once, on whichever event finally makes all conditions true.
- How does Kafka guarantee an operation isn't completed twice when a timeout and an HW-advance event race?DelayedOperation holds an AtomicBoolean 'completed'; forceComplete uses compareAndSet, so exactly one of the event path or the timer path wins and runs onComplete/onExpiration. The loser is a no-op, preventing a double response to the client.
saying these in an interview costs you the question
- Saying an operation is only ever woken by its timer — the normal path is event-driven via checkAndComplete on a watcher key.
- Claiming a multi-partition op completes as soon as any one watched key fires.
- Forgetting the exactly-once completion guard (AtomicBoolean/CAS), which risks describing double client responses.
- Ignoring the reaper/purge, implying completed ops are freed instantly from every watcher list.