What is the request-reply (RPC-style) messaging pattern in Kafka, and what extra pieces does it need on top of plain pub/sub?
answer
- reply topic + correlation ID
- REPLY_TOPIC and CORRELATION_ID headers
- ReplyingKafkaTemplate.sendAndReceive -> future
- pending-request map keyed by correlation ID
- RPC shape over async log
basics
~20 sRequest-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 sPlain 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
Know the one-line shape: send a request, get a reply back, use a correlation ID to match them.
Be able to name the REPLY_TOPIC and CORRELATION_ID headers and describe ReplyingKafkaTemplate.sendAndReceive returning a future.
Explain the full mechanism including the pending-request map and timeouts, and when to choose this over events.
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.