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 ##########################################