skip to content

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%

answer

  1. @KafkaListener returns value -> auto reply
  2. factory.setReplyTemplate is required
  3. copy CORRELATION_ID verbatim; never hardcode reply topic
  4. honor REPLY_TOPIC and REPLY_PARTITION from request
  5. encode errors in reply, else caller just times out

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.

solid answer

~50 s

On the responder, a @KafkaListener consumes the request topic. The simplest correct implementation returns a value (or a Message) from the listener method; Spring's MessagingMessageListenerAdapter then sends that value to the topic named in the incoming KafkaHeaders.REPLY_TOPIC header and copies KafkaHeaders.CORRELATION_ID (and honors REPLY_PARTITION if present) onto the reply — for this you configure the listener container factory with a reply KafkaTemplate (setReplyTemplate). If you build it manually, you must read REPLY_TOPIC/CORRELATION_ID from the request headers and set them on the outgoing ProducerRecord yourself, or the caller can never match the reply. To return a different destination or wrap the payload, return a Message<?> and set headers explicitly. For errors, you can return a value carrying an error indicator, or configure error handling so the caller still gets a reply (otherwise it just times out). The contract: preserve the correlation ID exactly, and never invent your own reply topic — always use the one the request advertised.

go deeper

for a junior

Know the responder reads the request, computes an answer, and sends it back on the reply topic with the same correlation ID.

for a middle

Implement it: return a value from @KafkaListener with setReplyTemplate, or manually copy CORRELATION_ID and REPLY_TOPIC.

for a senior

Cover error replies vs silent timeouts, REPLY_PARTITION routing, and responder-side idempotency under at-least-once.

for a principal

Define the reply schema/error contract and idempotency strategy as platform standards so every responder behaves consistently.

## The responder's job The responder is the service that *answers* requests. It must: 1. **Consume** the request topic. 2. **Do the work** and produce a result. 3. **Send the result to the right place with the right correlation** so the caller's pending-request map resolves the correct future. ## The Spring 'just return a value' path (recommended) Declare a listener on the request topic whose method **returns** the reply: ``` @KafkaListener(topics = "requests", containerFactory = "rrFactory") @SendTo // optional; default reply destination falls back to REPLY_TOPIC header public PriceQuote handle(PriceRequest req) { return price(req); } ``` For this to route correctly: - The **listener container factory** must have a reply `KafkaTemplate` set (`factory.setReplyTemplate(kafkaTemplate)`), because the adapter uses it to send the reply. - Spring's `MessagingMessageListenerAdapter` reads the **incoming** `KafkaHeaders.REPLY_TOPIC` and sends the returned object there, automatically **copying `KafkaHeaders.CORRELATION_ID`** from request to reply (and `REPLY_PARTITION` if set). `@SendTo` can override the destination, but with no static target it uses the header — which is what request-reply needs since each caller may use a different reply topic/partition. ## The manual path If you don't return a value (e.g. you send asynchronously after some I/O), you must do the wiring yourself: read `record.headers()` for `REPLY_TOPIC` and `CORRELATION_ID`, build a `ProducerRecord` to that topic, and set the **same** correlation header on it. Forgetting to copy the correlation ID is the classic bug: the reply lands on the reply topic but matches no pending future, so the caller times out and the reply is dropped as an orphan. ## Headers you must respect - `KafkaHeaders.CORRELATION_ID` — **copy verbatim**; this is how the caller matches. - `KafkaHeaders.REPLY_TOPIC` — the destination the caller is listening on; **use it, don't hardcode**. - `KafkaHeaders.REPLY_PARTITION` — if present, send to that exact partition so the right caller instance receives it. ## Errors and the silent-timeout trap If the responder throws and produces no reply, the caller simply **times out** with no information. Better practice: catch and return a reply payload that encodes failure (an error code / problem object), so the caller distinguishes 'responder said no' from 'responder vanished'. Spring can also be configured so exceptions still produce a reply. Either way, design the reply schema to carry success/error explicitly. ## Idempotency on the responder Because delivery is at-least-once, the responder may see the **same request twice** (rebalance, retry). If handling has side effects, key on a business idempotency token from the request so reprocessing is safe; otherwise duplicate work and duplicate replies result. ## Summary contract Preserve the correlation ID exactly, route to the advertised reply topic/partition, encode errors in the reply rather than going silent, and make handling idempotent.

  • What is the most common bug when hand-rolling the responder?
    Not copying KafkaHeaders.CORRELATION_ID (and/or ignoring REPLY_TOPIC) onto the reply. The reply reaches the reply topic but matches no pending future in the caller, so it's dropped as an orphan and the caller times out.
  • What configuration on the listener container factory enables automatic replies in Spring?
    Setting a reply KafkaTemplate via factory.setReplyTemplate(...). The MessagingMessageListenerAdapter uses it to send the listener's return value to the request's REPLY_TOPIC, copying the correlation header automatically.

saying these in an interview costs you the question

  • Hardcoding a reply topic on the responder instead of reading REPLY_TOPIC from each request.
  • Forgetting to copy CORRELATION_ID, so every reply becomes an orphan and callers time out.
  • Letting the responder throw without replying, leaving the caller to time out with no error detail.
  • Assuming the responder sees each request exactly once (it's at-least-once; make handling idempotent).

context