From cd398b40ee3fb5b9013b2541730c4e4252ecf28c Mon Sep 17 00:00:00 2001 From: Rangsh <2570024918@qq.com> Date: Mon, 13 Jul 2026 14:13:22 +0800 Subject: [PATCH] =?UTF-8?q?feat(rabbitmq):=20=E6=96=B0=E5=A2=9E=E6=AD=BB?= =?UTF-8?q?=E4=BF=A1=E9=98=9F=E5=88=97=EF=BC=88DLQ=EF=BC=89=E6=94=AF?= =?UTF-8?q?=E6=8C=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 当前 RabbitMQ 已实现 Confirm + Returns + Retry,但缺少 DLQ 导致重试耗尽后消息直接丢失,无法追溯和人工补偿。 新增: - RabbitMqDeadLetterConfig:声明 DLX 交换机及 send/recall 死信队列 - RabbitMqDeadLetterReceiver:消费死信消息,记录 headers 与 body 日志 - 主队列绑定 x-dead-letter-exchange / x-dead-letter-routing-key - 配置 default-requeue-rejected=false,确保重试耗尽后进入 DLQ 而非 requeue 消息流向: austin.queues.send/recall → 消费失败 → 重试 3 次 → 耗尽 → austin.dlx → austin.queues.send.dead/recall.dead → 日志告警 --- .../rabbit/RabbitMqDeadLetterReceiver.java | 38 +++++++++++ .../receiver/rabbit/RabbitMqReceiver.java | 19 +++++- .../config/RabbitMqDeadLetterConfig.java | 66 +++++++++++++++++++ .../src/main/resources/application.properties | 8 +++ 4 files changed, 129 insertions(+), 2 deletions(-) create mode 100644 austin-handler/src/main/java/com/java3y/austin/handler/receiver/rabbit/RabbitMqDeadLetterReceiver.java create mode 100644 austin-support/src/main/java/com/java3y/austin/support/config/RabbitMqDeadLetterConfig.java diff --git a/austin-handler/src/main/java/com/java3y/austin/handler/receiver/rabbit/RabbitMqDeadLetterReceiver.java b/austin-handler/src/main/java/com/java3y/austin/handler/receiver/rabbit/RabbitMqDeadLetterReceiver.java new file mode 100644 index 0000000..951d856 --- /dev/null +++ b/austin-handler/src/main/java/com/java3y/austin/handler/receiver/rabbit/RabbitMqDeadLetterReceiver.java @@ -0,0 +1,38 @@ +package com.java3y.austin.handler.receiver.rabbit; + +import com.java3y.austin.support.constans.MessageQueuePipeline; +import lombok.extern.slf4j.Slf4j; +import org.springframework.amqp.core.Message; +import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +import java.nio.charset.StandardCharsets; + +/** + * 死信队列消费者 + * 消费重试耗尽后进入 DLQ 的消息,记录日志便于排查与人工补偿 + * + * @author Rangsh + */ +@Slf4j +@Component +@ConditionalOnProperty(name = "austin.mq.pipeline", havingValue = MessageQueuePipeline.RABBIT_MQ) +public class RabbitMqDeadLetterReceiver { + + @RabbitListener(queues = "${austin.rabbitmq.queues.send.dead}") + public void onSendDeadLetter(Message message) { + log.error("[DLQ][send] headers={}, body={}", + message.getMessageProperties().getHeaders(), + new String(message.getBody(), StandardCharsets.UTF_8)); + // TODO: 后续可扩展落库 / 告警 / 人工补偿入口 + } + + @RabbitListener(queues = "${austin.rabbitmq.queues.recall.dead}") + public void onRecallDeadLetter(Message message) { + log.error("[DLQ][recall] headers={}, body={}", + message.getMessageProperties().getHeaders(), + new String(message.getBody(), StandardCharsets.UTF_8)); + // TODO: 后续可扩展落库 / 告警 / 人工补偿入口 + } +} diff --git a/austin-handler/src/main/java/com/java3y/austin/handler/receiver/rabbit/RabbitMqReceiver.java b/austin-handler/src/main/java/com/java3y/austin/handler/receiver/rabbit/RabbitMqReceiver.java index cdf2c70..8418e81 100644 --- a/austin-handler/src/main/java/com/java3y/austin/handler/receiver/rabbit/RabbitMqReceiver.java +++ b/austin-handler/src/main/java/com/java3y/austin/handler/receiver/rabbit/RabbitMqReceiver.java @@ -9,6 +9,7 @@ import com.java3y.austin.support.constans.MessageQueuePipeline; import org.apache.commons.lang3.StringUtils; import org.springframework.amqp.core.ExchangeTypes; import org.springframework.amqp.core.Message; +import org.springframework.amqp.rabbit.annotation.Argument; import org.springframework.amqp.rabbit.annotation.Exchange; import org.springframework.amqp.rabbit.annotation.Queue; import org.springframework.amqp.rabbit.annotation.QueueBinding; @@ -33,7 +34,14 @@ public class RabbitMqReceiver implements MessageReceiver { private ConsumeService consumeService; @RabbitListener(bindings = @QueueBinding( - value = @Queue(value = "${spring.rabbitmq.queues.send}", durable = "true"), + value = @Queue( + value = "${spring.rabbitmq.queues.send}", + durable = "true", + arguments = { + @Argument(name = "x-dead-letter-exchange", value = "${austin.rabbitmq.dlx.exchange}"), + @Argument(name = "x-dead-letter-routing-key", value = "${austin.rabbitmq.routing.send.dead}") + } + ), exchange = @Exchange(value = "${austin.rabbitmq.exchange.name}", type = ExchangeTypes.TOPIC), key = "${austin.rabbitmq.routing.send}" )) @@ -49,7 +57,14 @@ public class RabbitMqReceiver implements MessageReceiver { } @RabbitListener(bindings = @QueueBinding( - value = @Queue(value = "${spring.rabbitmq.queues.recall}", durable = "true"), + value = @Queue( + value = "${spring.rabbitmq.queues.recall}", + durable = "true", + arguments = { + @Argument(name = "x-dead-letter-exchange", value = "${austin.rabbitmq.dlx.exchange}"), + @Argument(name = "x-dead-letter-routing-key", value = "${austin.rabbitmq.routing.recall.dead}") + } + ), exchange = @Exchange(value = "${austin.rabbitmq.exchange.name}", type = ExchangeTypes.TOPIC), key = "${austin.rabbitmq.routing.recall}" )) diff --git a/austin-support/src/main/java/com/java3y/austin/support/config/RabbitMqDeadLetterConfig.java b/austin-support/src/main/java/com/java3y/austin/support/config/RabbitMqDeadLetterConfig.java new file mode 100644 index 0000000..2c96d97 --- /dev/null +++ b/austin-support/src/main/java/com/java3y/austin/support/config/RabbitMqDeadLetterConfig.java @@ -0,0 +1,66 @@ +package com.java3y.austin.support.config; + +import com.java3y.austin.support.constans.MessageQueuePipeline; +import org.springframework.amqp.core.Binding; +import org.springframework.amqp.core.BindingBuilder; +import org.springframework.amqp.core.DirectExchange; +import org.springframework.amqp.core.Queue; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +/** + * RabbitMQ 死信队列配置 + * 重试耗尽后消息进入 DLQ,避免消息丢失,便于人工排查与补偿 + * + * @author Rangsh + */ +@Configuration +@ConditionalOnProperty(name = "austin.mq.pipeline", havingValue = MessageQueuePipeline.RABBIT_MQ) +public class RabbitMqDeadLetterConfig { + + @Value("${austin.rabbitmq.dlx.exchange}") + private String dlxExchangeName; + + @Value("${austin.rabbitmq.queues.send.dead}") + private String sendDeadQueueName; + + @Value("${austin.rabbitmq.queues.recall.dead}") + private String recallDeadQueueName; + + @Value("${austin.rabbitmq.routing.send.dead}") + private String sendDeadRoutingKey; + + @Value("${austin.rabbitmq.routing.recall.dead}") + private String recallDeadRoutingKey; + + @Bean + public DirectExchange deadLetterExchange() { + return new DirectExchange(dlxExchangeName, true, false); + } + + @Bean + public Queue sendDeadLetterQueue() { + return new Queue(sendDeadQueueName, true); + } + + @Bean + public Queue recallDeadLetterQueue() { + return new Queue(recallDeadQueueName, true); + } + + @Bean + public Binding sendDeadLetterBinding() { + return BindingBuilder.bind(sendDeadLetterQueue()) + .to(deadLetterExchange()) + .with(sendDeadRoutingKey); + } + + @Bean + public Binding recallDeadLetterBinding() { + return BindingBuilder.bind(recallDeadLetterQueue()) + .to(deadLetterExchange()) + .with(recallDeadRoutingKey); + } +} diff --git a/austin-web/src/main/resources/application.properties b/austin-web/src/main/resources/application.properties index fa7cc20..a8aab2b 100644 --- a/austin-web/src/main/resources/application.properties +++ b/austin-web/src/main/resources/application.properties @@ -71,11 +71,19 @@ spring.rabbitmq.listener.simple.retry.enabled=true spring.rabbitmq.listener.simple.retry.initial-interval=1000 spring.rabbitmq.listener.simple.retry.multiplier=2 spring.rabbitmq.listener.simple.retry.max-attempts=3 +# retry exhausted → reject (no requeue) → DLQ +spring.rabbitmq.listener.simple.default-requeue-rejected=false austin.rabbitmq.exchange.name=austin.point spring.rabbitmq.queues.send=austin.queues.send spring.rabbitmq.queues.recall=austin.queues.recall austin.rabbitmq.routing.send=austin.send austin.rabbitmq.routing.recall=austin.recall +# dead letter +austin.rabbitmq.dlx.exchange=austin.dlx +austin.rabbitmq.queues.send.dead=austin.queues.send.dead +austin.rabbitmq.queues.recall.dead=austin.queues.recall.dead +austin.rabbitmq.routing.send.dead=austin.send.dead +austin.rabbitmq.routing.recall.dead=austin.recall.dead ########################################## RabbitMq end ##########################################