ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

RabbitMQ消息不丢失的7层防线:从生产到消费全链路可靠性实践

RabbitMQ消息不丢失的7层防线:从生产到消费全链路可靠性实践 做消息中间件最怕的不是慢而是消息丢。三年前我接手一个订单回调链路凌晨两点被告警叫醒查了一圈生产者日志显示发送成功Broker 也确认消息进队列了消费者却始终没拿到消息。最后发现网关进程在发送后还没来得及落盘就崩溃了Broker 只是把消息收进了内存节点一重启整条消息凭空消失。从那一刻我就明白消息不丢失从来不是一个点的事而是生产者、Broker、消费者三个阶段必须同时设防任何一段的“想当然”都会变成线上事故。这篇文章从一个很常见的部署形态说起Windows 上的网关服务作为生产者接入请求通过 RabbitMQ 投递给 Linux 上的消费者服务处理结果写入 MySQL。这个架构看起来简单真要做到消息全程不丢其实需要从发送端一直防御到消费端我把它拆成 7 层防线。适合正在用 RabbitMQ 但又对可靠性心里没底的同学也适合准备从零搭一套消息对账体系的团队。1. 消息到底是在哪一环丢的先看懂链路再谈防线1.1 一次告警把我从“发送成功”的幻觉里打醒那次事故的现场很典型。Windows 网关收到了用户下单请求生成一条 JSON 消息通过 RabbitMQ 客户端发到 Broker日志里打印了“send success”。Linux 上的消费者日志干干净净什么都没打。MySQL 里的订单状态一直停在“待支付”没有消费者去更新它。我先怀疑消费者挂了登录上去一看进程活着再去 RabbitMQ 管理台看队列发现队列不见了。继续翻 broker 日志确认节点在凌晨重启过一次重启之后很多非持久化队列直接消失。当时生产端用的还是 RabbitMQ 默认的 fire-and-forget 发送方式没有开启发布确认消息只要写进 socket 缓冲区客户端就认为发送成功。而 Broker 这边队列和消息都没做持久化内存一清数据全完。这条链路暴露了三个问题发送端对“发送成功”的定义太乐观Broker 端把可靠性完全寄托在内存上消费端因为没有数据进来连失败的机会都没有。很多人以为消息丢失是自己代码 bug实际上是把网络、进程、磁盘、脑裂等系统级问题都忽略了。1.2 全链路上最容易丢消息的四个位置我习惯把 RabbitMQ 的消息旅程分成四段每一段都有对应的丢失场景位置丢失场景默认行为生产者 → Broker 网络传输网络闪断、Broker 宕机、连接池超时客户端不关心消息直接没Exchange → Queue 路由routing key 不匹配、Exchange 不存在消息被静默丢弃Queue 内部存储消息未持久化、磁盘坏、队列未声明 durable进程或节点重启后消失Queue → 消费者处理消费者自动 ACK、业务异常被吞、处理进程崩溃消息被确认后删除这四段说到底是三个信任假设被打破你相信“send 方法返回了就成功”但 TCP 写入成功不代表 Broker 落盘你相信“队列都存在就不会丢”但内存队列挡不住宕机你相信“消费端 try-catch 包住就安全”但自动 ACK 在方法执行完就已经认为处理成功。1.3 7 层防线整体框架基于上面四个位置我把“消息不丢失”拆成 7 层相互独立的防线。每一层解决一个具体问题层与层之间不互相替代。防线序号防线名称所在阶段解决的核心问题第一层Publisher Confirm生产者确认消息真的被 Broker 接收第二层Mandatory Return Listener生产者路由不出去的消息有人接管第三层发送重试与幂等生产者网络抖动导致的重发不引发重复第四层Exchange/Queue/Message 三级持久化Broker进程重启后消息还在第五层镜像队列/仲裁队列Broker单点故障后消息仍可读第六层消费者手动 ACK消费者处理成功前消息不删除第七层消费幂等与对账消费者重复投递不会造成重复业务后文会按这三段逐一展开每一层我会把原理、配置、代码和坑都讲透。2. 生产端三道防线把“发出去了”变成“确认落盘了”2.1 第一层Publisher Confirm——Broker 说收到才算收到生产端第一件事是把“发送成功”的定义从“客户端写入了 socket”改成“Broker 返回了 ACK”。RabbitMQ 的 Publisher Confirm 机制干的就是这件事。从使用角度Java 原生客户端有两个步骤channel.confirmSelect(); channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, props, body); boolean confirmed channel.waitForConfirms(3000);如果返回 true说明 Broker 已经收到这条消息并写入了队列具体是否落盘取决于消息和队列的持久化配置。如果超时或返回 false发送端就要决定重发还是转入降级流程。要注意的是waitForConfirms是同步等待在高吞吐场景下频繁调用会拉低发送速度。更合适的做法是用异步 ConfirmListenerchannel.addConfirmListener(new ConfirmListener() { Override public void handleAck(long deliveryTag, boolean multiple) throws IOException { // 处理确认成功的消息从待确认集合中移除 } Override public void handleNack(long deliveryTag, boolean multiple) throws IOException { // 处理确认失败的消息记录并重发 } });在 Spring Boot 里则要设置spring: rabbitmq: publisher-confirm-type: correlated然后用 CorrelationData 关联消息CorrelationData cd new CorrelationData(messageId); rabbitTemplate.convertAndSend(EXCHANGE_NAME, ROUTING_KEY, message, cd); Confirm confirm cd.getFuture().get(3, TimeUnit.SECONDS); if (confirm ! null confirm.isAck()) { // 更新发送记录为 CONFIRMED } else { // 更新发送记录为 UNCONFIRMED交给补偿任务 }我建议维护一个“已发送未确认”的内存集合确认回调回来再移除。原因很简单消息确认是异步的如果你发一条消息就 new 一个 CorrelationData之后不做任何集合清理未来回调到达顺序不一致时你根本不知道哪条消息超时了补偿也就无从谈起。2.2 第二层Mandatory Return Listener——消息路由不出去也要有人管不少人觉得 Confirm 有了就万事大吉但 Confirm 只保证 Broker 收到了消息不保证消息进了正确队列。如果发送时路由键敲错或者 Exchange 和 Queue 的绑定关系因为配置变更丢失了Broker 的默认动作是静默把消息丢掉。你从发送端看Confirm 是 ACK队列里却什么都没有只能靠对账去发现。解决这个问题要靠 mandatory 参数。发送时把 mandatory 设为 trueBroker 找不到匹配队列时不再直接丢弃而是通过 Basic.Return 把消息退回生产者。Spring Boot 里开两个配置spring: rabbitmq: publisher-returns: true template: mandatory: true再注册 ReturnsCallbackBean public RabbitTemplate.ReturnsCallback returnsCallback() { return returned - { String messageId returned.getMessage().getMessageProperties().getMessageId(); String exchange returned.getExchange(); String routingKey returned.getRoutingKey(); log.error(消息路由失败, messageId{}, exchange{}, routingKey{}, messageId, exchange, routingKey); // 将消息落库为 ROUTE_FAILED 状态等待人工或脚本修正路由配置后重发 }; }实战里这条防线的主要价值不是“事后补救”而是“尽早暴露”。我遇到过一次因为灰度环境的交换机绑定关系被另一个发布任务覆盖导致生产环境部分消息路由失败如果没有 Return Listener这类问题通常要等业务方反馈说“单子没下来”才能发现而那已经过去几个小时了。把路由失败的消息落库同时打错误日志就能在几分钟内收到告警。2.3 第三层发送重试与幂等——网络抖动和重复消息一起治Confirm 超时、Broker 拒绝连接、网络闪断这些都是发送端要面对的不确定性。有了第一第二层防线后第三层要做的是决定“失败后怎么办”。失败的常规处理是重试。Spring AMQP 自带发送重试机制spring: rabbitmq: template: retry: enabled: true initial-interval: 1000 max-attempts: 3 max-interval: 10000 multiplier: 2但注意这套重试只覆盖客户端抛出 IOException 的场景不会覆盖 Confirm 超时。Confirm 超时本质上不是发送异常而是发送后“没有结果”所以要么自己写重试要么靠定时补偿。我的做法是维护一张发送记录表用 status 字段标记消息当前状态由定时任务扫描未确认和失败的消息按指数退避重发重发次数超过阈值后转入人工处理队列。重试引出的第二个问题是幂等。网络重发、超时重发都可能让同一条消息被发送多次。如果每次重发都生成一个新的 messageId消费端就不知道这些其实是同一条消息后面的幂等防线会直接失效。所以重发时一定要复用原来的 messageId所有消息内容保持完全一致。发送记录表里要把 messageId 作为唯一键从根上保证同一条消息只有一个业务身份。3. Broker 端两道防线持久化和高可用的真实含义3.1 第四层Exchange、Queue、Message 三级持久化缺一不可Broker 端的防线只有一句话你声明的资源和你发的消息全都要能扛住重启。RabbitMQ 的持久化是三个独立开关拼起来的缺一个等于没做。第一个开关是 Exchange 持久化channel.exchangeDeclare(order.exchange, BuiltinExchangeType.DIRECT, true);第二个开关是 Queue 持久化channel.queueDeclare(order.queue, true, false, false, null);第三个开关是消息持久化AMQP.BasicProperties properties new AMQP.BasicProperties.Builder() .deliveryMode(2) .messageId(messageId) .build(); channel.basicPublish(order.exchange, order.key, properties, body);这里最容易误导人的是很多人以为 Queue 设成 durable 就够了结果消息发出去时没有显式设置 deliveryMode默认仍然是 transient。RabbitMQ 重启后 Queue 还在Queue 里的消息却没了。反过来也一样消息是 durable 的但 Queue 是 transient 的消息同样没地方待。我见过线上事故清单里排第一的就是这个组合错配。还要注意一点如果你声明队列用的参数和已有队列不一致RabbitMQ 会报错或者静默忽略你的新声明。比如初始化代码里把 durable 从 true 改成 false而旧的真持久化队列还在你的代码不会收到任何异常但线上队列仍然保持旧的持久化属性。所以在做资源声明变更时一定要确认旧队列已下线。3.2 第五层镜像队列与仲裁队列跨节点冗余怎么选单节点持久化只能防进程崩溃防不了磁盘损坏或整机故障。要做到高可用下的消息不丢失必须把消息副本放到多个节点。RabbitMQ 3.8 之前的主流方案是镜像队列3.8 之后官方推荐使用仲裁队列两者的取舍值得单独说。镜像队列通过 policy 配置rabbitmqctl set_policy ha-all ^order\. {ha-mode:all}它的工作方式是主从复制客户端只连接主节点主节点把消息复制到所有镜像节点。看起来简单但实际运行中有两个风险一是复制是同步的主节点压力大了整体吞吐明显下降二是主节点故障时镜像节点的晋升和消息对齐需要时间极端情况下存在脑裂风险官方后来也不再建议在新项目里使用。仲裁队列Quorum Queue是新一代方案基于 Raft 共识实现。声明方式MapString, Object args new HashMap(); args.put(x-queue-type, quorum); args.put(x-quorum-initial-group-size, 3); channel.queueDeclare(order.quorum.queue, true, false, false, args);仲裁队列的写入需要多数派节点确认比如三个副本至少有 2 个节点落盘成功才返回结果。一致性比镜像队列更强而且 leader 自动切换消息不会在切换期间丢失。代价是 Raft 的日志复制会带来额外开销相同硬件下吞吐比普通镜像队列要低一些但对于订单、交易这类对可靠性要求极高的场景这个代价值得付。对比项镜像队列仲裁队列一致性模型主从复制存在脑裂风险Raft 多数派确认自动选主支持但切换需要时间支持选举更快吞吐表现同步复制开销大日志复制开销也大但稳定性更好特性限制大部分普通队列特性可用不支持 exclusive 队列、transient 消息推荐场景老项目维护新项目首选3.3 存储层容易踩的坑fsync、磁盘水位和惰性队列持久化开关打开了还要知道 RabbitMQ 底层怎么持久化否则容易产生“我已经开了持久化就不可能丢”的幻觉。RabbitMQ 对持久化消息的写入流程并不是每条消息同步 fsync而是先写内存再异步批量落盘。如果节点在数据尚未 fsync 时突然断电最后几毫秒的持久化消息仍然可能丢失。仲裁队列在这方面更稳能保证返回 ACK 时数据已经写入多数派节点的磁盘。磁盘水位也是 Broker 端很容易被忽略的配置。当磁盘剩余空间低于 low watermark默认是 1GB 的绝对值或内存的某个比例RabbitMQ 会主动阻塞生产者的写操作。这是保护机制不是故障。但有人为了“避免阻塞”把水位调成 0结果磁盘写满后整个节点更彻底地不可用甚至影响其他队列的读取。我的经验是让水位告警早点发生同时在监控大盘上盯着磁盘使用率而不是试图绕过保护。消息堆积场景下普通队列把所有未消费消息全部放在内存里几千条大消息就能让堆内存吃紧。如果队列的特性决定消费速度跟不上生产速度建议把队列声明为惰性模式让消息尽量落盘减少内存压力MapString, Object args new HashMap(); args.put(x-queue-mode, lazy); channel.queueDeclare(slow.consumer.queue, true, false, false, args);惰性队列的读写性能会下降但换来的是 Broker 不会被积压消息拖垮。真正的生产级方案还是要靠消费速度提升和背压控制不能只靠存储层硬扛。4. 消费端两道防线ACK 语义决定消息到底删不删4.1 第六层手动 ACK、nack、reject——确认前消息不会丢Broker 把消息投递给消费者后是否立即从队列删除取决于消费者使用什么确认模式。RabbitMQ 的自动确认模式会在消息发出后直接标记为已确认如果消费者在处理中途进程崩溃消息就丢了。这看起来很不合理但默认模式恰恰是 auto这也是很多人第一天上生产就遇到丢消息的根源。改用手动确认模式是消费端第一道防线。Spring Boot 配置spring: rabbitmq: listener: simple: acknowledge-mode: manual消费者方法里显式控制RabbitListener(queues order.queue) public void onMessage(Message message, Channel channel) throws IOException { long deliveryTag message.getMessageProperties().getDeliveryTag(); String messageId message.getMessageProperties().getMessageId(); try { handleBusiness(message); channel.basicAck(deliveryTag, false); log.info(消息处理成功, messageId{}, messageId); } catch (DataIntegrityViolationException e) { // 幂等冲突说明已经处理过直接确认即可 channel.basicAck(deliveryTag, false); } catch (BizRetryableException e) { log.warn(业务可重试异常, messageId{}, ex{}, messageId, e.getMessage()); channel.basicNack(deliveryTag, false, true); } catch (Exception e) { log.error(业务处理异常, 进入死信, messageId{}, messageId, e); channel.basicNack(deliveryTag, false, false); } }这里的三种动作含义完全不同basicAck处理成功Broker 可以删除消息basicNack(requeuetrue)处理失败但希望重新投递消息会回到队列头可能立即再投给同一个消费者basicNack(requeuefalse)处理失败且不重投消息进入死信队列如果配置了 DLX或被直接丢弃。有一个细节很多人会忽略requeuetrue不是“延迟重试”而是“立刻重投”。如果不加任何限制一个处理总是失败的消费者会陷入同一个线程连续收到同一条消息的死循环日志被打爆CPU 飙升。真正可控的延迟重试要靠死信队列加延迟插件或者配合消息头里的重试次数计数器由手动 ACK 决定第几次失败后不再 requeue。4.2 第七层消费幂等与本地事务——重复投递不等于重复成功即使你用了手动 ACK重复消费仍然不可避免。最典型的场景消费者处理完业务但还没发送 basicAck进程崩溃Broker 会把这条消息重新投递给下一个消费者。又比如上一条消息 requeue 之后再次投递同一个业务会被执行两次。解决重复消费的核心是幂等。对 MySQL 做写操作的业务我强烈推荐用“本地事务 唯一键”方案而不是先查 Redis 再判断。举例消费“订单支付成功”消息开启本地事务先向consume_record表插入一条记录唯一键是业务幂等键插入成功继续更新订单状态插入冲突说明这条消息已经处理过直接提交空事务并返回成功本地事务提交后再执行 basicAck。Transactional(rollbackFor Exception.class) public void handleOrderPaid(OrderPaidMessage msg) { ConsumeRecord record new ConsumeRecord( msg.getMessageId(), msg.getOrderId(), ORDER_PAID, LocalDateTime.now() ); consumeRecordMapper.insert(record); // 唯一索引order_id event_type orderMapper.updatePaidStatus(msg.getOrderId(), PAID); }这里的关键是“先插入消费记录再更新业务表”利用数据库唯一约束保证只有一个线程能成功其他线程的插入会抛唯一键冲突从而避免重复更新订单状态。如果把顺序反过来先更新订单再插入消费记录两个并发线程都通过了不存在检查各自执行了一次订单更新幂等就破了。幂等键的设计也不能想当然。如果直接用 messageId 做唯一键只能拦住“同一条消息重投”拦不住“两个不同消息但业务语义相同”的情况。比如订单表里同一个订单先收到“创建”消息后收到“状态同步”消息这两条消息的 messageId 不一样但业务上都对应同一个订单必须用order_id event_type这类业务键做唯一约束。唯一键选错等线上出现重复数据再改就麻烦了。4.3 并发消费和 ACK 的配合手动 ACK 方案在前并发部署在后两者叠加会产生一些有趣的时序问题。我趟过一次这样的坑消费者服务用默认的 SimpleMessageListenerContainer每个队列开了 10 个并发消费者prefetch 默认是 250。高峰时段队列里积压了几万条消息每个消费者线程的内存里都挂着大量 unacked 消息处理一条慢 SQL 卡住几秒整个服务的 unacked 数就破千。这时候某个消费者和 Broker 之间的连接断掉Broker 会把这批 unacked 消息全部重新入队而其他消费者再次拉取形成了雪球效应。解决办法是把 prefetch 调小并且结合业务处理耗时设置合理值spring: rabbitmq: listener: simple: prefetch: 50 concurrency: 5 max-concurrency: 10prefetch 的含义是每个消费者在收到 ACK 之前最多能持有的未确认消息数。它越大Broker 投递效率越高但消费者内存压力和重复投递的风险也越大。我对耗时在 100ms 以内的处理prefetch 一般压到 50 以内耗时超过 1s 的prefetch 只给 5~10。并发场景下还要注意数据库连接和唯一索引冲突。两个线程同时收到同一条消息时一个线程插入消费记录成功另一个线程插入时抛唯一键冲突这个冲突本身也是正常业务路径要 catch 住并且确认消息而不是把它当系统异常 nack否则又会出现死循环重投。5. 把7层防线串起来Windows 网关到 Linux 消费者的全链路落地5.1 一次完整的消息轨迹messageId 贯穿三个进程前面讲的是 7 层独立防线但真正落地时还要把整条链路串起来。我在项目里定义了一个约定网关发送消息时必须生成全局唯一的 messageId放进消息头和业务 body同时写入 MySQL 发送记录表。这个 messageId 会一直贯穿网关、RabbitMQ、消费者三个进程成为全链路追踪的唯一凭证。一次正常订单消息的轨迹应该是Windows 网关注入层收到 HTTP 请求生成 messageId落库发送记录 status0发送中网关通过 RabbitMQ 客户端发送消息到 order.exchangerouting key 为 order.paid生产者开启 ConfirmBroker 返回 ACK网关把发送记录 status 更新为 1已确认Linux 消费者从 order.queue 拉取消息日志里通过 MDC 把 messageId 打进每一条上下文消费者完成业务处理插入消费记录更新订单状态返回 ACK消费记录 status2成功。如果链路中任何一步断了messageId 都能帮助我们快速定位。比如发送记录 status1 但消费记录不存在说明消息在 Broker 中积压或者投递失败发送记录 status0 超时说明网络或 Broker 出现异常。消费者侧用 MDC 打印消息轨迹也很简单RabbitListener(queues order.queue) public void onMessage(Message message, Channel channel) throws IOException { String messageId message.getMessageProperties().getMessageId(); try (var ignored MDC.putCloseable(messageId, messageId)) { log.info(消费者开始处理消息); // 业务处理 } }这样查日志时直接grep messageId就能把三个进程的日志串成一条时间线。没有这个前提对账再多也讲不清消息到底卡在哪一环。5.2 对账表设计用 MySQL 兜住发送与消费两端7 层防线都做了也不能保证没有极端情况让消息悄悄丢失。所以最后一道保底是“可对账”。我的做法是设计两张 MySQL 表一张记录发送端事实一张记录消费端事实。CREATE TABLE msg_send_record ( message_id VARCHAR(64) PRIMARY KEY, biz_key VARCHAR(64) NOT NULL, exchange VARCHAR(64) NOT NULL, routing_key VARCHAR(64) NOT NULL, payload TEXT, status TINYINT NOT NULL DEFAULT 0 COMMENT 0发送中 1已确认 2路由失败 3重试中, retry_count INT NOT NULL DEFAULT 0, create_time DATETIME NOT NULL, update_time DATETIME NOT NULL, KEY idx_status_create_time (status, create_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; CREATE TABLE msg_consume_record ( consume_id BIGINT AUTO_INCREMENT PRIMARY KEY, message_id VARCHAR(64) NOT NULL, biz_key VARCHAR(64) NOT NULL, event_type VARCHAR(64) NOT NULL, status TINYINT NOT NULL DEFAULT 0 COMMENT 0待处理 1处理中 2成功 3死信, fail_reason VARCHAR(512), create_time DATETIME NOT NULL, update_time DATETIME NOT NULL, UNIQUE KEY uk_business (biz_key, event_type) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;对账任务的核心 SQL 是查“发送端已确认但消费端没有成功消费”的记录SELECT r.message_id, r.biz_key, r.payload, r.create_time FROM msg_send_record r LEFT JOIN msg_consume_record c ON r.message_id c.message_id WHERE r.status 1 AND r.create_time DATE_SUB(NOW(), INTERVAL 5 MINUTE) AND (c.message_id IS NULL OR c.status ! 2) LIMIT 100;这里要特别留意两个边界一是延迟队列的业务本身就可能延迟消费对账系统只告警不自动重发避免误伤二是对账的扫描周期和业务容忍时间要匹配我一般设置为每 1 分钟扫一次超过 5 分钟还没消费成功才告警既能有较快的发现速度又不会频繁触发误报。5.3 几个容易被忽略的时序细节全链路落地还涉及消息与数据库之间的先后顺序顺序错了前面的防线等于白做。先落库还是先发消息一定要先落发送记录再发消息。这样生产者在发送前崩溃发送记录 status0对账任务能发现并重发。反过来先发消息再落库消息一旦发出还未来得及记录Broker 和消费者都处理完了发送记录表里却什么都没有这条消息就是“黑户”对账无从谈起。先提交本地事务还是先 ACK应该先提交本地事务再发 basicAck。如果先 ACK消费者进程紧接着崩溃本地事务还没来得及提交消息已经被 Broker 删掉会造成业务丢失。只有先确保业务落库再 ACK才能保证 Broker 只有在业务成功之后才把消息清掉。补偿重发时用什么 messageId必须复用原始 messageId。对账任务根据原始 messageId 重发消费者通过幂等键拦截重复。如果重发时重新生成 messageId消费者会把它当成新消息处理幂等防线会直接失效。6. 实战中遇到的几个典型问题与排查思路6.1 队列消息堆积但不消费消费者日志却一直“正在处理”这个现象很迷惑队列 depth 一路上涨消费者进程的日志却一直在打“开始处理”一条条看起来都很正常。我还遇到过消费者日志里能看到处理成功但 ACK 一直没有发出的情况。排查链路先从 RabbitMQ 侧看rabbitmqctl list_queues name messages messages_unacknowledged consumers如果messages_unacknowledged很大说明消息已经发给消费者但消费者没有确认。再看消费者进程的线程 dump大概率会发现大量线程阻塞在数据库连接、外部 RPC 或者 Redis 命令上。处理慢不是问题问题是没有超时控制一旦下游缓慢整个消费线程池都被拖住。解决思路有两个方向一是把消费者与下游交互的接口统一加超时防止线程无限等待二是把 prefetch 压小避免单个消费者积压太多 unacked 消息。如果业务允许还可以把“消息拉取”和“消息处理”拆成两个线程池前者的 ACK 速度不受后者拖累。6.2 重复消费禁而不止问题出在唯一键设计有段时间我们的对账系统总是收到“消费重复”告警加了唯一索引也不管用。后来查下来问题确实出现在唯一键设计上。那张消费记录表的唯一键是message_id按道理一条消息只会成功一次。但业务方在一次补偿操作里人工修改了同一笔订单的金额并重新投递了一条 payload 不同的消息messageId 自然也是新的。但从业务角度看这是同一个订单的同一类事件应该被幂等拦截实际上却执行了两次。后来改成了biz_key event_type的唯一键比如order_no ORDER_PAID不管是重投还是人工补偿只要同一笔订单的同一类操作已经成功后续到达的消息都会在插入消费记录时被拦下。这个案例给我们的启示是幂等键一定要用业务语义而不是消息技术字段。messageId 只能证明“同一条消息”不能证明“同一件业务”。6.3 重启后消息消失检查持久化时最容易被忽略的两个开关平时检查持久化配置时大家都会记得看 Queue 是不是 durable但有两个“开关”经常被漏掉导致重启后消息没了排查半天才发现是它。第一个开关是消息本身的 deliveryMode。很多人发的消息用默认属性根本没设置 PERSISTENT队列哪怕声明了 durable消息也只是 transient。重启后布丁还在布丁里的葡萄干全没了。生产发送代码里一定要显式设置MessageDeliveryMode.PERSISTENT而且不要依赖框架默认值不同版本 RabbitMQ 客户端默认值可能会变。第二个开关是 Exchange 的 durability。Exchange 如果不 durableBroker 重启后交换机消失再往这个交换机发消息会直接报 404。更麻烦的是如果交换机已经存在且非持久化你怎么改声明代码都不生效除非删掉重建。换句话讲持久化检查不仅看 Queue还要把 Exchange、Message 两个维度一起查三个全绿才算数。最后再分享一点我自己的体会。7 层防线的每一层看起来都不复杂真正难的是一旦上了生产每个环节都可能被并发、超时、重启、版本变更这些外部因素击穿。我后来每次设计消息方案都会拿这 7 条当检查清单发布前逐条过一遍线上巡检时也盯着这 7 个指标看。RabbitMQ 本身提供了足够强大的可靠性能力能不能用成“消息不丢”取决于我们在每一层是否都不抱侥幸心理。
返回列表