DEV Community

Cover image for Dead letters in Spring Boot: why RabbitMQ messages end up in the DLQ, and how to replay them
Viktor
Viktor

Posted on

Dead letters in Spring Boot: why RabbitMQ messages end up in the DLQ, and how to replay them

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
Enter fullscreen mode Exit fullscreen mode

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");
}
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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();
}
Enter fullscreen mode Exit fullscreen mode

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");
}
Enter fullscreen mode Exit fullscreen mode

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:

  1. 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.
  2. Send it where it came from, not to one queue for all: the exchange and routing key from x-death, or the x-original-* headers after a republish.
  3. Strip the bookkeeping. A consumer that gives up after three deaths rejects a replayed message on arrival if its old x-death is still there.
  4. 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;
    }
}
Enter fullscreen mode Exit fullscreen mode

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.
  • RepublishMessageRecoverer if 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)