欢迎光临
我们一直在努力

Day37-数据层 × 中间件AI化篇:RabbitMQ消息可靠性:生产者确认+消费者ACK+持久化

 RabbitMQ 的经典坑:消费者在处理到一半时,JVM 抛了一个 OOM,整个进程挂了。MQ 以为消息没被消费,又把消息投回队列,循环往复。但因为没开 ACK 确认,消息早就从队列里被"拿走"了,等消费者重启起来,那条消息就这么消失了。

这篇文章,把 RabbitMQ 消息可能丢的三个位置、怎么兜底、怎么排查,一次讲明白。


一、先画清楚:一条消息从发送到消费,路上有三道关

很多人写代码只调convertAndSend和@RabbitListener,从来不关心中间发生了什么。我用一张图把链路展开:

对应三个可能丢消息的环节:

  • 生产者到 Exchange 阶段:网络抖动、Broker 端连接断开,发送失败但生产者不知道。
  • Exchange 到 Queue 阶段:队列没做持久化,Broker 重启后队列连同消息一起没了。
  • Queue 到 Consumer 阶段:消费者拿到消息后没处理完就崩溃,或者没 ACK 就被踢下线。
  • 下面三个代码片段,分别针对这三个环节打补丁。


    二、第一道关:生产者确认(Publisher Confirms)

    Spring AMQP 默认的convertAndSend是"发完就忘"的模式。要保证消息真到了 Broker,必须开启Publisher Confirms机制。

    配置开启

    # application.yml
    spring:
    rabbitmq:
    host: rabbitmq-prod
    port: 5672
    username: admin
    password: ${RABBIT_PWD}
    publisher-confirm-type: correlated # 关键:开启发送方确认
    publisher-returns: true # 消息无法路由时回调
    template:
    mandatory: true # 配合 publisher-returns 使用

    版本依赖:Spring Boot 3.2.x / spring-rabbit 3.x,旧版本是publisher-confirms(布尔型),新版本细分成了none/simple/correlated三种。

    代码实现:发送方确认回调

    @Component
    public class OrderPublisher {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    /**
    * 启动时注册回调
    */
    @PostConstruct
    public void init() {
    // 消息成功到达 Exchange 回调
    rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
    if (!ack) {
    // 消息没到 Exchange,记录日志并告警
    log.error("消息发送失败, correlationId={}, cause={}",
    correlationData.getId(), cause);
    // TODO: 落库到重发表,定时任务扫描重发
    } else {
    log.debug("消息发送成功, correlationId={}", correlationData.getId());
    }
    });

    // 消息从 Exchange 路由不到任何 Queue 时回调
    rabbitTemplate.setReturnsCallback(returned ->
    log.error("消息路由失败, exchange={}, routingKey={}, replyText={}",
    returned.getExchange(), returned.getRoutingKey(), returned.getReplyText())
    );
    }

    public void sendOrder(Order order) {
    // correlationData 用于在回调里识别是哪条消息
    CorrelationData cd = new CorrelationData(order.getId());
    rabbitTemplate.convertAndSend("order.direct", "order.created", order, cd);
    }
    }

    实战要点:

    • ConfirmCallback告诉你消息有没有到 Exchange;
    • ReturnsCallback告诉你消息从 Exchange 出来是不是有 Queue 收;
    • 两者配合才完整。只看 Confirm 不看 Return,可能出现"消息进了 Exchange 但掉进黑洞"的情况。

    三、第二道关:Broker 端持久化(消息不丢的物理基础)

    Publisher Confirms 只能保证消息到了 Broker,但Broker 重启之后消息还在不在,要靠持久化。RabbitMQ 的持久化分三个层级,缺一不可。

    1. 交换机持久化

    @Bean
    public DirectExchange orderExchange() {
    // durable=true 表示重启后交换机还在
    return new DirectExchange("order.direct", true, false);
    }

    2. 队列持久化

    @Bean
    public Queue orderQueue() {
    // QueueBuilder.durable() 内部就是 durable=true
    return QueueBuilder.durable("order.created.queue").build();
    }

    避坑提醒:如果之前用 Spring 默认的匿名队列(spring.gen-xxx那种),重启后队列就没了,消息自然蒸发。生产环境全部用固定名 + durable。

    3. 消息持久化

    @Service
    public class PersistentPublisher {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void send(Order order) {
    rabbitTemplate.convertAndSend("order.direct", "order.created", order, msg -> {
    // deliveryMode=2 表示消息持久化到磁盘
    msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
    return msg;
    });
    }
    }

    RabbitMQ 持久化原理(敲黑板):

    只有 durable Queue + PERSISTENT Message 才会落盘。哪怕其中一个条件不满足,重启后消息就丢。

    镜像队列:Broker 层面的高可用

    光持久化还不够——单节点磁盘坏了照样完蛋。生产环境必须部署 RabbitMQ 集群 + 镜像队列(mirrored queue)。

    # 在任一节点执行:把 order.* 队列设置为镜像模式,副本数2
    rabbitmqctl set_policy ha-order "^order\\." '{"ha-mode":"all","ha-sync-mode":"automatic"}' 1

    或者在管理后台 → Policies 里加一条:

    • Pattern: ^order\\.
    • Definition: ha-mode=all, ha-sync-mode=automatic
    • Priority: 1

    这样任何一台 Broker 宕机,消息都不会丢。代价是写入性能会下降 20~30%,对绝大多数业务可以接受。


    四、第三道关:消费者 ACK + 手动确认

    消费者这一关最容易被忽视。Spring AMQP 的@RabbitListener默认是自动 ACK——消息从队列里一拿出来就算消费成功,方法抛异常也算成功。这等于把可靠性交给运气。

    改造为手动 ACK 模式

    yaml

    spring:
    rabbitmq:
    listener:
    simple:
    acknowledge-mode: manual # 关键:手动确认
    prefetch: 10 # 预取数量,避免一次拉太多

    代码实现

    关键参数:

    • basicAck(deliveryTag, false):第二个参数multiple=false表示只确认当前消息。
    • basicNack(deliveryTag, false, requeue):最后一个参数控制是否重回队列,生产环境千万别写成true循环重试,要靠死信队列兜底。

    进阶:Spring 原生手动 ACK(更推荐)

    如果你不想直接操作Channel,Spring 提供了声明式注解:

    @Component
    public class OrderConsumer {

    @RabbitListener(queues = "order.created.queue", ackMode = "MANUAL")
    public void onMessage(Order order, Channel channel,
    @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
    try {
    processOrder(order);
    channel.basicAck(tag, false);
    } catch (Exception e) {
    channel.basicNack(tag, false, false);
    }
    }
    }


    五、消息积压:5 种处理方案

    可靠性兜底做好后,下一个高频问题就是消息积压。我列一个排查清单,按顺序处理:

    方案 1:紧急扩容消费者(5分钟见效)

    bash

    # 把消费者服务从2个实例扩到10个,先把积压消化掉
    kubectl scale deploy order-consumer –replicas=10

    临时方案,但很有效。前提是下游业务能扛住并发压力。

    方案 2:扩大预取数量

    yaml

    spring:
    rabbitmq:
    listener:
    simple:
    prefetch: 50 # 默认1太保守,提高到50让消费者吃饱

    注意:prefetch太大会导致单消费者压力大、消息分配不均,要根据机器配置调。

    方案 3:清理无用消息

    用管理后台或rabbitmqctl查看队列里都是什么消息:

    bash

    # 导出队列里的消息(测试环境用)
    rabbitmqadmin get queue=order.created.queue count=10 ackmode=ack_requeue_true

    如果发现是脏数据(比如测试时塞进去的)导致消费者一直失败,临时清空队列比查根因更实际。

    方案 4:增加死信队列单独兜底

    积压时千万别让死信队列也爆了,否则整个 RabbitMQ 集群都会卡住。给死信队列单独设置消息数量上限和TTL:

    @Bean
    public Queue orderDlq() {
    return QueueBuilder.durable("order.created.dlq")
    .withArgument("x-max-length", 10000) // 最多存1万条
    .withArgument("x-message-ttl", 7 * 24 * 3600 * 1000) // 7天过期
    .build();
    }

    方案 5:异步消费 + 限流削峰

    如果积压是周期性的(比如每天晚上8点秒杀),长期方案是把消费端做成异步:

    @RabbitListener(queues = "order.created.queue")
    public void onMessage(Order order) {
    // 把任务丢到线程池里异步处理,立刻返回 ACK
    orderProcessPool.submit(() -> processOrder(order));
    }

    这样消费者永远不会被慢业务拖死,最多积压在内存里。


    六、消息丢失排查清单(收藏版)

    最后给你一张故障排查 Checklist,出问题时按顺序查:

    排查项检查命令期望结果
    Publisher Confirms 是否开启 application.yml 查publisher-confirm-type 必须是correlated
    Queue 是否 durable 管理后台 → Queues Features 列显示D
    消息是否持久化 抓包看deliveryMode 必须为 2
    消费者 ACK 模式 acknowledge-mode配置 生产环境用manual
    是否有镜像队列 rabbitmqctl list_queues name policy 关键队列有ha-all策略
    死信队列是否满了 管理后台看 DLQ 长度 不应超过 1 万条
    消费者并发数 concurrency参数 至少 3-5 个

    七、实战建议:三条老兵经验

  • 生产环境默认开启 Publisher Confirms + 手动 ACK。 哪怕业务说"丢一两条没关系",也坚持开。线上事故里 90% 的"丢一两条"最后都是"丢一万条"。

  • 给每条消息带一个全局唯一 ID,贯穿全链路日志。 生产者用订单ID做correlationData,消费者日志里打印messageId,ELK 一查就知道哪条消息在哪个环节出问题。AI 故障复盘时也省事。

  • 可靠性方案做"双 11 演练",不要等真出问题才测试。 每年大促前,模拟三种场景:Broker 节点宕机、消费者 OOM 崩溃、数据库主从切换。观察消息是否真的不丢。只靠文档是验证不出可靠性的。


  • 结尾:一句金句 + 下篇预告

    消息丢失的锅,10 次有 9 次是配置问题,不是代码问题。

    把 Publisher Confirms、durable Queue、PERSISTENT Message、手动 ACK 这四样配齐,RabbitMQ 消息可靠性就到位了。剩下的边角问题,遇到了再补,但主线不能省。

    下一篇 Day 38,我们聊 Kafka 为什么快:顺序 I/O / 零拷贝 / 页缓存一文讲清,对比 RabbitMQ / Kafka / RocketMQ 的技术选型。消息队列家族里,Kafka 是另一条完全不同的路,值得花一整篇讲透。


    参考版本:

    • Spring Boot 3.2.x
    • spring-boot-starter-amqp 3.2.x
    • RabbitMQ 3.12+
    赞(0)
    未经允许不得转载:171主机测评 » Day37-数据层 × 中间件AI化篇:RabbitMQ消息可靠性:生产者确认+消费者ACK+持久化
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址