skip to content

Request-Reply and Correlation Patterns

Doing request/reply over Kafka with correlation IDs and reply topics, and knowing when you should not. Interviewers use it to check that you can argue against RPC over a log.

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

questions

5

What is the request-reply (RPC-style) messaging pattern in Kafka, and what extra pieces does it need on top of plain pub/sub?

level: juniorimportance: must knowfreq 55%

answer

  1. reply topic + correlation ID
  2. REPLY_TOPIC and CORRELATION_ID headers
  3. ReplyingKafkaTemplate.sendAndReceive -> future
  4. pending-request map keyed by correlation ID
  5. RPC shape over async log

basics

~20 s

Request-reply makes one service send a message and wait for a matching response. On top of pub/sub you add a reply topic to send the answer back, and a correlation ID so the requester knows which reply belongs to which request.

solid answer

~40 s

Plain Kafka pub/sub is one-way: a producer writes to a topic, consumers read it, and there is no built-in notion of a 'response'. Request-reply layers an RPC shape on top: the requester sends a record to a request topic and then waits for a reply. To make that work you need (1) a reply topic the responder writes the answer to, advertised on each request via a KafkaHeaders.REPLY_TOPIC header, and (2) a correlation ID (KafkaHeaders.CORRELATION_ID) copied from the request onto the reply so the requester can match reply to request. Spring Kafka's ReplyingKafkaTemplate implements this: sendAndReceive() returns a RequestReplyFuture, stores the pending request in an in-memory map keyed by correlation ID, and completes the future when a reply with that ID arrives on the reply topic.

go deeper

for a junior

Know the one-line shape: send a request, get a reply back, use a correlation ID to match them.

for a middle

Be able to name the REPLY_TOPIC and CORRELATION_ID headers and describe ReplyingKafkaTemplate.sendAndReceive returning a future.

for a senior

Explain the full mechanism including the pending-request map and timeouts, and when to choose this over events.

for a principal

Frame the trade-off: request-reply re-introduces temporal coupling Kafka was built to avoid; justify or reject it per use case and design the topic/correlation contract.

## The problem Kafka is fundamentally a **publish-subscribe log**: a *producer* appends *records* (messages) to a *topic* (a named, partitioned, append-only stream), and *consumers* read records from that topic at their own pace. There is no built-in concept of a 'response' — communication is one-directional and asynchronous (the producer does not wait for anyone to read). **Request-reply** (also called *RPC-style* messaging — Remote Procedure Call, i.e. 'call a function on another service and get a value back') is a pattern that simulates a synchronous call over this async substrate. Service A wants to ask Service B a question and get an answer back. ## What you must add on top of pub/sub 1. **Two topics (or two roles):** a *request topic* that the caller writes to and the responder consumes, and a *reply topic* that the responder writes to and the original caller consumes. The caller tells the responder where to reply by putting the reply-topic name in a header on each request — in Spring Kafka this is the `KafkaHeaders.REPLY_TOPIC` header. 2. **A correlation ID.** Because many requests flow concurrently and Kafka delivers records in bulk, when a reply lands on the reply topic the caller must know *which* outstanding request it answers. The caller stamps each request with a unique **correlation ID** (`KafkaHeaders.CORRELATION_ID`, typically a UUID's bytes). The responder must **copy that exact ID** onto the reply record's header. The caller keeps a **pending-request map** (correlation ID → a future/promise) and, when a reply arrives, looks up the ID and completes the matching future. 3. **A timeout.** Since the responder might be down or slow, the caller cannot wait forever. Each pending request has a timeout; on expiry the future is completed exceptionally (e.g. `KafkaReplyTimeoutException` in Spring) and the entry is evicted from the map. ## In Spring Kafka concretely `ReplyingKafkaTemplate<K,V,R>` wraps a producer plus a `KafkaMessageListenerContainer` on the reply topic. You call `sendAndReceive(record)`; it sets the correlation header, registers the future in an internal `ConcurrentHashMap`, sends the request, and returns a `RequestReplyFuture`. Its listener side receives replies, reads the correlation ID, and completes the corresponding future. The responder side is usually a `@KafkaListener` method that `return`s a value, which Spring routes back to the reply topic carried in the request headers. ## Why not just use this everywhere? Request-reply re-introduces *temporal coupling* (caller blocks waiting) and operational complexity (reply topics, correlation, timeouts) that pub/sub was designed to remove. It is appropriate when a caller genuinely needs the result before proceeding; otherwise plain async events are simpler and more resilient.

  • Which two Kafka headers carry the request-reply wiring in Spring Kafka?
    KafkaHeaders.REPLY_TOPIC (where to send the answer) and KafkaHeaders.CORRELATION_ID (to match reply to request). KafkaHeaders.REPLY_PARTITION can also be set to pin the reply to a specific partition.
  • Where does the correlation ID get matched back to the caller?
    In the caller's in-memory pending-request map: ReplyingKafkaTemplate stores correlationId -> RequestReplyFuture, and its reply-topic listener completes the future when a reply carrying that ID arrives.

saying these in an interview costs you the question

  • Saying Kafka has native request-reply / RPC support (it does not — it must be layered on).
  • Forgetting the correlation ID and assuming reply order matches request order (Kafka gives no such cross-topic ordering guarantee).
  • Thinking the responder invents its own reply topic instead of reading REPLY_TOPIC from the request.

context

open as a page

When should you choose Kafka request-reply (RPC-style) over plain pub/sub event choreography, and what are the architectural trade-offs?

level: seniorimportance: must knowfreq 50%

basics

~20 s

Use request-reply only when the caller truly needs the answer before it can proceed and you specifically want it to flow over Kafka. It re-adds temporal coupling, blocking, reply topics, correlation, and timeouts. Prefer pub/sub events when the caller can react later — it's looser, more resilient, and replayable. Often a synchronous HTTP/gRPC call is a better RPC than Kafka.

open as a page

Walk through how ReplyingKafkaTemplate.sendAndReceive() works end to end, including how the reply container, correlation, and the returned future fit together.

level: middleimportance: should knowfreq 45%

basics

~20 s

You build ReplyingKafkaTemplate from a producer factory plus a listener container on the reply topic. sendAndReceive() stamps a correlation ID, sets the reply-topic header, stores a future in a map keyed by that ID, sends the request, and returns the future. When the reply container receives a record with that ID, it completes the future.

open as a page

How do you implement the responder side of a Kafka request-reply exchange so replies are correctly correlated and routed, and what must you preserve from the request?

level: middleimportance: should knowfreq 38%

basics

~20 s

The responder consumes the request topic, does its work, and sends the result to the reply topic named in the request's REPLY_TOPIC header, copying the request's CORRELATION_ID onto the reply. In Spring, a @KafkaListener method that returns a value does this automatically via the container factory's reply template.

open as a page

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?

level: seniorimportance: should knowfreq 35%

basics

~20 s

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

open as a page