Spring Boot集成RabbitMQ消息可靠投递与死信队列延迟任务实战

RabbitMQ是企业级消息中间件中应用最广泛的方案之一,支持多种消息模型和可靠性保障机制。在微服务架构中,消息丢失、重复消费、延迟任务处理是后端开发面临的三大核心问题。Spring Boot通过spring-boot-starter-amqp提供了RabbitMQ的集成支持,配合消息确认机制、死信队列、幂等消费等策略,可以构建高可靠的消息处理系统。本文从RabbitMQ可靠投递的各个环节出发,结合Spring Boot实战代码,给出完整配置方案。

RabbitMQ消息可靠投递链路分析与配置

消息从生产者到消费者的完整链路包含三个环节,每个环节都可能发生消息丢失:

生产者到Exchange:网络故障或Exchange不存在导致消息丢失。通过Publisher Confirm机制解决——生产者发送消息后等待Broker的确认回执,未确认则重发。

Exchange到Queue:Routing Key不匹配或Queue不存在导致消息被丢弃。通过Mandatory标志和Return回调解决——消息无法路由时触发Return回调,生产者可以记录日志或重发。

Queue到消费者:消费者处理失败或异常退出导致消息未正确处理。通过手动ACK机制解决——消费者处理完成后手动确认,处理异常时拒绝消息使其重新入队或进入死信队列。

Spring Boot配置Publisher Confirm和Return回调:

# application.yml
spring:
  rabbitmq:
    host: 192.168.1.100
    port: 5672
    username: admin
    password: admin123
    virtual-host: /production
    publisher-confirm-type: correlated
    publisher-returns: true
    template:
      mandatory: true
    listener:
      simple:
        acknowledge-mode: manual
        prefetch: 10
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000

配置类定义Exchange、Queue和绑定关系:

@Configuration
public class RabbitMQConfig {

    // 业务队列
    public static final String ORDER_QUEUE = "order.queue";
    public static final String ORDER_EXCHANGE = "order.exchange";
    public static final String ORDER_ROUTING_KEY = "order.create";

    // 死信队列
    public static final String ORDER_DLX_EXCHANGE = "order.dlx.exchange";
    public static final String ORDER_DLX_QUEUE = "order.dlx.queue";
    public static final String ORDER_DLX_ROUTING_KEY = "order.dead";

    // 延迟队列(通过TTL+死信实现)
    public static final String DELAY_QUEUE = "delay.queue";
    public static final String DELAY_EXCHANGE = "delay.exchange";

    // 死信交换机
    @Bean
    public DirectExchange orderDlxExchange() {
        return new DirectExchange(ORDER_DLX_EXCHANGE, true, false);
    }

    @Bean
    public Queue orderDlxQueue() {
        return QueueBuilder.durable(ORDER_DLX_QUEUE).build();
    }

    @Bean
    public Binding dlxBinding() {
        return BindingBuilder.bind(orderDlxQueue())
            .to(orderDlxExchange())
            .with(ORDER_DLX_ROUTING_KEY);
    }

    // 业务队列(绑定死信交换机)
    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable(ORDER_QUEUE)
            .withArgument("x-dead-letter-exchange", ORDER_DLX_EXCHANGE)
            .withArgument("x-dead-letter-routing-key", ORDER_DLX_ROUTING_KEY)
            .build();
    }

    @Bean
    public DirectExchange orderExchange() {
        return new DirectExchange(ORDER_EXCHANGE, true, false);
    }

    @Bean
    public Binding orderBinding() {
        return BindingBuilder.bind(orderQueue())
            .to(orderExchange())
            .with(ORDER_ROUTING_KEY);
    }

    // 延迟队列:消息TTL到期后转发到业务队列
    @Bean
    public Queue delayQueue() {
        return QueueBuilder.durable(DELAY_QUEUE)
            .withArgument("x-dead-letter-exchange", ORDER_EXCHANGE)
            .withArgument("x-dead-letter-routing-key", ORDER_ROUTING_KEY)
            .withArgument("x-message-ttl", 30000)  // 30秒TTL
            .build();
    }

    @Bean
    public DirectExchange delayExchange() {
        return new DirectExchange(DELAY_EXCHANGE, true, false);
    }

    @Bean
    public Binding delayBinding() {
        return BindingBuilder.bind(delayQueue())
            .to(delayExchange())
            .with("delay.routing");
    }
}

生产者消息确认回调实现

生产者发送消息时通过ConfirmCallback和ReturnsCallback监控投递结果:

