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.
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.
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
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();
}
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.
The bug is fixed, 300 messages wait in orders.dlq. A replay has to get four things right:
x-death, or the x-original-* headers after a republish.x-death is still there.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.
default-requeue-rejected: false, so a failure is rejected instead of looping.x-death.RepublishMessageRecoverer if you want the exception next to the message.