From the Warren blog: dead-letter operations for RabbitMQ.
A @RabbitListener throws. What happens next depends on three Spring AMQP settings most services never set on purpose: whether the message is requeued, whether it is retried, and what the last retry does with it. Get them wrong and one bad message loops forever or silently disappears. Get them right and it waits in a dead-letter queue until you have fixed the bug. Here is the whole chain, and how to put the messages back afterwards.
What Spring does when a listener throws
By default the listener container catches the exception and rejects the message with requeue. RabbitMQ puts it back at the head of the queue, the same consumer gets it again a millisecond later, throws again, and so on. A single poison message keeps a consumer busy, fills the log, and holds back everything behind it.
Two exceptions to that: a MessageConversionException (the payload could not be turned into your type) is treated as fatal by the default error handler and rejected without requeue. And if your code throws AmqpRejectAndDontRequeueException, the message is rejected without requeue too. For everything else, switch the default off:
spring:
rabbitmq:
listener:
simple:
default-requeue-rejected: false
Now a failed message is rejected with requeue=false. Whether that means "gone" or "dead-lettered" depends on the queue.
Give the queue a dead-letter route
A rejected message only reaches a dead-letter queue if its queue has a dead-letter exchange. In Spring you declare it with the queue:
@Bean
Queue orders() {
return QueueBuilder.durable("orders.process")
.deadLetterExchange("orders.dlx")
.deadLetterRoutingKey("orders.dlq")
.build();
}
@Bean
DirectExchange deadLetterExchange() {
return new DirectExchange("orders.dlx");
}
@Bean
Queue deadLetterQueue() {
return QueueBuilder.durable("orders.dlq").build();
}
@Bean
Binding deadLetterBinding() {
return BindingBuilder.bind(deadLetterQueue()).to(deadLetterExchange()).with("orders.dlq");
}
RabbitMQ then republishes every rejected, expired or overflowing message from orders.process to orders.dlx with routing key orders.dlq, and adds an x-death header that says which queue it came from, why, and its original exchange and routing key. (Reading x-death explains every field.)
One trap: queue arguments are fixed when a queue is created. If orders.process already exists without them, Spring's declaration fails with PRECONDITION_FAILED and the listener does not start. For existing queues use a policy instead, which can be added and changed at any time:
rabbitmqctl set_policy orders-dlx "^orders\.process$" \
'{"dead-letter-exchange":"orders.dlx","dead-letter-routing-key":"orders.dlq"}' \
--apply-to queues
Retry before giving up
Many failures are transient: a timeout, a lock, a service that restarts. Spring Boot retries inside the listener before it rejects:
spring:
rabbitmq:
listener:
simple:
default-requeue-rejected: false
retry:
enabled: true
max-attempts: 4
initial-interval: 1s
multiplier: 2
max-interval: 10s
Two things to know about this retry. It is stateless and in memory: the consumer thread sleeps between attempts while the message stays unacknowledged, and with a prefetch of 250 the other 249 wait too. Keep the total backoff short, seconds rather than minutes. And because the attempts never left the consumer, RabbitMQ does not see them: the message arrives in the DLQ with x-death count 1, although it failed four times.
For longer delays (a downstream system that is down for ten minutes), let the broker wait instead: a wait queue with a TTL whose dead-letter exchange points back at the work exchange. Then every cycle shows up in x-death, and the consumer can give up after a number of cycles by reading it:
@SuppressWarnings("unchecked")
static long deathsIn(Message message, String queue) {
var deaths = (List<Map<String, Object>>) message.getMessageProperties().getHeaders().get("x-death");
if (deaths == null) {
return 0;
}
return deaths.stream()
.filter(death -> queue.equals(String.valueOf(death.get("queue"))))
.mapToLong(death -> ((Number) death.get("count")).longValue())
.sum();
}
What the last attempt does: reject, or republish
When the retries are used up, a MessageRecoverer decides. The default rejects the message without requeue, so it is dead-lettered by the broker with x-death as described above. The alternative is RepublishMessageRecoverer:
@Bean
MessageRecoverer messageRecoverer(RabbitTemplate rabbitTemplate) {
// failed messages go to orders.error with the exception in the headers
return new RepublishMessageRecoverer(rabbitTemplate, "orders.error", "orders.process.failed");
}
It publishes a copy to an exchange of your choice and acknowledges the original. The copy carries x-exception-message, x-exception-stacktrace (trimmed to fit the frame size), x-original-exchange and x-original-routingKey. The broker never dead-lettered it, so there is no x-death.
| Reject (default) | RepublishMessageRecoverer | |
|---|---|---|
| Where it lands | The queue's dead-letter exchange | The exchange you name |
| Why it failed |
x-death reason rejected, no exception |
The exception message and stack trace |
| Where it came from |
x-death: queue, exchange, routing keys |
x-original-exchange, x-original-routingKey
|
| Needs topology | Dead-letter exchange on the queue | Error exchange and queue |
If you want to know why a message failed without searching the logs, republish. If you want the broker to stay in charge (TTL, length limits and quorum delivery limits dead-letter the same way), reject. Many teams do both: rejected messages of a TTL wait queue, republished ones from the listener.
Putting the messages back
The bug is fixed, 300 messages wait in orders.dlq. A replay has to get four things right:
- Acknowledge only after the broker confirmed the copy. Publish with confirms, wait, then ack the dead letter. Ack first and a broken connection loses the message.
-
Send it where it came from, not to one queue for all: the exchange and routing key from
x-death, or thex-original-*headers after a republish. -
Strip the bookkeeping. A consumer that gives up after three deaths rejects a replayed message on arrival if its old
x-deathis still there. - Go slowly if the consumer just recovered. Three hundred messages in a second can be the second incident.
With the client factory underneath Spring's and the plain channel API it looks like this; originalRoute reads the first x-death entry's exchange and routing key, or the x-original-* pair:
/** Moves up to [limit] messages from the DLQ back to where they came from. */
int replay(CachingConnectionFactory connectionFactory, String dlq, int limit) throws Exception {
// a connection of its own from the underlying client factory: Spring's factory caches channels,
// and confirm mode must not leak into that cache
try (var connection = connectionFactory.getRabbitConnectionFactory().newConnection("dlq-replay");
var channel = connection.createChannel()) {
channel.confirmSelect();
int moved = 0;
while (moved < limit) {
GetResponse response = channel.basicGet(dlq, false);
if (response == null) {
break; // the queue is empty
}
var props = response.getProps();
Map<String, Object> headers = new HashMap<>(props.getHeaders() == null ? Map.of() : props.getHeaders());
Route route = originalRoute(headers); // from x-death, or x-original-* after a republish
headers.keySet().removeIf(name ->
name.startsWith("x-death") || name.startsWith("x-first-death") ||
name.startsWith("x-last-death") || name.startsWith("x-exception") ||
name.startsWith("x-original"));
channel.basicPublish(route.exchange(), route.routingKey(), true,
props.builder().headers(headers).build(), response.getBody());
channel.waitForConfirmsOrDie(5_000); // the broker has it now
channel.basicAck(response.getEnvelope().getDeliveryTag(), false); // only then remove it here
moved++;
}
return moved;
}
}
Not in this sketch, and worth adding before you run it in production: a return listener for unroutable messages (mandatory=true sends them back instead of dropping them), a rate limit, selecting messages by exception rather than taking the first N, and a record of what was moved where. Do not use RabbitTemplate.receive() for this: it acknowledges on receipt, before anything has been published.
Checklist
-
default-requeue-rejected: false, so a failure is rejected instead of looping. - A dead-letter exchange on every work queue, by argument for new queues, by policy for existing ones.
- Short in-memory retries; long delays through a TTL wait queue that shows up in
x-death. -
RepublishMessageRecovererif you want the exception next to the message. - Replays that confirm before they ack, route back to the original exchange, and strip the death headers.
Where Warren fits
Warren reads both conventions: x-death from broker dead-lettering and the x-exception-* and x-original-* headers of Spring's RepublishMessageRecoverer. It groups a dead-letter queue by exception, shows each message's story, and replays the ones you select to their original exchange and routing key, death and exception headers stripped, confirmed message by message and throttled if you want. Every replay is in the audit log. Free for one cluster.
Try Warren in a minute (one compose file, demo broker with real dead letters included) · README on GitHub · 4-minute walkthrough on YouTube
Top comments (0)