skip to content

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%

answer

  1. ProducerFactory + reply-topic container
  2. UUID -> CORRELATION_ID header (16 bytes)
  3. ConcurrentHashMap correlationKey -> RequestReplyFuture
  4. replyTimeout default 5s -> KafkaReplyTimeoutException
  5. remove-then-complete; evict on timeout

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.

solid answer

~50 s

ReplyingKafkaTemplate<K,V,R> is constructed from a ProducerFactory (for sending requests) and a GenericMessageListenerContainer subscribed to the reply topic. On sendAndReceive(ProducerRecord): it generates a correlation ID (16-byte UUID by default), writes it into the KafkaHeaders.CORRELATION_ID header, sets KafkaHeaders.REPLY_TOPIC (and optionally REPLY_PARTITION) from the template's defaults or the record, creates a RequestReplyFuture, registers it in an internal futures map keyed by the correlation ID's CorrelationKey, then sends the record. The returned future also exposes getSendFuture() for the send-side ack. The reply container's listener, for each incoming reply record, extracts the correlation header, removes the matching future from the map, and completes it with the reply ConsumerRecord. If no reply arrives within the configured replyTimeout (default 5s), a scheduled task completes the future exceptionally with KafkaReplyTimeoutException and evicts the entry. The responder side is typically a @KafkaListener whose return value Spring routes back using the request's reply headers.

go deeper

for a junior

Know that sendAndReceive returns a future that completes when the reply arrives.

for a middle

Walk the correlation-id stamping, pending-map registration, reply-container completion, and timeout path.

for a senior

Discuss the shared-reply-topic mis-routing problem and REPLY_PARTITION / per-instance reply topics as fixes.

for a principal

Reason about at-least-once RPC semantics, idempotency of retried requests on timeout, and scatter-gather via AggregatingReplyingKafkaTemplate.

## Components `ReplyingKafkaTemplate<K, V, R>` (K=key, V=request value, R=reply value) extends `KafkaTemplate` and additionally **is** a `ConsumerSeekAware` message listener. You wire it from two things: - a **ProducerFactory** → used to send requests (the KafkaTemplate part); - a **listener container** (usually `KafkaMessageListenerContainer` built from a `ConsumerFactory`) subscribed to the **reply topic**. You pass this container into the `ReplyingKafkaTemplate` constructor, and the template registers *itself* as the container's message listener. ## sendAndReceive step by step 1. **Correlation ID.** The template generates a correlation value via its `CorrelationIdStrategy` (default: a random `UUID` serialized to 16 bytes) and writes it to the `KafkaHeaders.CORRELATION_ID` header on the outgoing record. 2. **Reply address.** It ensures `KafkaHeaders.REPLY_TOPIC` is present (from the record, or the template's `replyTopic`/default). It may also set `KafkaHeaders.REPLY_PARTITION` so replies land on a partition this instance consumes — important when scaling out callers (see below). 3. **Register the future.** It creates a `RequestReplyFuture<K,V,R>` and puts it into an internal `ConcurrentHashMap<CorrelationKey, RequestReplyFuture>`. The `CorrelationKey` wraps the correlation bytes with proper equals/hashCode. 4. **Send.** It sends the record. The returned `RequestReplyFuture` exposes `getSendFuture()` (the producer-send result) and is itself the receive future. 5. **Receive.** When the reply container delivers a record, the template's `onMessage` reads the correlation header, **removes** the future from the map, and completes it with the reply `ConsumerRecord`. (Removing first means a late duplicate reply is logged and dropped.) 6. **Timeout.** At registration the template schedules a timeout task (`replyTimeout`, default 5000 ms) that, if the future is still pending, completes it exceptionally with `KafkaReplyTimeoutException` and evicts the entry so the map does not leak. ## The responder side The service answering requests is normally a `@KafkaListener` method on the request topic that **returns** a value. Spring's `MessagingMessageListenerAdapter` sees a non-void return and, using a `KafkaTemplate` (the listener container factory's `replyTemplate`), sends the returned value to the topic named in the request's `REPLY_TOPIC` header, copying the `CORRELATION_ID` onto the reply. If you build the responder by hand you must copy the correlation header yourself. ## Scaling and the shared-reply-topic problem If multiple caller *instances* share one reply topic and one consumer group, replies are load-balanced across instances by partition — so an instance might receive a reply for a correlation ID it never sent, and the rightful owner never sees it. Fixes: (a) give each instance a **unique reply topic / unique group** so it consumes all of its own replies; or (b) use `ReplyingKafkaTemplate` with **embedded reply headers** + `REPLY_PARTITION` so each instance owns specific partitions; or (c) the `AggregatingReplyingKafkaTemplate` for scatter-gather. ## Reply timeout vs header timeout The per-send timeout can be overridden per request; some setups also propagate a `KafkaHeaders.REPLY_TOPIC`-side timeout. On timeout the future fails fast — the caller should treat it like any RPC timeout (retry carefully, since the request may still be processed: request-reply over Kafka is at-least-once, not exactly-once-RPC).

  • How do you stop one caller instance from stealing another instance's replies on a shared reply topic?
    Give each instance its own reply topic or unique consumer group so it consumes all of its own replies, or use REPLY_PARTITION so each instance owns specific partitions. Otherwise group load-balancing routes replies to the wrong instance and futures time out.
  • What is the default replyTimeout and what happens when it fires?
    5000 ms by default. A scheduled task completes the still-pending RequestReplyFuture exceptionally with KafkaReplyTimeoutException and removes it from the pending map so it does not leak.

saying these in an interview costs you the question

  • Claiming the reply container is optional — ReplyingKafkaTemplate needs a container on the reply topic to receive replies.
  • Saying late/duplicate replies complete the future twice; the future is removed from the map on first completion, so duplicates are dropped.
  • Assuming multiple instances can naively share one reply topic + group without replies being mis-routed.

context