Spring Boot集成RabbitMQ死信队列:消息可靠投递实战方案

消息中间件在高并发架构中承担异步解耦和削峰填谷职责。RabbitMQ作为主流消息中间件,在业务中台建设中广泛应用。消息消费失败时的处理策略直接影响系统可靠性。死信队列(Dead Letter Queue)将消费失败的消息路由到备用队列,配合重试机制实现消息可靠投递。本文记录Spring Boot集成RabbitMQ死信队列的完整配置和代码实现,涵盖服务治理和消息中间件运维要点。

RabbitMQ死信队列原理与架构设计

消息在以下三种情况下会成为死信:

  • 消息被拒绝(basic.reject/basic.nack),且requeue参数为false
  • 消息TTL过期:消息在队列中存活时间超过设定的TTL值
  • 队列达到最大长度:队列消息数量超过x-max-length限制

架构设计:业务交换机将消息路由到业务队列,业务队列绑定死信交换机。消息消费失败后进入死信交换机,路由到死信队列。死信队列消费者进行告警通知和人工补偿处理。

业务交换机(exchange.business) --> 业务队列(queue.business)
                                        |
                                   消费失败/TTL过期
                                        |
                                   死信交换机(exchange.dlx)
                                        |
                                   死信队列(queue.dlx)

Spring Boot项目初始化与依赖配置

Maven依赖配置:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
</dependency>

application.yml配置RabbitMQ连接参数:

spring:
  rabbitmq:
    host: 192.168.1.50
    port: 5672
    username: admin
    password: ${RABBITMQ_PASSWORD}
    virtual-host: /production
    publisher-confirm-type: correlated
    publisher-returns: true
    listener:
      simple:
        acknowledge-mode: manual
        prefetch: 10
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000
          multiplier: 2.0

队列与交换机声明配置

使用Java Config声明业务队列、死信交换机和死信队列的绑定关系:

import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class RabbitMQConfig {

    // === 业务队列配置 ===
    public static final String BUSINESS_EXCHANGE = "exchange.business";
    public static final String BUSINESS_QUEUE = "queue.business";
    public static final String BUSINESS_ROUTING_KEY = "routing.business";

    // === 死信队列配置 ===
    public static final String DLX_EXCHANGE = "exchange.dlx";
    public static final String DLX_QUEUE = "queue.dlx";
    public static final String DLX_ROUTING_KEY = "routing.dlx";

    // 业务交换机
    @Bean
    public DirectExchange businessExchange() {
        return ExchangeBuilder.directExchange(BUSINESS_EXCHANGE)
                .durable(true)
                .build();
    }

    // 死信交换机
    @Bean
    public DirectExchange dlxExchange() {
        return ExchangeBuilder.directExchange(DLX_EXCHANGE)
                .durable(true)
                .build();
    }

    // 业务队列(绑定死信交换机)
    @Bean
    public Queue businessQueue() {
        return QueueBuilder.durable(BUSINESS_QUEUE)
                .withArgument("x-dead-letter-exchange", DLX_EXCHANGE)
                .withArgument("x-dead-letter-routing-key", DLX_ROUTING_KEY)
                .withArgument("x-message-ttl", 60000) // 消息TTL 60秒
                .withArgument("x-max-length", 10000)  // 队列最大长度
                .build();
    }

    // 死信队列
    @Bean
    public Queue dlxQueue() {
        return QueueBuilder.durable(DLX_QUEUE).build();
    }

    // 绑定关系
    @Bean
    public Binding businessBinding() {
        return BindingBuilder.bind(businessQueue())
                .to(businessExchange())
                .with(BUSINESS_ROUTING_KEY);
    }

    @Bean
    public Binding dlxBinding() {
        return BindingBuilder.bind(dlxQueue())
                .to(dlxExchange())
                .with(DLX_ROUTING_KEY);
    }
}

消息生产者实现与确认机制

生产者发送消息时启用Confirm回调,确保消息到达交换机和队列:

import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import lombok.extern.slf4j.Slf4j;
import java.util.UUID;

@Slf4j
@Component
public class OrderMessageProducer {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void sendOrderMessage(OrderDTO order) {
        CorrelationData correlationId = new CorrelationData(UUID.randomUUID().toString());

        rabbitTemplate.convertAndSend(
            RabbitMQConfig.BUSINESS_EXCHANGE,
            RabbitMQConfig.BUSINESS_ROUTING_KEY,
            order,
            message -> {
                message.getMessageProperties().setDeliveryTag(1);
                message.getMessageProperties().setExpiration("30000"); // 单条消息TTL
                return message;
            },
            correlationId
        );

        log.info("订单消息已发送: orderId={}, correlationId={}",
                order.getOrderId(), correlationId.getId());
    }
}

配置Confirm回调处理消息投递结果:

@Component
@Slf4j
public class RabbitConfirmCallback implements RabbitTemplate.ConfirmCallback {

