ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

消息队列消息不丢失的七层防线:从生产者到消费者的全链路保障

消息队列消息不丢失的七层防线:从生产者到消费者的全链路保障 1. 一次“消息静默消失”的线上事故问题到底出在哪一环先讲一个真实发生过的案例这个案例几乎包含了我后面要说的所有坑。凌晨两点线上告警群里突然热闹起来。用户反馈支付成功扣了款但积分没到账、短信也没收到。订单服务日志显示支付成功消息已发出消息中间件控制台也显示消息已存储消费端却什么都没干。我第一反应是消费代码挂了拉日志发现消费端压根没收到这条消息。于是顺着链路排查最后定位到三个事实生产者用的是异步发送且没有回调处理发送失败时只打了一行 debug 日志没人看见。Broker 采用默认异步刷盘策略消息写入页缓存就返回成功当时所在的物理机恰好断电重启未落盘的数据全部丢失。消费者开启了自动提交 offset业务线程在处理消息时抛异常退出offset 却已经提交了。这条消息既不重试也不告警直接“蒸发”。这个案例里生产者、Broker、消费者三端各漏了一道防线消息就彻底找不回来了。做消息队列的人常说“消息不丢失”是个系统工程不是调某一个参数就能解决。消息从业务系统产生开始到最终被消费方成功处理并落库中间要经历生产者发送、Broker 存储、消费者拉取与确认三个大阶段每个阶段都有独立的丢失风险。也就是说消息不丢失 生产者不丢 Broker 不丢 消费者不丢三者缺一不可。这篇文章我不会只讲概念而是把七个关键防护点一层层拆开每层对应什么风险、什么配置、什么代码写法以及在真实项目中踩过的坑和验证方法一次说透。读这篇文章之前假定你已经了解消息队列的基本概念比如主题、分区、消费组、offset 这些词。如果对 Kafka 和 RocketMQ 的配置不熟也没关系我会把两套主流 MQ 的对应实现都拉出来对比你只需要掌握思路换到任何 MQ 体系都能用。2. 生产者端的三道闸门把“发后即忘”变成“确认收到”绝大多数消息丢失事故的源头在生产者侧因为开发者最容易在这里使用“发后即忘”的方式。发后即忘并不是说消息一定丢而是当发送失败发生时你没有任何感知。我见过不止一个项目生产者发送消息的代码就是一行producer.send(record)异常全部吞掉。这在低并发、网络稳定的环境里可能谁也发现不了问题但一旦 Broker 重启、网络抖动或者 topic 不存在消息就会悄悄丢掉。生产者端要做三道防线核心思想是每一次发送都必须有明确的成功或失败结论并且对失败要有补偿手段。2.1 第一层同步发送或用回调感知结果所谓“确认收到”是要求生产者能够拿到 Broker 的确认结果。Kafka 的 Producer 有两种常见写法// 不推荐发后即忘 producer.send(new ProducerRecord(order-event, orderId, message)); // 推荐带回调的异步发送 producer.send(new ProducerRecord(order-event, orderId, message), (metadata, exception) - { if (exception ! null) { log.error(消息发送失败topic{}, key{}, order-event, orderId, exception); // 进入补偿流程见第二层 } else { log.info(消息发送成功partition{}, offset{}, metadata.partition(), metadata.offset()); } });如果是同步发送也简单直接判断返回值。RocketMQ 里写法更直观SendResult sendResult producer.send(message); if (sendResult.getSendStatus() ! SendStatus.SEND_OK) { // 处理失败 }用回调或者同步发送核心是拿到发送结果而不是干等或者不管。有人会问高性能场景下同步发送不是会很慢吗这就是典型的需求权衡问题。对于订单、支付这类必须保证不丢的消息哪怕多付出 1 毫秒的延迟也值得对于日志采集、行为埋点这类允许少量丢失的数据用发后即忘可以换取吞吐这是完全合理的。关键是你要意识到自己舍弃了什么而不是无意识地丢掉关键数据。2.2 第二层重试机制与补偿表不放过一次瞬时失败拿到发送失败结果之后下一步是重试。Kafka 生产者自带重试机制核心配置是retries和retry.backoff.msacksall retries3 retry.backoff.ms300 max.in.flight.requests.per.connection1max.in.flight.requests.per.connection1很关键它限制了在单个连接上未确认请求的最大数量。如果不设这个值重试可能会导致消息顺序错乱——第一条消息发送失败第二条消息却先发出去了等第一条重试成功时顺序就颠倒了。对于强顺序要求的场景必须设置为 1。RocketMQ 的同步发送失败后你也可以自己封装重试for (int retry 0; retry 3; retry) { try { SendResult result producer.send(message); if (result.getSendStatus() SendStatus.SEND_OK) { break; } } catch (Exception e) { Thread.sleep(300L * (retry 1)); } }但重试不是万能的。Broker 若真正宕机重试十次也没用。所以我还习惯在数据库里建一张消息补偿表结构大概是这样的字段说明id主键business_key业务唯一键如订单号topic目标主题payload消息内容status0: 待发送, 1: 已发送, 2: 已确认, 3: 发送失败retry_count已重试次数next_retry_time下次重试时间发送消息前先落入本地事务把业务操作和消息状态更新放在同一个数据库事务里也就是经典的“本地消息表”方案。再通过一个定时任务扫描 status 为 0 或 3 且超过 next_retry_time 的记录重新投递。这保证了只要本地业务成功了消息至少不会丢最坏情况是延迟到达。2.3 第三层事务消息解决“业务成功但消息发送失败”的原子性问题本地消息表需要额外建表、写定时任务有些人觉得麻烦于是 RocketMQ 直接提供了事务消息机制。它的流程我简单描述一下生产者发送 half message此时对消费者不可见。执行本地业务事务。如果本地事务成功提交消息失败则回滚消息。如果步骤 3 因为宕机等原因没有执行Broker 会回查生产者本地事务状态根据结果决定提交或回滚。代码大致长这样以 RocketMQ 为例TransactionMQProducer producer new TransactionMQProducer(); producer.setTransactionListener(new TransactionListener() { Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { // 执行本地业务比如写订单表 try { orderService.createOrder((OrderDO) arg); return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { return LocalTransactionState.ROLLBACK_MESSAGE; } } Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // 回查本地事务状态 return orderService.isOrderCreated(msg.getKeys()) ? LocalTransactionState.COMMIT_MESSAGE : LocalTransactionState.ROLLBACK_MESSAGE; } });Kafka 也有事务 API但用起来复杂而且业界用 Kafka 做事务消息的比例明显低于 RocketMQ。如果你用的是 Kafka我更推荐本地消息表方案因为它通俗易懂也不依赖特定版本特性。这一层防线做的事用一句话总结把消息发送和业务写库放进同一个“事务决策”里杜绝“钱扣了消息没发出去”的惨案。3. Broker 端的两道关卡写入可靠与存储可靠缺一不可过了生产者这一关消息已经送到 Broker 手里了。但 Broker 并不是保险柜它自身也面临两个风险第一收到消息后只写了内存/页缓存就返回成功机器断电就丢第二数据虽然写入了本地磁盘但磁盘损坏或者机器报废数据依然跟着没了。这两类风险对应两层防线。3.1 第四层同步刷盘让消息真正落在磁盘上先讲一个小原理。几乎所有的 MQ 写入消息时并不是直接写磁盘文件而是先写入操作系统的页缓存Page Cache然后由操作系统异步刷盘。这样做的目的是利用内存的高性能来换取吞吐量但代价是如果写入页缓存还没刷盘时进程崩溃或机器断电数据就丢了。以 RocketMQ 为例Broker 的刷盘策略有两种刷盘策略行为可靠性性能ASYNC_FLUSH写入页缓存即返回成功低断电可能丢高SYNC_FLUSH写入磁盘文件后才返回成功高较低但可接受配置位于 Broker 的broker.confflushDiskTypeSYNC_FLUSHKafka 的刷盘配置不太一样它更依赖副本机制来保证可靠性而本机刷盘由log.flush.interval.messages和log.flush.interval.ms控制。默认值是比较宽泛的如果希望更可靠可以调小这些值但代价是引入更多次磁盘写入。我在生产环境里的经验是订单、支付、账户等核心链路用同步刷盘日志、行为数据用异步刷盘。不要一刀切全上同步刷盘那样会把日志型的高吞吐场景拖垮也不要在核心链路上贪图性能用异步刷盘因为一次断电就能让你损失一批关键消息业务恢复成本远比省下的那点性能高。3.2 第五层多副本与 ISR坏一台机器也不丢数据刷盘只能防断电防不了磁盘损坏、机器报废。这时候需要副本机制。Kafka 的副本机制核心概念是 ISRIn-Sync Replicas同步中的副本。Leader 分区的数据会同步到多个 Follower 副本Producer 写入时可以通过acks参数控制需要多少个副本确认acks0不等待确认最多丢。acks1Leader 写入成功就返回Leader 宕机可能丢。acksall所有 ISR 副本都写入成功才返回最安全。配合min.insync.replicas参数可以设定最少需要几个副本同步成功。最稳妥的组合是acksall min.insync.replicas2意思是至少要有一个 Follower 和 Leader 保持同步写入才算成功。如果 Follower 全部挂掉写入会失败而不是静默接受这逼着生产者走重试或补偿。RocketMQ 也有对应概念主从模式下 Broker 的配置brokerRoleSYNC_MASTERSYNC_MASTER表示主节点需要等待从节点复制成功后才返回。换成ASYNC_MASTER则会丢掉这条保护。这里特别提醒一个 Kafka 的隐藏坑unclean.leader.election.enable一定要设为 false。如果设为 true当所有同步副本都挂了Kafka 可以选一个不同步的副本当 Leader——这样做的好处是服务可用性提高了但代价是消息大量丢失。对“消息不丢失”有刚需的业务这个参数必须关掉。宁可短暂不可用也不能让数据悄悄丢。4. 消费者端的手动确认自动提交是丢消息的隐形杀手我在最开始那个案例里说消费者这端也丢了消息很多人不理解消费者不是只负责收消息吗怎么会丢问题就出在 offset 的提交时机上。偏移量offset可以理解成书签。消费者读完一批消息后要把书签记录到 Broker下次拿着书签继续读后面的消息。如果在消息处理完之前就更新了书签等于书签已经翻页内容还没看懂。此时消费者进程崩溃重启后会从书签位置继续消费——处理失败的那批消息永远不会再读到了。Kafka 消费者默认是自动提交 offset 的相关配置是enable.auto.committrue auto.commit.interval.ms5000每 5 秒自动提交一次。如果你的业务逻辑处理时间超过 5 秒消息处理到一半offset 已经提交了或者代码在poll()之后for循环里处理消息第 1 条处理失败抛异常后续的几条压根没执行但 offset 还是被自动提交了。这就是不丢消息的大忌。正确做法是把自动提交关掉改成手动提交而且要在消息处理成功后再提交props.put(enable.auto.commit, false); while (running) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { try { process(record); // 业务处理 consumer.commitSync(); // 处理成功后再提交 } catch (Exception e) { log.error(消费失败等待下次重试或进入死信流程, e); // 注意这里 continue不要提交 } } }RocketMQ 的消费者机制不太一样它是通过返回状态来决定是否确认的consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) - { try { businessService.process(msgs); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } catch (Exception e) { return ConsumeConcurrentlyStatus.RECONSUME_LATER; } });返回RECONSUME_LATER的消息会在之后的重试队列中再次尝试而不是直接被丢弃。这里最忌讳的就是不管业务是否成功一律返回消费成功这在 RocketMQ 中极为常见尤其是消费逻辑里 catch 住了异常并“吞掉”从现象上看消息全部消费成功实际业务数据全没落库。所以消费者这一层防线的核心是你必须在业务处理真正成功时才确认消息。如果拿不准宁可返回失败让它重试也不要模棱两可地确认。重试顶多带来重复消费而你还有最后一道防线兜底。5. 第七层防线重试、死信与幂等把“最后一公里”焊死即使你做到了生产者确认、Broker 同步刷盘、多副本、消费者手动提交消息依然有可能重试。为什么因为“不丢失”和“不重复”是两回事。几乎所有主流 MQ 保证的是 At Least Once至少一次投递也就是说消息可能不丢但可能重复。当消费者的业务代码处理消息后没有来得及提交 offset进程就崩溃了重启后会重新消费这一条——这就产生了重复。应对重复的方法不是去消灭它而是让重复消费变得无害这正是第七层防线的意义。5.1 消费失败的重试策略与死信队列消费者拿到消息后处理失败怎么办第一步是重试。Kafka 场景下需要自己维护重试逻辑。最简单的方法是把消费失败的消息写入一个本地待重试表用定时任务扫描重发。更讲究的做法是利用 Kafka 的重试主题设计不同延迟级别的重试 Topic比如 1 秒、10 秒、60 秒分别投递一次。RocketMQ 则内置了重试队列消费失败的消息会按照延迟等级自动重试默认最多 16 次。重试次数用完之后消息去哪不能直接丢掉而应该送进死信队列DLQ。RocketMQ 有自动死信队列消息重试耗尽后就进入%DLQ%消费组名你可以在控制台查看、手动介入。Kafka 没有死信队列的概念需要自己实现消费者在多次重试失败后把消息发给一个专门的dlq-topic再由人工处理程序或者告警介入。我在项目中给死信队列配了一套告警规则只要死信主题有消息产生就触发企业微信/钉钉通知并附上最后一批消费失败的原始消息内容。因为死信队列里躺着的基本都是脏数据或者业务 bug不及时处理就会演变成资损事故。有些团队会把死信队列的消息定期删掉那是非常危险的等于把日志销毁了后面根本没法追溯。5.2 幂等兜底同一个消息重复消费也不怕这里必须聊幂等。所谓幂等就是同一个操作执行一次和执行多次结果一致。常见的做法有四种基于数据库唯一索引。比如消费消息后要插入一张积分流水表给业务的唯一键比如 orderId建唯一索引重复插入会被数据库拦住直接跳过。基于 Redis SETNX。消费前先SETNX consume_lock_{orderId} 1拿到锁才处理处理完再删掉。注意设置锁的过期时间防止服务宕机锁不释放。基于状态机。比如订单状态流转待支付 → 已支付 → 已发货。重复消费时如果发现订单状态已经不是“待支付”了说明已经处理过直接跳过。基于消息内业务键去重。在本地维护一张已处理消息表记录 messageId处理前先查这个表。我做消息消费的第一原则是不管前面几层做得有多好消费者的落库操作都必须有幂等保护。因为重复消息是客户端无法完全避免的哪怕你前面配置拉满在极端情况下比如消费者处理完、提交 offset 前的网络分区依然会产生重复。幂等是最后一道防线这道防线没守住前面的努力都可能白费。这里还要提醒一个很容易踩的坑如果你在消费逻辑里先做了各种判断最后发现是重复消息就 return那你需要确保 return 之前把这条消息的商品事务、缓存、统计全部一致地跳过。我见过一个系统数据库唯一索引防住了重复插入但漏了 Redis 缓存的更新导致重复消息插入被挡住但缓存状态被覆盖反而产生了数据不一致。幂等方案必须覆盖到所有可能触发的旁路逻辑不只是主库写操作。6. 7 层防线的落地清单从配置到故障演练的完整方案讲到这里七层防线全部出现我先把它们汇总成一张检查表然后在下面给出落地建议。层级防线名称关键操作常出问题的默认值1生产者发送确认同步发送或回调异步无回调2生产者重试补偿设置 retries、本地消息表无重试、无补偿3业务事务原子性本地消息表或事务消息业务成功后消息丢失4Broker 持久化同步刷盘异步刷盘5多副本同步副本数≥2、同步复制单副本或异步复制6消费者手动确认业务成功后提交 offset自动提交7重试、死信与幂等重试策略死信告警幂等表无限重试或直接丢弃这张表可以当成你接手任何一个消息项目的第一份排查清单。上生产环境之前我建议你按这个顺序逐项对照检查配置。6.1 一套推荐的 Kafka 配置模板如果你用的是 Kafka并且消息属于“一定不能丢”的类型可以直接参考这套配置# 生产者端 acksall retries5 retry.backoff.ms300 max.in.flight.requests.per.connection1 enable.idempotencetrue # Broker 端 min.insync.replicas2 unclean.leader.election.enablefalse log.flush.interval.messages10000 log.flush.interval.ms1000 # 消费者端 enable.auto.commitfalse auto.offset.resetearliest这里多说一句enable.idempotencetrue的作用。它让生产者具备幂等能力即使客户端重试发送Broker 也能识别重复消息并避免重复写入。它和acksall配合使用能基本消除生产者重试导致的重复消息。注意开启幂等发送后max.in.flight.requests.per.connection会默认被设置为 5如果你同时有顺序性要求还是需要显式设回 1。6.2 一套推荐的 RocketMQ 配置模板RocketMQ 场景下重点在 Broker 和消费者# Broker 配置 brokerRoleSYNC_MASTER flushDiskTypeSYNC_FLUSH # 消费者注意别吃异常重试次数按业务设置 consumer.setConsumeTimeout(15); consumer.setMaxReconsumeTimes(16);RocketMQ 还有一点和 Kafka 不一样虽然 Broker 用了同步刷盘、同步复制但 producer 如果用了单向发送sendOneWay或者异步发送没有正确处理SendCallback前面的努力等于白费。所以 RocketMQ 的核心链路我全部使用producer.send()同步发送并捕获异常。6.3 怎么验证消息真的不丢配置完之后怎么验证你的防线是否有效光看配置不叫验证要做故障演练。我常用的验证方式分为三个步骤第一步断开 Broker 网络验证生产者补偿逻辑。挑一台测试环境上的应用用iptables禁用其访问 MQ 的端口观察生产者是否按预期报错、重试、进入本地补偿表。恢复网络后检查补偿任务是否把积压的消息全部补发成功且消息顺序没有乱。第二步kill -9 模拟 Broker 宕机验证消息已刷盘。往队列里连续写入一批核心消息在数据刚写入但可能未刷盘的时刻强制断电测试环境可以直接断虚拟机电源重启 Broker 后检查消息是否存在。同步刷盘模式下写入成功的消息必须全部还在。如果你用的是异步刷盘你会发现确实丢了最后一批——这就是为什么我反复强调核心链路用同步刷盘。第三步消费者 kill -9 模拟处理中崩溃验证 offset 提交时机。开启一个大批量消费任务在处理一批消息的中途 kill -9 消费者进程重启后观察之前那条正在处理的消息是否被再次投递。如果配置了手动提交且没有提前提交 offset那这条消息会被重新消费此时你的幂等逻辑应该保证数据不重复。这三步做完心里才有底。我见过太多团队把生产者的acksall配上就觉得消息一定不丢了直到断电演练才傻眼——Broker 刷盘策略还是默认的异步模式一把电闸就丢了几千条消息。所以验证永远不是可选项。6.4 日志、监控与告警让消息丢在“明处”最后一小节说说比技术配置更重要的东西可观测性。消息丢没丢你的系统自己要有能力感知。我的统一做法是生产者端每次发送成功打一条包含 topic、partition、offset、耗时、messageId 的日志失败打 ERROR 日志并附上完整消息内容。消费者端在消息进入和离开业务处理逻辑处各打一条日志带上 messageId 和业务唯一键。Broker 层面监控未确认消息数、死信队里堆积量、消费组滞后量Lag。设置两条告警规则死信队列有消息进来说明消费链路出问题了需要立刻介入消费组 Lag 持续增长说明消费速度跟不上需要扩容或者排查阻塞。日志的意义在于当消息真的丢了你能在第一时间定位到是在哪一环丢的。如果生产者日志显示已发送、Broker 控制台显示已写入、消费者日志显示没收到那问题大概率出在消费者订阅或者分组配置上。只要有一环的日志缺失或不准排查就会变成大海捞针。我在实际项目中还养成一个习惯对每个核心消息都会生成全局唯一的 messageId从生成到发送、存储、消费、落库全程透传到日志和数据库表。这样不管哪一环出问题都可以用 messageId 把整条链路的日志串起来看。这个习惯帮我节省过太多排查时间。最后分享一点个人习惯文章写到这里7 层防线和落地方法都讲完了。最后还是想补充一点非常个人的经验接手任何一套带消息队列的系统我都会先向团队问三个问题——消息丢了有没有感知有没有补偿和重试兜底重复消费有没有幂等保护如果三问里有一问答不上来那这个系统就没有达到“消息不丢失”的标准。七层防线不是让你全部堆满而是让你在每一层的取舍上都有明确依据。核心业务我建议一层都不要省非核心业务可以按风险与成本做裁剪。这样一来既避免了“一刀切”的性能浪费也杜绝了“裸奔”式的数据丢失。希望这篇文章能给你带来一点实实在在的参考。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进