@Service
@Slf4j
public class OrderMessageProducer {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    @PostConstruct
    public void init() {
        // 消息确认回调
        rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
            if (!ack) {
                log.error("消息未到达Exchange, correlationData={}, cause={}",
                    correlationData, cause);
                // 可以在此触发重发逻辑
            }
        });

        // 消息路由失败回调
        rabbitTemplate.setReturnsCallback(returned -> {
            log.error("消息路由失败: exchange={}, routingKey={}, replyCode={}, replyText={}",
                returned.getExchange(),
                returned.getRoutingKey(),
                returned.getReplyCode(),
                returned.getReplyText());
            // 记录到数据库供后续补偿
        });
    }

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

        rabbitTemplate.convertAndSend(
            RabbitMQConfig.ORDER_EXCHANGE,
            RabbitMQConfig.ORDER_ROUTING_KEY,
            order,
            message -> {
                message.getMessageProperties().setMessageId(
                    correlationData.getId()
                );
                message.getMessageProperties().setDeliveryMode(
                    MessageDeliveryMode.PERSISTENT
                );
                return message;
            },
            correlationData
        );

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

    // 发送延迟消息
    public void sendDelayMessage(OrderDTO order, long delayMs) {
        rabbitTemplate.convertAndSend(
            RabbitMQConfig.DELAY_EXCHANGE,
            "delay.routing",
            order,
            message -> {
                message.getMessageProperties().setExpiration(
                    String.valueOf(delayMs)
                );
                return message;
            }
        );
    }
}

注意 setExpiration 设置单条消息的TTL,与Queue级别的 x-message-ttl 不同。单条消息TTL存在队头阻塞问题——如果队首消息的TTL很长,后面TTL短的消息也不会及时过期。对于精确延迟场景,推荐使用RabbitMQ延迟插件 rabbitmq_delayed_message_exchange

消费者手动确认与幂等消费实现

消费者使用手动ACK模式,确保消息处理成功后才确认:

@Component
@Slf4j
public class OrderMessageConsumer {

    @Autowired
    private OrderService orderService;

    @Autowired
    private StringRedisTemplate redisTemplate;

    @RabbitListener(queues = RabbitMQConfig.ORDER_QUEUE)
    public void handleOrderMessage(
            OrderDTO order,
            Channel channel,
            Message message
    ) throws IOException {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        String messageId = message.getMessageProperties().getMessageId();

        try {
            // 幂等检查:基于Redis SETNX
            Boolean isFirst = redisTemplate.opsForValue()
                .setIfAbsent("msg:consumed:" + messageId,
                    "1", 24, TimeUnit.HOURS);

            if (Boolean.FALSE.equals(isFirst)) {
                log.info("消息已消费,跳过: messageId={}", messageId);
                channel.basicAck(deliveryTag, false);
                return;
            }

            // 业务处理
            orderService.processOrder(order);

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

        } catch (BusinessException e) {
            // 业务异常:拒绝消息,不重新入队,进入死信队列
            log.error("业务异常: {}", e.getMessage());
            redisTemplate.delete("msg:consumed:" + messageId);
            channel.basicReject(deliveryTag, false);
        } catch (Exception e) {
            // 系统异常:拒绝消息,重新入队重试
            log.error("系统异常", e);
            redisTemplate.delete("msg:consumed:" + messageId);

            // 判断重试次数
            Integer retryCount = getRetryCount(message);
            if (retryCount >= 3) {
                // 超过重试次数,进入死信队列
                channel.basicReject(deliveryTag, false);
            } else {
                // 重新入队
                channel.basicNack(deliveryTag, false, true);
            }
        }
    }

    private Integer getRetryCount(Message message) {
        Map<String, Object> headers =
            message.getMessageProperties().getHeaders();
        Integer count = (Integer) headers.get("x-retry-count");
        return count != null ? count : 0;
    }
}

死信队列消费者与延迟任务处理

死信队列中的消息通常是处理失败或过期的消息,需要单独消费处理:

@Component
@Slf4j
public class DeadLetterConsumer {

    @Autowired
    private OrderFailService orderFailService;

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

        try {
            // 记录失败订单,触发人工处理或补偿逻辑
            orderFailService.recordFailedOrder(order, "超过最大重试次数");

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

        } catch (Exception e) {
            log.error("死信消息处理失败", e);
            // 死信消息处理失败直接丢弃或记录到数据库
            channel.basicAck(deliveryTag, false);
        }
    }
}

延迟任务场景(如订单30分钟未支付自动取消)通过TTL+死信队列实现:生产者发送消息到延迟队列,消息TTL到期后自动转发到业务队列触发消费处理。延迟队列方案适用于固定延迟场景,对于动态延迟建议使用延迟插件 rabbitmq_delayed_message_exchange,在Exchange层面实现延迟投递,不存在队头阻塞问题。

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

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

相关推荐