    @Override
    public void confirm(CorrelationData correlationData, boolean ack, String cause) {
        if (ack) {
            log.info("消息投递成功: correlationId={}", correlationData.getId());
        } else {
            log.error("消息投递失败: correlationId={}, cause={}",
                    correlationData.getId(), cause);
            // 投递失败处理:记录到数据库,定时任务重发
        }
    }
}

消息消费者实现与异常处理

消费者手动确认消息,处理失败时reject到死信队列:

import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.core.ChannelAwareMessageListener;
import com.rabbitmq.client.Channel;
import org.springframework.stereotype.Component;
import lombok.extern.slf4j.Slf4j;

@Slf4j
@Component
public class OrderMessageConsumer {

    @RabbitListener(queues = RabbitMQConfig.BUSINESS_QUEUE)
    public void handleMessage(OrderDTO order, Message message, Channel channel) throws Exception {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();

        try {
            // 业务处理
            processOrder(order);

            // 处理成功,手动ACK
            channel.basicAck(deliveryTag, false);
            log.info("订单处理成功: orderId={}", order.getOrderId());

        } catch (BusinessException e) {
            // 业务异常,判断重试次数
            Integer retryCount = getRetryCount(message);
            if (retryCount >= 3) {
                // 超过重试次数,reject到死信队列
                log.error("订单处理失败,超过重试次数,进入死信队列: orderId={}, retry={}",
                        order.getOrderId(), retryCount);
                channel.basicReject(deliveryTag, false);
            } else {
                // 未超过重试次数,重新入队重试
                log.warn("订单处理失败,重试中: orderId={}, retry={}",
                        order.getOrderId(), retryCount);
                channel.basicNack(deliveryTag, false, true);
            }
        } catch (Exception e) {
            // 未知异常,直接进入死信队列
            log.error("订单处理未知异常,进入死信队列: orderId={}", order.getOrderId(), e);
            channel.basicReject(deliveryTag, false);
        }
    }

    private void processOrder(OrderDTO order) {
        // 扣减库存、创建物流单、发送通知等业务逻辑
        inventoryService.deduct(order.getProductId(), order.getQuantity());
        logisticsService.createShipment(order);
    }

    private Integer getRetryCount(Message message) {
        Object count = message.getMessageProperties()
                .getHeader("x-retry-count");
        return count != null ? (Integer) count : 0;
    }
}

死信队列消费者与告警补偿

死信队列消费者负责告警通知和人工补偿处理:

@Slf4j
@Component
public class DeadLetterConsumer {

    @Autowired
    private AlertService alertService;

    @Autowired
    private OrderFailedRepository orderFailedRepo;

    @RabbitListener(queues = RabbitMQConfig.DLX_QUEUE)
    public void handleDeadLetter(Message message, Channel channel) throws Exception {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();

        try {
            OrderDTO order = (OrderDTO) message.getMessageProperties()
                    .getConverter()
                    .fromMessage(message);

            // 持久化失败订单到数据库,供人工排查
            OrderFailedRecord record = new OrderFailedRecord();
            record.setOrderId(order.getOrderId());
            record.setPayload(JSON.toJSONString(order));
            record.setReason("消费失败,超过最大重试次数");
            record.setCreateTime(LocalDateTime.now());
            orderFailedRepo.save(record);

            // 发送告警通知
            alertService.sendAlert(String.format(
                "订单消息进入死信队列,需人工处理。订单号: %s",
                order.getOrderId()
            ));

            channel.basicAck(deliveryTag, false);
            log.warn("死信消息已处理: orderId={}", order.getOrderId());

        } catch (Exception e) {
            log.error("死信消息处理异常", e);
            // 死信处理失败不reject,避免死信队列消息丢失
            channel.basicNack(deliveryTag, false, false);
        }
    }
}

API接口规范与监控治理

对外提供消息发送的REST API接口:

@RestController
@RequestMapping("/api/v1/orders")
public class OrderController {

    @Autowired
    private OrderMessageProducer producer;

    @PostMapping
    public Result<String> createOrder(@RequestBody @Valid OrderDTO order) {
        // 保存订单到数据库
        orderService.save(order);
        // 异步发送消息
        producer.sendOrderMessage(order);
        return Result.success(order.getOrderId());
    }
}

RabbitMQ管理端监控关键指标:

# 通过RabbitMQ HTTP API获取队列状态
curl -u admin:password http://192.168.1.50:15672/api/queues/production/queue.dlx

# 关注指标:
# messages:死信队列积压消息数(应为0,持续增长说明消费端异常)
# messages_ready:待消费消息数
# consumers:消费者数量(应>=1)
# messages_unacknowledged:未确认消息数

配置Prometheus采集RabbitMQ指标,设置告警规则:死信队列消息数大于0时触发告警;业务队列消息积压超过5000条触发告警;消费者连接数降为0触发紧急告警。配合Grafana面板可视化消息投递成功率、消费延迟、重试次数等指标,构建完整的消息中间件监控体系。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-ji-cheng-rabbitmq-si-xin-dui-lie-xiao-xi-ke-kao/

(0)
小编小编
上一篇 4小时前
下一篇 4小时前

相关推荐