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?
answer
- @KafkaListener returns value -> auto reply
- factory.setReplyTemplate is required
- copy CORRELATION_ID verbatim; never hardcode reply topic
- honor REPLY_TOPIC and REPLY_PARTITION from request
- encode errors in reply, else caller just times out
basics
~20 sThe 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 sOn 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
Know the responder reads the request, computes an answer, and sends it back on the reply topic with the same correlation ID.
Implement it: return a value from @KafkaListener with setReplyTemplate, or manually copy CORRELATION_ID and REPLY_TOPIC.
Cover error replies vs silent timeouts, REPLY_PARTITION routing, and responder-side idempotency under at-least-once.
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).