Fixes: https://github.com/spring-cloud/spring-cloud-stream/issues/2939 The RabbitMQ 4.0 does not deal with client side `x-*` headers. Therefore, an `x-death.count` is not incremented anymore when message is re-published from client back to the broker. * Spring AMQP 3.2 has introduced an `AmqpHeaders.RETRY_COUNT` custom header. Use `messageProperties.incrementRetryCount()` in the `RabbitMessageChannelBinder` when we re-published message back to the broker for server-side retries * Fix docs respectively Resolves #3019
229 lines
7.5 KiB
Plaintext
229 lines
7.5 KiB
Plaintext
[[rabbit-dlq-processing]]
|
|
= Dead-Letter Queue Processing
|
|
|
|
Because you cannot anticipate how users would want to dispose of dead-lettered messages, the framework does not provide any standard mechanism to handle them.
|
|
If the reason for the dead-lettering is transient, you may wish to route the messages back to the original queue.
|
|
However, if the problem is a permanent issue, that could cause an infinite loop.
|
|
The following Spring Boot application shows an example of how to route those messages back to the original queue but moves them to a third "`parking lot`" queue after three attempts.
|
|
The second example uses the https://www.rabbitmq.com/blog/2015/04/16/scheduling-messages-with-rabbitmq/[RabbitMQ Delayed Message Exchange] to introduce a delay to the re-queued message.
|
|
In this example, the delay increases for each attempt.
|
|
These examples use a `@RabbitListener` to receive messages from the DLQ.
|
|
You could also use `RabbitTemplate.receive()` in a batch process.
|
|
|
|
The examples assume the original destination is `so8400in` and the consumer group is `so8400`.
|
|
|
|
[[non-partitioned-destinations]]
|
|
== Non-Partitioned Destinations
|
|
|
|
The first two examples are for when the destination is *not* partitioned:
|
|
|
|
[source, java]
|
|
----
|
|
@SpringBootApplication
|
|
public class ReRouteDlqApplication {
|
|
|
|
private static final String ORIGINAL_QUEUE = "so8400in.so8400";
|
|
|
|
private static final String DLQ = ORIGINAL_QUEUE + ".dlq";
|
|
|
|
private static final String PARKING_LOT = ORIGINAL_QUEUE + ".parkingLot";
|
|
|
|
public static void main(String[] args) throws Exception {
|
|
ConfigurableApplicationContext context = SpringApplication.run(ReRouteDlqApplication.class, args);
|
|
System.out.println("Press enter to exit");
|
|
System.in.read();
|
|
context.close();
|
|
}
|
|
|
|
@Autowired
|
|
private RabbitTemplate rabbitTemplate;
|
|
|
|
@RabbitListener(queues = DLQ)
|
|
public void rePublish(Message failedMessage) {
|
|
long retries = failedMessage.getMessageProperties().getRetryCount();
|
|
if (retries < 3) {
|
|
failedMessage.getMessageProperties().incrementRetryCount();
|
|
this.rabbitTemplate.send(ORIGINAL_QUEUE, failedMessage);
|
|
}
|
|
else {
|
|
this.rabbitTemplate.send(PARKING_LOT, failedMessage);
|
|
}
|
|
}
|
|
|
|
@Bean
|
|
public Queue parkingLot() {
|
|
return new Queue(PARKING_LOT);
|
|
}
|
|
|
|
}
|
|
----
|
|
|
|
[source, java]
|
|
----
|
|
@SpringBootApplication
|
|
public class ReRouteDlqApplication {
|
|
|
|
private static final String ORIGINAL_QUEUE = "so8400in.so8400";
|
|
|
|
private static final String DLQ = ORIGINAL_QUEUE + ".dlq";
|
|
|
|
private static final String PARKING_LOT = ORIGINAL_QUEUE + ".parkingLot";
|
|
|
|
private static final String DELAY_EXCHANGE = "dlqReRouter";
|
|
|
|
public static void main(String[] args) throws Exception {
|
|
ConfigurableApplicationContext context = SpringApplication.run(ReRouteDlqApplication.class, args);
|
|
System.out.println("Press enter to exit");
|
|
System.in.read();
|
|
context.close();
|
|
}
|
|
|
|
@Autowired
|
|
private RabbitTemplate rabbitTemplate;
|
|
|
|
@RabbitListener(queues = DLQ)
|
|
public void rePublish(Message failedMessage) {
|
|
long retries = failedMessage.getMessageProperties().getRetryCount();
|
|
if (retries < 3) {
|
|
failedMessage.getMessageProperties().incrementRetryCount();
|
|
Map<String, Object> headers = failedMessage.getMessageProperties().getHeaders();
|
|
headers.put("x-delay", 5000 * retriesHeader);
|
|
this.rabbitTemplate.send(DELAY_EXCHANGE, ORIGINAL_QUEUE, failedMessage);
|
|
}
|
|
else {
|
|
this.rabbitTemplate.send(PARKING_LOT, failedMessage);
|
|
}
|
|
}
|
|
|
|
@Bean
|
|
public DirectExchange delayExchange() {
|
|
DirectExchange exchange = new DirectExchange(DELAY_EXCHANGE);
|
|
exchange.setDelayed(true);
|
|
return exchange;
|
|
}
|
|
|
|
@Bean
|
|
public Binding bindOriginalToDelay() {
|
|
return BindingBuilder.bind(new Queue(ORIGINAL_QUEUE)).to(delayExchange()).with(ORIGINAL_QUEUE);
|
|
}
|
|
|
|
@Bean
|
|
public Queue parkingLot() {
|
|
return new Queue(PARKING_LOT);
|
|
}
|
|
|
|
}
|
|
----
|
|
|
|
[[partitioned-destinations]]
|
|
== Partitioned Destinations
|
|
|
|
With partitioned destinations, there is one DLQ for all partitions.
|
|
We determine the original queue from the headers.
|
|
|
|
[[republishtodlq-false]]
|
|
=== `republishToDlq=false`
|
|
|
|
When `republishToDlq` is `false`, RabbitMQ publishes the message to the DLX/DLQ with an `x-death` header containing information about the original destination, as shown in the following example:
|
|
|
|
[source, java]
|
|
----
|
|
@SpringBootApplication
|
|
public class ReRouteDlqApplication {
|
|
|
|
private static final String ORIGINAL_QUEUE = "so8400in.so8400";
|
|
|
|
private static final String DLQ = ORIGINAL_QUEUE + ".dlq";
|
|
|
|
private static final String PARKING_LOT = ORIGINAL_QUEUE + ".parkingLot";
|
|
|
|
private static final String X_DEATH_HEADER = "x-death";
|
|
|
|
public static void main(String[] args) throws Exception {
|
|
ConfigurableApplicationContext context = SpringApplication.run(ReRouteDlqApplication.class, args);
|
|
System.out.println("Press enter to exit");
|
|
System.in.read();
|
|
context.close();
|
|
}
|
|
|
|
@Autowired
|
|
private RabbitTemplate rabbitTemplate;
|
|
|
|
@SuppressWarnings("unchecked")
|
|
@RabbitListener(queues = DLQ)
|
|
public void rePublish(Message failedMessage) {
|
|
Map<String, Object> headers = failedMessage.getMessageProperties().getHeaders();
|
|
long retries = failedMessage.getMessageProperties().getRetryCount();
|
|
if (retries < 3) {
|
|
failedMessage.getMessageProperties().incrementRetryCount();
|
|
List<Map<String, ?>> xDeath = (List<Map<String, ?>>) headers.get(X_DEATH_HEADER);
|
|
String exchange = (String) xDeath.get(0).get("exchange");
|
|
List<String> routingKeys = (List<String>) xDeath.get(0).get("routing-keys");
|
|
this.rabbitTemplate.send(exchange, routingKeys.get(0), failedMessage);
|
|
}
|
|
else {
|
|
this.rabbitTemplate.send(PARKING_LOT, failedMessage);
|
|
}
|
|
}
|
|
|
|
@Bean
|
|
public Queue parkingLot() {
|
|
return new Queue(PARKING_LOT);
|
|
}
|
|
|
|
}
|
|
----
|
|
|
|
[[republishtodlq-true]]
|
|
=== `republishToDlq=true`
|
|
|
|
When `republishToDlq` is `true`, the republishing recoverer adds the original exchange and routing key to headers, as shown in the following example:
|
|
|
|
[source, java]
|
|
----
|
|
@SpringBootApplication
|
|
public class ReRouteDlqApplication {
|
|
|
|
private static final String ORIGINAL_QUEUE = "so8400in.so8400";
|
|
|
|
private static final String DLQ = ORIGINAL_QUEUE + ".dlq";
|
|
|
|
private static final String PARKING_LOT = ORIGINAL_QUEUE + ".parkingLot";
|
|
|
|
private static final String X_ORIGINAL_EXCHANGE_HEADER = RepublishMessageRecoverer.X_ORIGINAL_EXCHANGE;
|
|
|
|
private static final String X_ORIGINAL_ROUTING_KEY_HEADER = RepublishMessageRecoverer.X_ORIGINAL_ROUTING_KEY;
|
|
|
|
public static void main(String[] args) throws Exception {
|
|
ConfigurableApplicationContext context = SpringApplication.run(ReRouteDlqApplication.class, args);
|
|
System.out.println("Press enter to exit");
|
|
System.in.read();
|
|
context.close();
|
|
}
|
|
|
|
@Autowired
|
|
private RabbitTemplate rabbitTemplate;
|
|
|
|
@RabbitListener(queues = DLQ)
|
|
public void rePublish(Message failedMessage) {
|
|
Map<String, Object> headers = failedMessage.getMessageProperties().getHeaders();
|
|
long retries = failedMessage.getMessageProperties().getRetryCount();
|
|
if (retries < 3) {
|
|
failedMessage.getMessageProperties().incrementRetryCount();
|
|
String exchange = (String) headers.get(X_ORIGINAL_EXCHANGE_HEADER);
|
|
String originalRoutingKey = (String) headers.get(X_ORIGINAL_ROUTING_KEY_HEADER);
|
|
this.rabbitTemplate.send(exchange, originalRoutingKey, failedMessage);
|
|
}
|
|
else {
|
|
this.rabbitTemplate.send(PARKING_LOT, failedMessage);
|
|
}
|
|
}
|
|
|
|
@Bean
|
|
public Queue parkingLot() {
|
|
return new Queue(PARKING_LOT);
|
|
}
|
|
|
|
}
|
|
----
|