
1. 先认清消息从生产到消费的完整链路再谈怎么防丢聊到 RabbitMQ 消息不丢失网上能搜到的答案绝大多数是老三样交换机持久化、队列持久化、消息持久化。背下来去面试确实够用但真到了线上环境你会发现消息丢得莫名其妙持久化都开了消费端也手动确认了消息照样消失。原因在于——你根本还没搞清楚消息是在哪一步丢的。一条消息从诞生到被业务消费实际要穿过三个环节生产者客户端把消息发到 BrokerBroker 内部完成路由并将消息写入队列消费者从队列拉取消息并处理。这三个环节里的每一个都存在丢消息的可能性。生产者发到交换机之前网络断了消息没发出去Broker 路由不到对应队列消息被静默丢弃队列收到了但 Broker 突然宕机内存里的消息没来得及落盘消费者收到消息后还没处理完进程就挂了。任何一个节点出问题整条链路的数据可靠性就归零。所以真正靠谱的思路是分而治之——把一条消息的完整旅程切成三段每一段都有各自的保底手段缺了哪一段都不行。这也是我写这篇文章的初衷不谈虚的把每个环节的防丢配置、代码写法和背后的设计原因讲清楚。先看一下如何快速定位当前消息是在哪个环节丢的。如果你用的是 RabbitMQ 自带的 Web 管理界面可以重点盯三个指标队列的 Ready 数量、Unacked 数量以及消息总数。如果队列的 Total 一直不涨说明消息压根没进来问题在生产者端或者交换机路由如果 Ready 不降但 Unacked 一直很高说明消费者每次消费都很慢甚至卡住了如果消息从 Ready 变成 Unacked 之后消费者一重启就被清空那十有八九是自动 ACK 导致的经典丢失场景。2. 生产端确认机制publisher confirm 才是第一道关口很多团队在讨论消息不丢失时第一反应是把注意力全放在 Broker 和消费者上忽略了最前面的生产者。但实际生产环境中生产者到 Broker 之间那一跳是最容易因为网络抖动、超时、连接被重置而丢消息的。2.1 开启 confirm 模式收到 Broker 的回执才算发送成功RabbitMQ 从 3.0 开始支持 Publisher Confirm 机制这是生产端防丢的根本手段。核心逻辑是生产者将消息发送到 Broker 之后Broker 会返回一个 ack 给生产者告诉它“我已经成功收到这条消息了”。如果 Broker 内部出现异常无法持久化会返回 nack。在 Spring Boot 环境下只需要在配置文件里显式开启 publisher-confirm-type。spring: rabbitmq: publisher-confirm-type: correlated publisher-returns: true第一个配置项 correlated 表示开启 confirm 回调并且回调能携带到具体的 correlatedData 对象。第二个配置项 publisher-returns 表示当消息无法路由到任何队列时Broker 要通知生产者执行失败回调。在代码里需要实现 RabbitTemplate.ConfirmCallback 和 RabbitTemplate.ReturnsCallbackrabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (!ack) { // 记录失败消息考虑重发或落库 log.error(消息发送确认失败: {}, cause); } }); rabbitTemplate.setReturnsCallback(returned - { // mandatory 消息路由不到队列时触发 log.error(消息路由失败交换机: {}, 路由键: {}, returned.getExchange(), returned.getRoutingKey()); // 这里可以把回退的消息保存到数据库后续做补偿 });重点说一个容易被忽略的细节publisher confirm 的 ack 只在消息成功写入交换机时触发并不代表消息已经写入队列。如果你用的是 default exchange空字符串交换机它会直接路由到同名队列这种场景下交换机写入基本就等同于队列写入。但如果是 fanout 或 topic 交换机交换机收到消息并不代表任何队列都存下来了。所以一定要同时开启 mandatory 机制配合 returns 回调才能在路由失败时拿到消息。2.2 发送端不要只发一次要有重试和补偿机制即使开启了 confirm网络瞬时抖动也会导致 ack 一直没回来或者连接直接断了。此时如果你只在 ack 回调里打个日志就结束消息照样会丢。正确姿势是发送消息时生成一个全局唯一的业务消息 ID把消息内容先持久化到本地数据库发送成功后把状态置为已确认。这样即使回调丢失、发送失败也可以启动一个定时任务扫描数据库中长时间处于“发送中”状态的消息做定时重发。这个思路在真实业务里被验证是非常好用的先把消息按业务单据落库状态为 PENDING发送并收到 ack 后更新为 SENT。定时任务每 30 秒扫描一次 5 分钟还没有确认的消息重新发送。这样生产者成功的标准不再是“发出去”而是“收到 Broker 的明确确认”两件事的区别是本质性的。还有一个点值得注意mandatory 参数设置为 true 之后如果路由不到队列消息就会通过 returns 回调返回给生产者。但这里有一个坑——returns 回调里的消息不是你发送时的原始类型而是 MessagePostProcessor 处理后的 Message 对象。所以你在回调里重新入库时要确保消息体里带有业务 ID或者把 correlationData 里塞上业务唯一键否则补偿很难做。3. Broker 端持久化durable 不只是开了就完事Broker 端的持久化最常被写成“三道配置交换机设 durable、队列设 durable、发送消息时 deliveryMode 设为 2”。这句话方向没错但关键在于怎么设置、在什么时机设置。3.1 三层持久化的正确配置方式交换机持久化相对简单声明交换机时指定 durable 即可Bean public DirectExchange orderExchange() { return new DirectExchange(order.exchange, true, false); }第二个参数 durable 为 true第三个参数 autoDelete 设为 false。如果之前已经用非持久化方式声明过同名交换机再改成 durable 是不生效的——只能删除重建。这也是很多人改了配置没效果的原因之一。队列持久化同理Bean public Queue orderQueue() { return QueueBuilder.durable(order.queue).build(); }消息持久化则是在发送时设置消息属性MessageProperties properties new MessageProperties(); properties.setDeliveryMode(MessageDeliveryMode.PERSISTENT); Message message new Message(订单数据.getBytes(StandardCharsets.UTF_8), properties); rabbitTemplate.convertAndSend(order.exchange, order.create, message);如果用 Spring 的 RabbitTemplate.convertAndSend 发送对象默认不会带上持久化属性建议要么显式传入 Message要么配置全局的 message converter。我看到很多线上项目用的是自定义 Jackson2JsonMessageConverter 配合 SimpleMessageConverter 的默认策略结果消息体是 JSON非持久化属性也是默认的——这种组合最容易让人误以为一切都正常实际发出的消息根本没做持久化标记。3.2 持久化到磁盘并不代表立刻刷盘这是最容易产生误解的一层。在 RabbitMQ 的默认行为中消息进入队列后并不会立即 fsync 到磁盘而是先写入操作系统页缓存由 RabbitMQ 的持久化进程在合适的时机批量刷盘。这意味着如果消息刚写入队列、还没刷盘的那一刻机器突然断电这部分数据依然会丢。这种丢失属于“没法从客户端代码层面完全规避”的范畴彻底解决需要 RabbitMQ 端调整持久化策略例如修改 rabbitmq.conf 中的 disk_fsync_interval 参数设置一个更小的刷盘间隔。代价是吞吐量下降——每一条消息都可能触发磁盘同步。个人经验是在高可靠性优先的业务场景如订单、支付对账可以接受这个性能折损但如果是日志、统计类数据没必要把所有消息都做成持久化使用 Lazy Queue 配合少量 topic 反而是更经济的方案。3.3 队列副本才是应对宕机的核心在单机模式下无论你刷盘多积极机器彻底宕机后系统重启也可能面临磁盘损坏或数据不一致风险。更合理的做法是让同一份消息在集群内保留多份副本任意一个节点挂了其他节点还能继续服务。RabbitMQ 中做到这一点需要配置镜像队列Mirrored Queue或者直接选用比较新的 Quorum Queue。经典镜像队列的做法是在声明队列时通过 x-ha-policy 参数开启MapString, Object args new HashMap(); args.put(x-ha-policy, all); Queue durableQueue new Queue(order.queue, true, false, false, args);all 表示在集群所有节点上创建镜像消息会被同步复制到每个节点的同名队列中。看起来简单但镜像队列存在明显的性能短板同步复制是同步阻塞式的每秒吞吐量上不去。更关键的是镜像队列在节点故障时存在确认丢失的窗口如果 master 节点收到消息后还没来得及同步到 slave 就宕机了这条消息依然会丢。因此我个人的建议是新项目一律优先考虑 Quorum Queue仲裁队列。它是 Raft 协议的实现要求多数副本写入成功才返回确认从存储设计上保证了消息不因单节点宕机而丢失。声明方式同样简单Bean public Queue orderQuorumQueue() { return QueueBuilder.durable(order.queue) .quorum() .build(); }4. 消费端手动 ACK重灾区中的重灾区服务端防住了消费端又是一个大雷区。默认情况下RabbitMQ 的消费端 ACK 模式是自动确认Consumer 收到消息后Broker 立刻把这条消息标记为已消费从队列中移除。如果业务代码在收到消息后崩溃、抛异常或者处理逻辑根本没执行完这条消息就永远消失了。4.1 改成手动 ACK 的两种姿势在 Spring Boot 下最典型的手动 ACK 配置是spring: rabbitmq: listener: simple: acknowledge-mode: manual然后消费者方法中显式调用 basicAck 或 basicNackRabbitListener(queues order.queue) public void handleOrder(Channel channel, Message message) { long deliveryTag message.getMessageProperties().getDeliveryTag(); try { // 业务逻辑处理 channel.basicAck(deliveryTag, false); } catch (Exception e) { // 处理失败 channel.basicNack(deliveryTag, false, true); } }注意 basicNack 的第三个参数是 requeue如果填 true消息会重新放回队列Broker 会把这条消息再次投递给消费者。这里有一个很反直觉的坑如果业务代码有 bug每次消费这条消息都会抛异常那么这个 requeuetrue 会导致消息在同一个死循环里无限重试把 CPU 打满队列里其他正常消息全部排队等待。正确做法是根据异常类型决定要不要 requeue。可重试的临时错误比如下游接口 5xx可以 requeue 几次不可重试的持久性错误比如参数校验失败、业务状态异常要 requeuefalse让消息进入死信队列或者至少记录下来让监控系统去告警。经验之谈线上很多事故都是因为无脑 requeue 把消费者实例拖垮了。4.2 basicReject 和 basicNack 的边界要分清很多初学者分不清 basicReject 和 basicNack 的区别。basicReject 只能用于单条消息拒绝不支持批量basicNack 可以同时拒绝多条。如果你使用 Spring AMQP默认情况下即使 method 参数只传了单条也建议统一用 basicNack因为它的语义更丰富后面扩展批量拒绝时不需要改代码。还有一个容易遗漏的细节requeuefalse 的 nack 只有在队列配置了死信交换机时消息才会被转到死信队列。如果没有配置这条消息会被直接丢弃不会再有补偿机会。所以实操中凡是用到手动 ACK几乎都会同时配上死信队列这是下一节要展开的方案。4.3 Unacked 消息堆积带来的风险手动 ACK 模式下消费者拉取到的消息如果一直不确认Broker 会把它标记为 Unacked 状态。Unacked 消息不算丢失但问题在于如果消费者进程被 kill这些 Unacked 消息会自动变回 Ready 重新投递如果消费者持续消费失败又不确认Unacked 数会越堆越高直到触发 channel 的流控或连接阻塞这时候连生产者都发不进新消息了引发连锁故障。所以手动 ACK 的正确姿势是“尽快速确认异常单独处理”。消息的消费不能理解为“拿到了就算数”而应该理解为“处理成功才算数”。每一处 ACK 的位置都要想清楚这条消息真的处理完了吗如果没处理完就 ACK消息丢失的责任就从 Broker 转移到了业务代码上。5. 集群高可用与仲裁队列能解决什么解决不了什么很多人会寄希望于“RabbitMQ 搞成集群消息就绝对不丢了”。这句话严谨地讲是错的——集群解决的是可用性问题也就是节点宕机后服务还能不能继续对外提供服务。它并不等于每条消息都百分百不丢。5.1 单节点宕机的恢复路径单机环境下即使你开了持久化如果机器硬件故障导致磁盘损坏无法恢复数据一样归零。集群环境下如果用的是镜像队列或仲裁队列消息在多个节点有副本一个节点挂掉后其他节点依然持有数据。这里的核心前提是消息在写入时已经被复制到足够多的副本。因此正确的高可用使用方式是集群节点数至少三个仲裁队列副本数至少三个默认行为这样任意一个节点挂掉都不会丢失写入成功的消息。如果你只有两个节点且副本数只有 1那和单机其实没什么本质区别。5.2 网络分区对消息可靠性的影响集群环境下另一个坑是网络分区split-brain。当 RabbitMQ 集群中的节点之间网络中断时各分区会各自为政可能出现两边都尝试处理同一批消息的情况。这时候消息不丢失和消息不重复不可兼得要保证不丢就可能产生重复投递要保证不重复就可能丢弃消息。实际操作中绝大多数业务场景更倾向于保证“至少一次传递”at least once而不是“恰好一次传递”exactly once。因为消费重复可以用幂等字段去重而消息丢失则很难弥补。所以配置集群时建议使用自动分区处理策略autoheal 或 pause-minority尽量避免手动干预。6. 端到端兜底方案死信、TTL 与对账补偿一个都不能少即使你把上面的每一个环节都配置到了位RabbitMQ 仍然存在一种隐蔽的丢消息路径——内存中的消息还没落盘时节点突然宕机或者 RabbitMQ 集群整体无法恢复。这时候必须依赖端到端的兜底方案而不是只靠 RabbitMQ 机制本身。6.1 死信交换机是消息的“安全网”给每个重要业务队列配置一个死信交换机DLX是防丢的第一层兜底。步骤是先声明一个普通队列给它附加 x-dead-letter-exchange 和 x-dead-letter-routing-key 参数。MapString, Object args new HashMap(); args.put(x-dead-letter-exchange, order.dlx.exchange); args.put(x-dead-letter-routing-key, order.dlx); Queue queue QueueBuilder.durable(order.queue).withArguments(args).build();当消息被消费端 nack 且 requeuefalse、或者消息 TTL 到期、或者队列消息数量超过最大长度时这些消息会被自动发送到死信交换机。死信队列的作用不是“让消息消失”而是“让处理不了的消息在另一个队列里等待补偿”。监控系统直接监听死信队列的堆积量就能第一时间发现业务异常。6.2 TTL 死信实现延迟重试真实业务中很多消息的消费失败只是下游临时性抖动比如数据库连接池满了、外部 API 超时。这种失败直接丢弃太过可惜无脑 requeue 又会造成死循环。更好的方案是把消息设置为延迟重试消费失败后 nack 且 requeuefalse配置一个带有 TTL 的延迟队列消息进入这个队列等待 10 秒后再转到业务队列重新投递。实现思路是使用两个队列一个业务队列一个延迟队列。延迟队列设置 x-dead-letter-exchange 指向业务交换机同时设置 x-message-ttl 为 10000 毫秒。当消息在延迟队列中存活 10 秒后RabbitMQ 自动将它转移到业务队列。此时消费者再次消费时如果还是失败可以继续走这条链路。建议限制重试次数比如消费失败次数超过 3 次后直接把消息丢进最终死信队列由人工介入处理。重试次数可以在消息头里用自定义属性记录每投递一次加一。6.3 落库对账消息可靠性最后的防线以上所有手段都在 RabbitMQ 机制和配置层面打转但真正的终极兜底方案从来不在 MQ 内部而在业务系统自己。我在多个高并发订单系统里实践过的方案是把消息本身当作一条业务数据落库同时把消费结果也落库然后通过另一个定时任务做对账。具体来说可以设计一张消息表字段说明msg_id全局唯一消息 ID由生产者生成biz_type业务类型如订单、支付、库存payload消息体 JSONstatus0发送中1已确认2路由失败3消费成功retry_count重试次数create_time / update_time时间戳生产者发送前先写入这张表状态为发送中确认回调收到 ack 后更新为已确认消费者消费完成后回调更新为消费成功。定时任务扫描两类异常数据一类是发送中超过 5 分钟的重新发送一类是已确认但消费超时未更新的人工或自动补偿。这张表的存在让 RabbitMQ 的任何异常都不会成为消息不可追踪的借口。这套方案的落库成本在实际业务中完全可以接受因为我只需要为关键业务消息建表而不是所有流水都做。对于不重要的消息可以走轻量级方案只做 confirm 手动 ACK不做本地状态表。最后聊一句实在话做消息可靠性设计不要迷信任何单一机制。RabbitMQ 的持久化、confirm、手动 ACK、镜像队列这些东西单独拿出来都不复杂但组合成一个完整方案时很多细节才真正暴露出来。我最后再分享一个经历过的小教训有一次线上消费者偶发报错我把异常吞了只打了 warn 日志没做 nack结果几百条消息在消费者进程重启后被自动 ACK 确认业务数据永久缺失只能靠凌晨的数据库备份回滚补救。从那之后我给自己立了一条规矩凡是手动 ACK 的消费者一定要有异常 notify 消息落盘宁可重复消费也不能静默吞掉。时刻记住MQ 只是消息搬运工真正的数据可靠性最终要靠业务系统自己盯着。