ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

RocketMQ消息丢失排查全链路:从sendResult到offset,一文讲透

RocketMQ消息丢失排查全链路:从sendResult到offset,一文讲透 遇到“MQ消息丢失”的反馈绝大多数人的第一反应是怀疑网络闪断但我在生产环境排查过的RocketMQ丢消息案例里网络恰恰是最不常见的原因。真正的问题往往藏在三个不起眼的地方send方法返回的状态码、Broker的刷盘与主从复制配置、以及消费者提交offset的顺序。这篇文章不是复述官方文档而是把我处理RocketMQ消息丢失问题时的排查思路、配置调整和踩坑经历整理出来。如果你正在用RocketMQ承载订单、支付、积分这类核心链路或者面试时被问到“MQ消息丢失怎么办”这篇文章应该能给你一些能直接落地的参考。我会从端到端的链路拆解开始把三个环节各自的丢失场景、关键配置和应对手段讲清楚最后用一个真实的排查案例把整个思路串起来。1. 消息在RocketMQ里到底可能丢在哪一环节1.1 一条消息从生产到消费要经过的几道门先看一条消息在RocketMQ里的完整旅程生产者调用send方法把消息发到BrokerBroker将消息顺序写入CommitLog文件同时写入PageCache内存缓存然后按刷盘策略落盘如果配置了主从复制主节点还要把消息同步给从节点。消费者通过长轮询从Broker拉取消息执行本地业务逻辑最后提交消费进度offset。这个过程可以粗暴地切成三段生产端、Broker端、消费端。每一段都有对应的丢失风险。做排查时我习惯先问自己一个问题业务抱怨的“消息丢了”到底是生产端根本没有发出去还是Broker收到了但没存住又或者是消费端拉到了但没消费成功这三个问题的排查路径完全不同如果一上来就翻消费者日志很可能白忙一场。有个类比我觉得很贴切消息从生产到消费就像寄快递。你把包裹交给快递员生产端到Broker快递公司负责运输和暂存Broker存储与复制快递员派送到你手上Broker到消费端。收件人说“我的快递丢了”可能是在任何一环丢的你需要看物流轨迹才能定位。1.2 三个环节里最常见的“假成功”陷阱在讲具体方案之前我想先把最常见的三种“假成功”现象列出来因为它们最容易让人误判环节典型丢失原因常见“假成功”表现核心规避手段生产端发送时Broker返回了非成功状态但代码没判读send()没有抛异常就当发送成功了判读sendResult非SEND_OK进入补偿流程Broker端异步刷盘时进程宕机PageCache数据没落盘客户端已经收到成功响应同步刷盘 主从同步复制或使用Dledger消费端先提交offset再处理业务业务失败后无法重试消息在控制台显示“已消费”实际业务没成功业务处理成功后再提交offset失败返回重试还有一点容易被忽略业务方嘴里说的“消息丢了”有时不是物理消失而是下游没拿到、或拿到了没消费成功最终数据对不上。排查前先和业务确认“丢”的具体表现比如订单有积分没到账、报表少了一条统计、短信没发出来这些现象对应的排查方向差异很大。2. 生产端丢消息别把sendResult不当回事2.1 不同发送方式对消息可靠性的影响生产端发送消息有三种常见方式可靠性从上到下递减发送方式是否关心结果丢消息风险适用场景同步send返回SendResult最低核心业务消息异步send回调触发中需处理回调对RT敏感但要求可靠sendOneway完全不关心最高日志、监控等允许丢失的消息我见过不少线上事故都是图省事用了sendOneway或者用了异步发送但SendCallback里只写了onSuccessonException为空甚至只打一行日志。生产环境里核心业务消息一律用同步send或者用异步send并确保回调里有完整的失败处理逻辑这个原则应该当成代码评审的硬性要求。同步发送看似简单其实也有隐患那就是对SendResult的判读。很多同事写的代码长这样try { producer.send(msg); // 没抛异常以为成功了 } catch (Exception e) { // 记录失败 }问题在于RocketMQ的同步发送在没有抛出异常的情况下返回结果不一定成功这就要看SendStatus了。2.2 sendResult里的四个状态只有SEND_OK能让你放心RocketMQ的SendStatus有四种分别是SEND_OK、FLUSH_DISK_TIMEOUT、FLUSH_SLAVE_TIMEOUT、SLAVE_NOT_AVAILABLE。很多人在这一步踩坑以为“能返回结果”就等于“发送成功”其实后三种状态都意味着消息没有真正安全落盘或同步。SEND_OK消息已经成功写入Broker的CommitLog这是唯一能让你放心的状态。FLUSH_DISK_TIMEOUTBroker配置了同步刷盘但刷盘超时了。此时消息虽然在内存里但没确认写入磁盘。FLUSH_SLAVE_TIMEOUTBroker配置了同步复制但主节点等待从节点同步超时。SLAVE_NOT_AVAILABLEBroker配置了同步复制但找不到可用的从节点。正确的发送判读逻辑应该是这样的SendResult sendResult producer.send(msg); if (sendResult.getSendStatus() ! SendStatus.SEND_OK) { // 不要吞掉这个状态记录日志并落到本地补偿表 saveToCompensateTable(msg, sendResult.getSendStatus()); }这里有一个容易被忽略的细节后三种状态是Broker“尽力处理”后的失败反馈消息可能已经写了CommitLog也可能没有。此时你不仅要在日志里记录状态还需要有补偿机制去核对消息到底有没有真正存下来。我之前在处理一个订单积分场景时就是因为发送方把SLAVE_NOT_AVAILABLE当成成功导致连续几个小时有订单没发积分这类问题不通过状态判读很难发现。2.3 重试与事务消息什么时候该上什么手段发送失败时生产者本身有重试机制默认情况下同步发送失败会重试2次加上首次发送一共3次。你可以通过setRetryTimesWhenSendFailed调整比如核心链路我会调到5次。但要注意这个重试只针对抛异常的情况对于上面提到的非SEND_OK状态RocketMQ不会自动重试需要业务方自己判断和处理。事务消息解决的是另一个更棘手的问题本地数据库操作和发MQ之间的一致性。比如订单系统先insert订单再发一条“订单创建成功”的消息如果insert成功了但发消息失败下游就永远不知道这笔订单存在。RocketMQ的事务消息通过半消息本地事务回查机制解决TransactionMQProducer producer new TransactionMQProducer(txOrderProducerGroup); producer.setTransactionListener(new TransactionListener() { Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { // 1. 执行本地事务插入订单数据 // 2. 根据本地事务结果返回COMMIT_MESSAGE或ROLLBACK_MESSAGE return LocalTransactionState.COMMIT_MESSAGE; } Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // 3. 如果长时间没收到确认RocketMQ会回调此方法 // 4. 根据本地事务是否存在/状态决定COMMIT或ROLLBACK return LocalTransactionState.COMMIT_MESSAGE; } });但不是所有场景都需要上事务消息。事务消息会引入额外的复杂度并且本地事务执行时间会影响整体的发送RT。如果业务能接受“先落库再用定时任务扫补偿表发消息”反而更简单可靠。我见过很多团队把事务消息当万能药最后因为回查逻辑没写好导致消息状态悬在那里。事务消息适合那种“数据库写操作和MQ发送强一致”的场景比如支付成功后的账务变更而对那些可以异步补偿的业务用本地消息表更合适。3. Broker端丢消息刷盘和主从复制怎么权衡3.1 PageCache与刷盘机制异步和同步的真实差异Broker端接收消息后并不是直接写磁盘而是先写入CommitLog对应的PageCache中然后由后台线程异步刷盘。这里有个关键点写入PageCache成功Broker就可以向客户端返回成功响应了在异步刷盘模式下。异步刷盘的性能优势明显但代价是如果Broker进程在刷盘前宕机比如断电、kill -9、云服务器强制重启PageCache里还没来得及落盘的消息就全部丢了。同步刷盘则在每次写入CommitLog后强制刷盘确认落盘后才返回成功牺牲了吞吐换来了更高的数据安全性。我自己的测试里相同配置的机器4K大小的消息同步刷盘的TPS大约只有异步刷盘的一半左右。所以这里没有绝对的“对”与“错”只有取舍核心交易链路用同步刷盘日志、统计、行为分析这类允许小概率丢失的消息用异步刷盘。即使是异步刷盘也强烈建议别把RocketMQ部署在IO能力太差的机器上否则刷盘线程堆积会引发更严重的问题。这里需要区分一个概念异步刷盘丢消息发生在“进程级故障”场景。如果是正常停机RocketMQ会主动把PageCache刷到磁盘但如果是断电、强杀进程那就没有任何机会了。所以生产环境一定要避免直接kill -9能用shutdown就走shutdown。3.2 主从复制模式SYNC_MASTER不是银弹Broker支持主从部署通过brokerRole配置区分角色。ASYNC_MASTER模式下主节点写完本地CommitLog就返回成功从节点异步拉取同步SYNC_MASTER模式下主节点必须等待从节点同步完成后才返回成功。从名字上看SYNC_MASTER似乎更可靠但它有两个实际问题第一性能损耗明显每次发送都要等一次网络RTT到从节点 第二如果从节点不可用SYNC_MASTER模式下发送会返回SLAVE_NOT_AVAILABLE状态但这个状态并不会让send()抛异常客户端如果只看“有没有抛异常”就会把这条消息当成成功然后消息就“丢”了。所以如果你配置了SYNC_MASTER一定要保证从节点可用并且客户端必须判读SendStatus。还有一个思路是使用RocketMQ 4.5版本之后引入的Dledger模式通过Raft协议在多个节点之间自动选主真正实现自动故障切换。Dledger通常需要至少3个节点适合对可用性和一致性要求都极高的场景。3.3 我常用的Broker端可靠性配置参考我在不同业务场景下会使用不同的Broker配置基线这里直接给出一套可以“抄作业”的组合业务类型刷盘策略flushDiskType主从模式brokerRole部署形态核心交易、支付账务SYNC_FLUSHSYNC_MASTER1主1从或1主2从常规业务数据ASYNC_FLUSHASYNC_MASTER1主1从资金/强一致场景SYNC_FLUSHDledger3节点及以上写入broker.conf的示例brokerClusterNameDefaultCluster brokerNamebroker-a brokerId0 flushDiskTypeSYNC_FLUSH brokerRoleSYNC_MASTER需要特别提醒一点同步刷盘在机械硬盘上很痛苦160左右的IOPS根本扛不住建议至少使用SSD。云服务器上还要留意云盘的IOPS波动如果刷盘经常超时你会在sendResult里看到FLUSH_DISK_TIMEOUT。我曾经遇到过云平台某个可用区磁盘IO抖动导致核心链路大量FLUSH_DISK_TIMEOUT的情况这种时候光调客户端没用得从存储层下手。4. 消费端丢消息offset提交姿势、重试与幂等4.1 消费端“丢消息”的真相offset的提交顺序消费端的“丢消息”场景和大多数人想的不太一样消息其实已经从Broker拉取到消费者本地了只是业务没处理成功但消费进度offset已经往前移动了。一旦offset提交Broker就认为这条消息已经消费完成后续不会再投递消息就“没了”。默认情况下RocketMQ的PushConsumer是处理完业务逻辑后自动提交offset也就是说你在MessageListener里返回CONSUME_SUCCESS它就认为消费成功并提交offset。问题出在一些不走默认通道的用法上比如consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) - { try { handle(msgs); // 业务处理 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } catch (Exception e) { log.error(消费异常, e); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 错误示范吞异常假装成功 } });这个错误示范很常见业务处理抛了异常但代码捕获后仍然返回CONSUME_SUCCESS结果就是消息offset被提交业务数据没落库看起来像消息丢了。正确的姿势是捕获异常后返回RECONSUME_LATER让RocketMQ重新投递。4.2 消费重试与死信队列失败消息的兜底网消费失败返回RECONSUME_LATER后消息会进入重试队列。默认情况下RocketMQ会重试16次具体版本和配置可能有差异重试间隔从秒级逐渐拉长直到小时级。如果重试结束后仍然失败消息会进入死信队列Topic名是%DLQ%消费组名。很多人不重视死信队列实际上它是消费端的最后一道兜底。我建议给死信队列单独做监控和告警因为消息能进死信队列说明已经有业务一直在失败。我遇到过一种特别坑的情况消费者代码升级后新增了一个字段反序列化失败所有消息都消费失败进入无限重试最后全部进死信队列。由于没有盯死信队列告警这个故障持续了大半天才被业务发现。如果你在下游系统无法快速修复的情况下可以先暂停消费者或者用工具把死信队列的消息重新投递到原Topic等修复后再消费。但这些都属于“事后补救”更重要的还是消费代码里做好异常分类可重试异常如下游超时返回RECONSUME_LATER不可重试异常如消息格式错误直接记录并舍弃或者发到专门的告警Topic让人工介入。4.3 至少一次投递语义下的幂等设计RocketMQ的投递语义是at least once至少一次而不是exactly once。也就是说消费者有可能重复收到同一条消息消费成功后还没来得及提交offset就宕机、重试队列重新投递、Consumer发生Rebalance等都可能导致重复消费。所以处理消息丢失问题时必须同步考虑幂等设计。一个合格的消费逻辑应该让“重复消费”对业务无感。常用的幂等方案有数据库唯一约束。以业务唯一键如订单号建唯一索引插入冲突时说明已处理过。Redis的setIfAbsent。利用消息ID或业务ID作为key处理前先写入写入成功才执行业务。状态机校验。先查业务状态符合预期才更新更新时带上状态条件。举个简单的Redis幂等示例String key msg:consumed: msgId; Boolean first redisTemplate.opsForValue() .setIfAbsent(key, 1, Duration.ofHours(24)); if (!Boolean.TRUE.equals(first)) { // 已消费过直接返回成功 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } // 执行业务逻辑需要注意Redis方案有一个前提setIfAbsent和业务执行要视为一个整体来设计否则可能会出现setIfAbsent成功了但业务失败的情况导致消息无法被重试。更稳妥的做法是在业务逻辑成功后再设置幂等标记或者用本地事务表记录处理状态。5. 丢消息定位实战一次从业务反馈到根因修复的完整链路5.1 先搞清楚业务口中的“丢消息”到底丢在哪之前处理过一个电商场景用户下单后订单系统会发一条MQ消息给积分服务用于累计积分。业务方反馈“偶尔有订单没有积分记录”看起来像消息丢了。我当时的第一反应不是去看网络而是先让业务方提供几个具体的单号然后用这些单号去RocketMQ控制台查询消息轨迹。如果你在控制台输入Message ID或业务Key能查到轨迹说明消息至少到了Broker如果轨迹完整但消费者没消费记录问题就在消费端如果连轨迹都没有那就要往生产端排查了。那次查单号的结果是部分订单的消息在控制台完全查不到。这说明问题大概率出在生产端而不是消费端排查方向瞬间清晰了很多。5.2 用消息轨迹把问题锁定到具体环节“消息轨迹”功能是排查丢消息最有力的工具。Broker端需要在broker.conf中开启traceTopicEnabletrue生产者和消费者在创建时都需要设置setEnableMsgTrace(true)这样控制台才能展示消息从生产、存储到消费的完整轨迹。这里要提醒一点消息轨迹本身也是一条消息也会有开销。如果Topic的TPS很高不建议全部开启trace可以只对核心Topic开启或者做采样。但像订单、支付这类核心链路我强烈建议开启因为排查问题省下的时间远超这一点点性能损耗。回到案例。用出问题的订单号查轨迹发现那些丢失的消息连“生产轨迹”都没有——这说明消息根本没有到达Broker问题被锁定在了生产端。接下来就是翻代码和日志了。5.3 顺着sendResult揪出根因SLAVE_NOT_AVAILABLE查看订单服务日志后发现发送消息时并没有抛异常SendResult里却出现了SLAVE_NOT_AVAILABLE。也就是说消息被发到了Broker主节点但主节点配置的是SYNC_MASTER需要同步到从节点才算成功而此时从节点不可用所以返回了这个状态。业务代码的问题也随之暴露只catch了Exception没有判读sendResult的SendStatus。因为在很多开发者的认知里“没抛异常发送成功”这个误区在这里直接导致了消息丢失。再往下查从节点不可用的原因是从节点磁盘满了Broker在检测到从节点落后或不可用后同步复制自然就无法完成。主节点本身是正常的它只是想同步给从节点但同步不了于是给客户端返回了SLAVE_NOT_AVAILABLE。从业务视角看这就是“消息丢了”。这个案例的根因是两层叠加第一层是基础设施故障从节点磁盘满第二层是应用代码缺陷不判读sendResult。如果任何一个环节做得好消息都不会丢。5.4 修复、补偿和验证修复分了三步走第一步先处理基础设施清理从节点磁盘空间确认主从同步恢复正常。这一步要快因为它决定了后续消息是否能正常同步。第二步修改生产端代码逻辑强制要求判读SendStatus不为SEND_OK时写入本地补偿表并记录完整的业务上下文。if (sendResult.getSendStatus() ! SendStatus.SEND_OK) { compensateService.save(msg, sendResult.getSendStatus()); }第三步处理已经丢失的消息。因为丢失前有业务日志我们按时间范围和订单状态捞出了受影响的消息清单写了个临时脚本重新发送补发了积分。验证环节也不能省我们连续观察了一个大促周期的数据对比了订单表和积分记录确认没有新的丢失后才彻底放心。整体来看这次事故从发现到完全解决花了半天大部分时间都花在定位问题上。如果一开始就开了消息轨迹定位会快很多。6. 几个容易误判成“丢消息”的场景与后续加固建议6.1 不是丢消息但看起来像丢消息的三种情况排查过程里我发现很多“丢消息”其实是假象以下几种情况很容易被误判。第一种是消费延迟或堆积。消息明明还在Broker的队列里只是消费者处理不过来。如果你在控制台看消费进度和堆积量会发现消息一条没少只是“还没轮到”。这种情况去改消费者逻辑或者增加消费者实例就行别往丢消息方向排查。第二种是消息进入了重试队列。消费失败后消息会不断延迟重试从业务侧看下游迟迟没反应像丢了。实际上消息一直在重试队列里等到重试成功或者最终进入死信队列才真正“被放弃”。第三种是顺序消息乱序。有时候消息没丢但因为处理乱序导致业务数据错乱看起来像少了某条消息。比如先消费了“取消订单”消息再消费“创建订单”消息数据库最终状态就是错误的导致业务方认为消息丢了。这种情况需要检查顺序消息的MessageQueueSelector实现以及是否真的把同一业务Key的消息发到了同一个队列。6.2 丢消息防护的日常巡检与兜底设计排查完一次丢消息之后我建议从“事后救火”升级到“事前预防”。日常巡检可以关注这几个关键指标巡检对象指标异常阈值参考生产端发送失败率、非SEND_OK比例统计值为0出现即排查Broker刷盘耗时、磁盘使用率、主从同步延迟刷盘耗时持续超过100ms要警惕消费端消费堆积量、重试队列数量、死信队列数量死信队列有消息就该告警除了巡检另一个很有用的兜底设计是“本地消息表”业务表和消息表放在同一个数据库事务里事务提交后由后台任务把状态为“待发送”的消息推送给MQ发送成功后再更新消息表状态。这样即使发送端程序崩溃重启后也能通过扫描消息表把漏掉的消息补发出去。我自己一直保持一个习惯每次修改Broker配置或者做集群升级之后都会主动往线上核心Topic发一条带唯一ID的探针消息然后把这条消息的完整轨迹截图存档。下次再有人喊“消息丢了”我至少能快速确认链路是否是通的。这个习惯帮我省过很多次排查时间也是一种很实用的基线验证手段。
返回列表