What failure modes does the in-memory pending-request map introduce in a Kafka request-reply client, and how do you keep it bounded and correct?
answer
- map in heap = leak + growth + orphan + no durability
- every entry needs an eviction path (timeout)
- no Kafka backpressure -> caller bulkhead/semaphore
- timeout = 'unknown outcome', make requests idempotent
- crash loses futures; replies strand on reply topic
basics
~20 sThe pending map holds a future per outstanding request in JVM memory. Risks: entries leak if replies never arrive and there's no timeout, the map grows under load, late or duplicate replies have no owner after eviction, and futures are lost on a process crash. Bound it with per-request timeouts, eviction on completion/timeout, and treat lost requests as RPC timeouts.
solid answer
~50 sThe map (correlationId -> future) lives only in the caller's heap, which creates several issues. (1) Leaks: without a timeout, any request whose reply is lost keeps its entry forever; ReplyingKafkaTemplate avoids this by scheduling a timeout that completes the future with KafkaReplyTimeoutException and removes the entry. (2) Unbounded growth: a burst of requests or a slow responder inflates the map and pending futures; you cap concurrency upstream (bulkhead/semaphore) since Kafka itself won't backpressure the caller. (3) Orphan replies: a reply arriving after eviction has no future to complete and is logged/dropped — so late replies are silently lost. (4) No durability: the map is in-memory, so a caller crash loses every pending future; on restart the responder's replies land on the reply topic with no waiter, and the original callers see timeouts. Because of (3) and (4), Kafka request-reply is at-least-once at the messaging layer but no-stronger-than at-most-once-success at the RPC layer; callers must make requests idempotent and treat timeouts as 'unknown outcome', not 'failed'.
go deeper
Just know there's a map of waiting requests and each needs a timeout so it doesn't wait forever.
Explain leaks without timeouts and that the map is in memory, lost on crash.
Cover all four failure modes plus the correctness rules: idempotent requests, timeout=unknown, caller bulkhead, unique reply consumption.
Tie it to the architectural cost of synchronous-over-async, define SLOs/metrics on pending size and timeout rate, and decide when reconciliation beats RPC entirely.
## What the pending-request map is In the request-reply pattern the caller cannot block a thread per request (that would not scale), so instead it keeps an **in-memory map** from **correlation ID → pending future** (`RequestReplyFuture`). The calling thread either blocks on `future.get(timeout)` or attaches a callback. The map is just a `ConcurrentHashMap` on the JVM heap of the caller process. Treating it as a real source of truth causes several failure modes: ### 1. Memory leak without timeouts If a reply never arrives (responder down, reply topic mis-config, network loss) and there is **no timeout**, the entry — and the future, plus whatever the calling code holds — stays in the map indefinitely. Over time this is a classic slow heap leak. `ReplyingKafkaTemplate` defends against this by scheduling, at registration time, a timeout task (`replyTimeout`) that completes the future with `KafkaReplyTimeoutException` and **removes** the entry. The lesson: *every* pending entry must have a guaranteed eviction path. ### 2. Unbounded growth under load Even with timeouts, a flood of requests against a slow responder makes the map (and the number of in-flight futures) balloon until each times out. Kafka does **not** backpressure the caller — producing to the request topic is fire-and-fast. So the caller must impose its own **bulkhead**: a bounded semaphore / concurrency limiter, a queue with rejection, or a circuit breaker that trips when the timeout rate spikes. Otherwise you convert a slow responder into a caller-side OOM. ### 3. Orphan / late / duplicate replies Once a future is completed (by reply or by timeout) the entry is removed. A reply that arrives **after** eviction — a slow responder that beat its own timeout, or a duplicate caused by at-least-once delivery / a responder retry — finds **no matching future** and is simply logged and dropped. Consequences: (a) a late successful reply is lost even though work was done; (b) you must not assume one reply per request. Design responders to be idempotent and tolerate the caller having moved on. ### 4. No durability across crashes The map is purely in heap. If the **caller crashes** while requests are outstanding: all pending futures vanish; the responder still processes the requests and writes replies to the reply topic; those replies have **no waiter** (or, after restart with a fresh consumer offset, may not even be read). The original business callers get timeouts/errors. There is **no exactly-once RPC** here — the request may have executed even though the caller never learned the result. ## Correctness rules that follow - **Idempotent requests:** include a business idempotency key so a retried request after a timeout does not double-apply. - **Timeouts mean 'unknown', not 'failed':** the work may have completed. Reconcile rather than blindly retry side-effecting operations. - **Bound concurrency** at the caller (semaphore/circuit breaker) because Kafka offers no caller backpressure. - **Unique reply consumption** per instance (own reply topic / group / partition) so a restarted or scaled instance does not strand replies. - **Observe** map size and timeout rate as first-class metrics; rising pending count is an early warning of responder degradation. ## Why this matters architecturally These failure modes are exactly the costs you accept when bolting synchronous RPC onto an async log. They are the concrete reason senior engineers prefer event choreography unless a true request/response is required.
- After a request-reply timeout, is it safe to just retry the request?Not blindly. A timeout means the outcome is unknown — the responder may have already done the work and replied late/after eviction. Retrying a side-effecting request can double-apply, so requests need a business idempotency key (or you reconcile state) before retry.
- Why can the pending map grow unbounded, and what prevents it?Kafka does not backpressure the producer, so a slow responder lets in-flight requests pile up faster than they time out. Prevent it with a caller-side bulkhead: a bounded semaphore/queue or a circuit breaker that trips on rising timeout rate, plus capping in-flight requests.
saying these in an interview costs you the question
- Calling Kafka request-reply 'exactly-once RPC' — the in-memory map and at-least-once delivery preclude that.
- Assuming a timeout means the request was not processed (it may have been; outcome is unknown).
- Relying on Kafka to backpressure the caller when the responder is slow — it does not.
- Ignoring late/duplicate replies; after eviction they are silently dropped.