skip to content

How do you correlate an asynchronous ack/nack back to the specific message that was sent, using CorrelationData?

level: seniorimportance: should knowfreq 42%

answer

  1. confirms async + out of order -> need identity
  2. CorrelationData(id) as last send arg
  3. ConfirmCallback receives same CorrelationData
  4. cd.getFuture() -> CompletableFuture<Confirm> (isAck/getReason)
  5. cd.getReturned() unifies unroutable check; getUnconfirmed for timeouts

basics

~20 s

Attach a CorrelationData (with a unique id) as an extra argument when you send. Because confirms are async and unordered, the ConfirmCallback receives that same CorrelationData, letting you look up which message the ack/nack belongs to. You can also await its CompletableFuture.

solid answer

~40 s

Confirms arrive asynchronously, batched, and not in send order, so a bare ack/nack is useless without knowing which message it refers to. With ConfirmType.CORRELATED you pass a CorrelationData carrying a unique id as the last argument to convertAndSend/send. The broker's confirm is then delivered to your ConfirmCallback with that exact CorrelationData, so you can update state (e.g. mark a DB row 'confirmed', or retry on nack). Even better, each CorrelationData exposes a CompletableFuture<CorrelationData.Confirm> via getFuture(): you can whenComplete/get on it per message, where Confirm.isAck() and getReason() tell the outcome. CorrelationData.getReturned() also lets that same object see whether the message was returned as unroutable, unifying confirm and return handling per message.

go deeper

for a junior

Know CorrelationData carries an id used to match the ack back to the message.

for a middle

Attach CorrelationData at send and read it in the ConfirmCallback to update state.

for a senior

Use the per-message CompletableFuture and getReturned(); handle amqp-thread concurrency and nack vs return.

for a principal

Architect an outbox/reconciliation with getUnconfirmed timeouts and durable per-id state surviving restarts.

## Why correlation is needed Publisher confirms are **asynchronous, batched, and can arrive out of order** relative to publishing. If your `ConfirmCallback` just gets `(ack=true)` with no identity, you cannot know *which* of the hundreds of in-flight messages it refers to. **`CorrelationData`** is the token that ties an outgoing message to its future confirm. ## Attaching CorrelationData Under `ConfirmType.CORRELATED`, you pass a `CorrelationData` as the final argument: ```java CorrelationData cd = new CorrelationData("order-42"); // unique id you choose rabbitTemplate.convertAndSend("orders.ex", "orders.created", payload, cd); ``` The `id` is yours to define — typically a message id, DB primary key, or UUID you can look up later. ## Two ways to consume the result **1. Template-level ConfirmCallback** — one callback for the whole template: ```java template.setConfirmCallback((correlationData, ack, cause) -> { String id = correlationData != null ? correlationData.getId() : null; if (ack) markConfirmed(id); else scheduleRetry(id, cause); // nack -> cause explains why }); ``` **2. Per-message CompletableFuture** — each `CorrelationData` owns a future: ```java cd.getFuture().whenComplete((confirm, ex) -> { if (ex != null) { /* future completed exceptionally, e.g. connection lost */ } else if (confirm.isAck()) { /* accepted */ } else { /* nack: confirm.getReason() */ } }); ``` `getFuture()` returns a `CompletableFuture<CorrelationData.Confirm>`; `Confirm` has `isAck()` and `getReason()`. This is cleaner for request-scoped code that wants to await one message's outcome without a global callback and lookup map. ## Unifying returns with the same object The same `CorrelationData` also captures an unroutable **return**: after a return, `cd.getReturned()` yields the `ReturnedMessage`. So a single object can answer both 'was it confirmed?' and 'was it returned as unroutable?', letting per-message logic treat 'acked but returned' correctly. ## Threading and state - Callbacks and future completions run on **amqp/listener threads**, not your publishing thread. Any shared state (maps of in-flight ids, DB updates) must be thread-safe. - A common pattern: persist the message with status PENDING keyed by the CorrelationData id, then flip to CONFIRMED/FAILED in the callback — this survives producer restarts better than an in-memory map. ## Timeouts and lost confirms — the hard edge case There is **no built-in timeout**: if the connection drops before a confirm arrives, the callback may never fire. Spring provides `rabbitTemplate.getUnconfirmed(long ageMillis)` to sweep CorrelationData objects older than an age so you can treat them as failed, and connection loss will complete outstanding futures exceptionally. You still must design a reconciliation/retry path; confirms give you the signal, not the recovery. ## Gotchas - Under `ConfirmType.NONE` or `SIMPLE`, the `ConfirmCallback`/future path won't deliver correlated results — you need `CORRELATED`. - A `null` CorrelationData is allowed but then you can't identify the message — always attach one if you care about the result. - Ids should be unique per in-flight message; reusing ids makes correlation ambiguous. ## When to use Use CorrelationData whenever you need per-message reliability bookkeeping — outbox/retry systems, exactly-tracked event publishing, or anything that reacts differently per message on nack vs return.

  • How can you await a single message's confirm without a global callback?
    Use the CorrelationData's own future: cd.getFuture() returns a CompletableFuture<CorrelationData.Confirm>; complete/await it and inspect isAck() and getReason().
  • A confirm never arrives because the connection dropped. How do you avoid hanging forever?
    There's no built-in timeout, so use rabbitTemplate.getUnconfirmed(ageMillis) to sweep old CorrelationData as failed, rely on the future completing exceptionally on connection loss, and build a reconciliation/retry path.

saying these in an interview costs you the question

  • Assuming confirms arrive in the same order as sends
  • Trying to correlate with ConfirmType.SIMPLE instead of CORRELATED
  • Storing in-flight state in a non-thread-safe map accessed from callbacks
  • Expecting a lost confirm to eventually time out on its own

context