ARTICLE · INTELLIGENCE

战地情报 · 详情页

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

RocketMQ消息堆积怎么办?从定位到根治的完整排查思路与实战

RocketMQ消息堆积怎么办?从定位到根治的完整排查思路与实战 “RocketMQ消息堆积了你怎么处理”这个问题我至少被问到过三次也看别人答过无数次。大部分人的第一反应是“加机器”第二反应是“重启消费者”——这两个答案在简单场景下看着没问题但真上了线特别是在大促或数据补偿任务跑批的时候往往会发现加了机器堆积还在涨重启完过两小时又开始涨。问题出在哪出在把“堆积”当成一个孤立现象而没有先搞清楚堆积的形态、位置和根因。如果面试官把问题限定在“RocketMQ”三个字上他想考察的其实不是你会不会扩容而是你有没有一套处理消息积压的完整方法论先判断堆积属于什么类型再定位是生产端太快、消费端太慢、还是Broker端读不动的瓶颈最后才是对应的处理动作。这篇文章我按自己的实战经验把整条链路拆开讲从堆积的本质、排查手段到不同场景下的处置方案再附上一个线上复盘的完整案例。无论你是正在准备面试还是真遇到了生产环境的大量积压我相信这套思路都能直接拿去用。1. 先把问题问清楚堆积是持续堆积还是瞬时积压1.1 堆积在RocketMQ里的准确定义先澄清一个基本概念。RocketMQ里消息写入CommitLogBroker端异步构建ConsumeQueue消费者提交消费位点offset消息的堆积量在数值上等于消息堆积量 Topic在Broker上的最大写入位点 - 消费者已提交的消费位点但“堆积量很大”不代表“出故障了”。RocketMQ本来就是靠内存队列来削峰填谷的生产端短暂冲高、消费端节奏稍慢Lag消费滞后量有一定波动是常态。真正的堆积指的是Lag在较长周期内持续增加并且消息延迟超出了业务可接受的范围——比如你预期用户下单后5秒内收到确认通知结果实际延迟了10分钟那才算堆积事故。1.2 两种堆积形态形态决定了处理手段我在排查堆积问题时有个习惯先不碰任何配置先看Lag曲线和消费TPS曲线的形态把堆积分成两类。第一类是瞬时积压突发型。比如上游做了批量数据修复、促销活动整点放量、或者是消费端依赖的下游数据库发生了30秒抖动。这类堆积的特征是Lag虽然在一瞬间冲得很高但消费者实例的TPS没有明显下滑甚至还在上升Lag曲线开始冲高后会在短时间内自然回落。这种情况的处理核心是“等”最多做一点临时扩容加速消化不需要改任何消费逻辑。第二类是持续堆积能力型。特征是Lag只增不减或者维持在一个高位平台不降同时消费者TPS远低于生产TPS甚至趋近于0。这说明消费端处理能力真的跟不上了要么是消费代码存在慢调用要么是线程被阻塞要么是Broker端读得太慢。这类堆积靠“等”只会越积越多必须介入处理。怎么区分这两种形态很简单看两个值lag曲线的斜率以及消费TPS是否跟着生产TPS波动。我一般直接到RocketMQ Dashboard上拉最近两小时的消费组Lag曲线和消费TPS曲线如果Lag冲高后斜率转负、TPS同步起来基本就是瞬时积压如果Lag是一条缓慢爬升的直线TPS趴在低位不动那就是持续堆积需要往下查。1.3 理解位点提交和ConsumeQueue排查才不会跑偏如果你想深入排查而不是停留在“加机器”层面还得理解消息从Broker到消费者的路径。RocketMQ的消息体写在CommitLog里这是顺序写所以写入性能极高。消费者的读取不是直接去CommitLog里按offset找消息而是先读ConsumeQueue——一个轻量级的索引文件里面存放的是消息在CommitLog里的物理偏移量。然后消费者再根据这个偏移量去CommitLog读取真实数据。问题就藏在“读取”这一步。如果消费速度追得上写入速度消费者读的数据基本都还在操作系统的PageCache里走的是内存速度快到飞起。一旦堆积量很大消费位点远远落后于最新写入位点消费者去读的数据大概率已经不在PageCache里了Broker就得走真实的磁盘IO去翻CommitLog文件。CommitLog是按顺序写的但消费端补数据时读的是很旧的物理位置本质上变成了一种随机读。堆积越深Broker的读放大越严重消费速度进一步下降形成恶性循环。这也是为什么很多堆积事故里消费者日志里全是拉取超时、Broker端磁盘IO被打满的原因。实际排查时一旦看到Broker的读IO飙升、等待时间拉长就得警惕是不是堆积深度已经把PageCache击穿了。2. 判断堆积的四类信号和一套可复制的排查链路2.1 四类信号从监控指标到业务感知处理堆积的第一步永远是“发现得够快”。我不建议等业务反馈“消息怎么还没到”再开始排查那样已经晚了。比较靠谱的监控体系至少覆盖以下四类信号消费组Lag监控这是最直接的指标。为每个核心消费组配置Lag的阈值告警比如超过10万或者持续5分钟增长就告警。更精细的做法是对核心Topic设置消费延迟时间告警消息从Broker写入到被成功消费的时间差这个比Lag更贴近业务感受。消费TPS监控消费端TPS掉到接近0或者长期低于生产端TPS说明消费链路有瓶颈。TPS曲线和大盘对比能快速区分瞬时积压和持续堆积。Broker层监控重点看磁盘读IO、IO util、PageCache剩余量、还有拉取请求的RT。堆积深度大的时候这些指标通常会出现异常抬升。消费端线程状态消费线程是否大面积BLOCKED、WAITING日志里是否频繁出现重平衡Rebalance、拉取消息超时等。2.2 一条可复制的排查链路真正遇到堆积我的排查顺序基本是固定的很少跳步。这套顺序你背下来也行遇到问题照着走至少不会像无头苍蝇一样乱撞。第一步确认堆积范围和程度拿到告警后先打开RocketMQ Dashboard或控制台找到对应消费组看Lag总量和每个队列的分布。同时用命令行确认一下Broker写入位点的情况# 查看某个消费组的消费进度和堆积情况 mqadmin consumerProgress -g consumerGroup # 查看Topic在各Broker上的写入位点 mqadmin topicStatus -n namesrvAddr -t topicconsumerProgress命令输出里有几个关键字段BROKER哪个Broker、QUEUE_ID队列编号、OFFSET当前消费位点、LAST_OFFSET最新写入位点、LAG堆积量。先把所有Broker上的LAG加起来判断总量再逐个队列看判断是不是某个队列特别多。第二步看消费TPS和消费端日志如果总量很大立刻去看消费端应用的日志和监控。重点看消费线程的TPS、异常堆栈、RPC调用的超时情况。这里有一个小技巧直接用jstack打印消费者进程的线程栈看消费线程到底卡在哪个方法上。# 打印消费者进程的线程栈查找消费线程状态 jstack consumer_pid thread_dump.txt # 在dump文件里搜索消费线程观察线程状态是RUNNABLE还是BLOCKED/WAITING grep -A 30 ConsumeMessageThread thread_dump.txt如果线程大量阻塞在某个RPC调用、数据库查询或者加锁操作上根因基本就在消费逻辑里如果线程状态正常但TPS就是上不去再看Broker端和网络。第三步查Broker状态到Broker机器上用top和iostat看磁盘IO情况确认PageCache是否已经耗尽。如果磁盘IO util长期接近100%基本实锤了“堆积导致读放大、读放大又拖慢消费”的恶性循环。# 查看磁盘IO等待和利用率 iostat -x 1 # 查看系统负载和内存缓存占用 top第四步确定是“哪个队列”和“哪个消费实例”的问题RocketMQ的消息是分队列Queue存储的一个Topic默认有读写队列数量一般是4或者8。消费组内的消费者实例会通过负载均衡分配队列。这里有个关键知识点同一个消费组内一个队列在同一时刻只能被一个消费者实例消费。所以如果你有16个队列但开启了20个消费者实例多出来的4个实例是闲置的反过来如果Topic只有2个队列你开10个消费者也没用并行度天花板就是2。排查时要看每个消费者实例的消费TPS分布。如果某个实例TPS特别高、其他实例基本空闲那是队列分配不均或者某个队列出现了热点如果所有实例TPS都上不去那是消费逻辑或者Broker读的共性问题。下面把排查步骤整理成一个表格方便你对照使用排查阶段使用手段重点观察的指标判断结论确认堆积总量Dashboard / mqadmin consumerProgressLAG、LAST_OFFSET总量多大、哪个队列堆积判断堆积形态Lag曲线 消费TPS曲线Lag斜率、TPS趋势瞬时积压还是持续堆积定位消费端问题日志、jstack、APM线程状态、异常堆栈、RPC耗时消费逻辑是否阻塞定位Broker瓶颈iostat、top、Broker监控磁盘IO util、读IO、PageCacheBroker是否已读不动确认队列分布consumerProgress逐队列查看各队列LAG、各实例TPS是否存在队列热点或分配不均3. 三种典型堆积现场消费慢、拉取慢、写不进3.1 消费端处理慢最常见的堆积现场这类堆积占了大部分生产事故。表象是消费者TPS很低但消费者的CPU可能不高因为线程都阻塞在等待外部资源上。我在线上见过最典型的两类**一是消费逻辑里做了重量级操作。**比如每条消息都触发一次远程RPC调用或一条慢SQL。系统平时的TPS撑得住但上游一旦放量下游服务先扛不住响应时间变长消费线程大量阻塞等待消费TPS立刻崩掉堆积随之而来。**二是幂等和重试机制引起的“重试风暴”。**有些消费代码没有做好幂等处理失败后消息进入重试队列重试时又把下游打挂了失败更多重试更多形成一个恶性循环。这种场景下消费日志里全是超时和重试的报错堆积持续增长重试队列也越堆越高。怎么判断是消费端的问题最简单的基准法用mqadmin consumerProgress看消费TPS如果TPS明显低于正常水位并且用jstack能看到消费线程阻塞在某个外部调用上基本可以确诊。3.2 拉取和分配不均队列决定并行度天花板这个场景很多人会忽略但它特别容易在扩容后发生。前面说过一个队列只能被一个消费实例消费所以当你面对堆积时如果方案的第一个动作是“无脑加机器”往往会踩到这个坑。我见过一个真实案例Topic只配了8个队列消费端部署了10台实例。业务同学说加了机器后堆积还是下不去一查发现有两台实例一直拿不到队列分配处于空转状态。真正在干活的只有8台。后来把Topic的读队列扩到32个RocketMQ支持在控制台或命令行扩容读写队列再让消费实例重平衡并行度才真正提上去。另外还有一种“分配不均”的情况。为了支持顺序消息生产端往往用固定的消息Key哈希取模选择队列这可能导致某个Broker上的某个队列消息特别多其他队列很闲。消费端按队列消费时热点队列就成了瓶颈整体TPS被个别的长尾队列拖住。3.3 Broker端读不动PageCache击穿与磁盘随机读这类堆积和消费端代码没关系问题出在Broker自己身上。当消费Lag很深的时刻消费者要读的数据已经不落在PageCache里了。从前面的读写路径可以看到CommitLog顺序写很快但消费端消费旧数据时相当于在一个很大的顺序写文件里去random read。磁盘随机读的性能远低于顺序读IO util会被打满消费拉取请求的RT飙升消费者端表现为持续超时、TPS上不去。检查Broker有没有被读拖垮iostat的输出最直观# 注意 %util 和 rkB/s读的值 iostat -x 1如果%util长期接近100%并且rkB/s异常高那Broker的磁盘读取就是瓶颈。此时的处理方式和消费端慢就不一样了单纯加消费者实例效果也有限——因为瓶颈在下游的磁盘而不是消费者的CPU。这种情况下更要先把堆积深度打下来把消费位点快速推进到较新的位置让消费者读的数据回到PageCache里Broker的IO自然就降下来了。3.4 三种现场的快速对照表堆积现场典型表现根因位置最有效的第一步消费端处理慢消费TPS低、线程阻塞在外部调用消费逻辑中有慢操作或重试风暴摘除异常依赖或暂停重试风暴先止损队列分配不均实例空转、部分队列LAG特别大Topic队列数少、顺序消息哈希不均扩容读写队列并触发重平衡Broker读瓶颈磁盘IO util高、拉取超时PageCache被击穿、随机读放大先把Lag打下来让读回到PageCache4. 从临时止损到长期根治分层次处理堆积4.1 第一步永远是止损而不是优化我见过有人一上来就改消费代码、加索引、优化SQL结果改了半小时堆积还在涨业务损失持续扩大。正确的顺序是先“止血”如果判断是消费逻辑的下游依赖出了问题先把消费端的那个外部调用降级或摘除让消息先被“快速消费掉”保住消费TPS业务延迟问题后面再通过补数据或重放来解决。现实中很多消费场景其实是异步通知类的——给用户发个推送、更新一下缓存。这类业务对消息内容不是强依赖即使消费时临时跳过某个下游调用消息也能正常消费堆积就会被迅速消化。但要注意跳过调用意味着这段时间的业务动作缺失事后需要做好补偿。对于不可丢失类型的消息绝对不能简单跳过只能在代码层面快速fail-safe比如把失败的写入本地表后续再异步补发。4.2 临时扩容加线程数 vs 加机器要分清止损之后如果堆积形态属于瞬时积压、确认消费逻辑本身没大问题可以考虑加速消费。先说加机器。加机器能提升并行度但前提是队列数足够。这就是前面反复强调的并行度上限 Topic的队列数 × 单实例消费线程数。如果你加的实例数已经超过队列数直接白加。所以加机器之前先看一眼消费组当前实例数和队列数。再说加线程数。RocketMQ的PushConsumer支持动态调整消费线程数配置项是consumeThreadMin和consumeThreadMax。默认线程数一般不高可以适当上调defaultMQPushConsumer.setConsumeThreadMin(20); defaultMQPushConsumer.setConsumeThreadMax(64);前提是你的消费逻辑是纯CPU或IO可控的盲目调大线程数在一些场景里反而会把下游数据库或下游服务压垮。调线程数的时候务必盯着下游系统的负载随时准备回退。还有一个常用的加速手段是批量消费。consumeMessageBatchMaxSize这个参数控制每次消费时拉取的最大消息条数配合理好消费端的批量处理逻辑比如批量插入、批量更新吞吐量往往能提升数倍。需要注意不是把这个参数调大就一定快如果消费代码本身是一条一条处理拉再多条也只是单条循环。4.3 定位根因后的长期优化堆积被消化完、业务恢复后才轮到真正的根治。根据不同的根因方向也不同消费逻辑慢把消费逻辑里的同步RPC改成异步化批量处理代替单条处理去掉不必要的加锁对下游调用加入超时和熔断机制。这是性价比最高的一类优化。数据库瓶颈消费端做批量写、合并请求或者引入本地缓存/结果表减少对热数据的重复查询。上游突增流量不一定要让消费端无限扛可以反过来对生产端做限流。在RocketMQ场景下给生产端加上流量控制削峰填谷不让消息洪峰一次性压垮下游。顺序消息瓶颈顺序消息天然会把某个队列的消费串行化这是业务需求决定的硬约束。优化空间主要在于减少单条消息的处理耗时别在顺序消费的代码里做重量级操作把不需要顺序保证的部分拆出去并行处理。4.4 两个危险操作重置消费位点和丢消息这里要重点提醒两个新手容易踩的坑。**第一个坑是轻易重置消费位点。**有些同学看到堆积量很大会想到用mqadmin或控制台把消费位点往前跳让消费者直接去消费最新消息把旧的堆积“跳过去”。这个操作在消息可丢失或可补偿的场景是有效的止损手段但如果在订单、支付这类要求精确不丢消息的业务里跳位点等于直接丢消息后续对账和补偿的成本会非常高。真要跳位点必须和业务确认可以接受消息丢弃并且做好补偿方案。**第二个坑是忽略死信队列。**RocketMQ中消息消费失败重试16次后会被投递到死信队列DLQ。堆积事故中如果消费端大量处理失败死信队列里的消息也会暴涨。排查时要同时看一眼DLQ的积压数量别只盯着主队列。死信队列里的消息需要单独建一个处理程序定期扫描、分析失败原因、修复后人工重放。我曾经处理过一个凌晨的告警主队列堆积量其实不大但DLQ里躺了四十万条消息就是因为消费代码里的一个空指针异常没被发现。5. 面试官到底在考什么答题框架与话术示例5.1 这个问题的核心考点面试官抛出“RocketMQ消息堆积了怎么处理”表面上是考你“会不会处理堆积”实际上考的是三件事第一你理不理解RocketMQ的基本运行机制。如果你能说出Queue是并行度上限、PageCache与消息堆积深度之间的关系、消费位点的作用说明你对这个中间件的理解不是停留在“能跑通Demo”的层面。第二你有没有排查思路。遇到问题第一步干什么、第二步干什么是背出来的方案还是有层次感的排查链路。一个有经验的人一定会先“定性”再“止损”再“根治”。第三你有没有生产经验的边界感。比如知道加机器不是万能的、知道跳位点有风险、知道死信队列要处理——这些是看文档看不出来的只有踩过坑才知道。5.2 一个可复用的三分钟回答框架如果在面试中被问到我会这样组织答案你可以参考这个框架替换成自己实际的经历首先我会确认堆积属于什么形态。看消费组的Lag曲线如果是瞬时冲高后会自然回落说明是突发流量或下游抖动导致的瞬时积压主要做临时扩容加速消费如果Lag持续增长、消费TPS跟不上生产TPS那就是消费能力不足需要进一步排查。接下来按链路排查先看消费端日志和线程状态用jstack看消费线程卡在哪里再看Broker的磁盘IO和PageCache确认是不是堆积太深导致读取落盘、IO被打满同时检查队列分布和消费者实例数确认是不是队列数限制了并行度。定位到根因后处理如果下游依赖有问题先降级止损如果消费逻辑慢做批量化和异步化改造如果是队列并行度不够考虑扩队列并加消费者实例。整个过程会注意不能盲目跳消费位点涉及核心业务宁可快速消费落本地表、再做补偿也不直接丢消息。事后会补监控给消费组加Lag告警、消费延迟告警分析这次堆积的触发因素避免下次再犯。这个回答里有几个亮点是面试官很吃的先说定性、再说止损、再说根因、最后说预防而且全程有具体的命令和参数听起来就像真正处理过线上问题的人。5.3 三个减分回答和三个加分细节减分的回答见过太多了总结一下只回答“加机器”说不出为什么有时候加机器没用。只回答“重启消费者”重启只能解决一时的线程卡死如果消费代码有慢SQL重启完很快又堆积。从头到尾不提监控和预防好像问题解决就结束了。加分的细节主动提到“消费者实例数不能超过队列数”这体现出你对负载均衡和数据分布的理解。主动提到“PageCache与堆积深度”说明你理解Broker底层的读写模型。主动提到“重置消费位点是最后的止损手段而不是常规手段”说明你有生产风险评估的思维。6. 一次线上堆积复盘的完整记录6.1 事故表象晚间高峰期消费组Lag暴涨去年负责的一个业务每天晚间有一个定时任务会把一批业务数据推送给用户。那晚监控突然告警核心消费组order-notify-group的Lag在一小时内从2万涨到120万用户侧开始出现“消息迟迟收不到”的投诉。我先做了第一步判断拉出Lag曲线和消费TPS曲线。Lag曲线一路陡增没有任何回落迹象消费TPS从正常的8000条/秒暴跌到500条/秒左右。这不是瞬时积压是持续堆积问题出在消费链路本身。6.2 定位过程线程栈里找出的真凶我按照排查链路一步步往下走。先看消费端应用的日志发现日志里大量出现调用下游接口超时的报错。再用jstack抓线程栈看到几十个消费线程全部阻塞在同一个下游HTTP调用的InputStream.read上。那个接口是我们自己维护的“用户偏好服务”当晚正好在做版本发布响应时间从20ms飙到了800ms并且批量超时把Tomcat连接池打满了其他请求也进不去。与此同时我检查了Broker状态磁盘IO尚可PageCache没有被打穿。问题基本锁定在消费端依赖的下游服务超时导致消费线程大面积阻塞消费TPS骤降Lag持续攀升。6.3 处理动作先止损再根治确认根因后的第一件事不是改代码而是止损。方案是快速把消费端对该下游接口的调用改为降级失败时直接返回默认值不让线程阻塞等待。这个改动上线后消费线程立刻恢复消费TPS在20分钟内从500条/秒恢复到8000条/秒以上Lag曲线开始快速下探大约一个多小时后堆积全部消化完。这里有个处理技巧因为消费场景是“发送订单状态变更通知”消息里本身带着订单状态字段即使下游偏好服务暂时降级通知也能正常发出只是个性化程度降低属于业务可接受的补偿范围。事后我们再对降级期间的个性化缺失做了部分补偿推送。根治阶段做三件事一是给下游调用加了超时熔断连接超时200ms、读超时300ms、错误率超过50%熔断10秒二是把单个消息串行调用下游改成批量接口一次传递多条消息三是消费线程池增加了动态调节。同时给这个消费组加了Lag告警和消费延迟告警阈值设为“持续5分钟Lag超过5万”才触发避免噪音。6.4 复盘里最重要的几个认知这次事故给我留下几个很深的教训**堆积的峰值数字不是最可怕的可怕的是消费TPS没有恢复。**只要消费TPS能重新拉起来哪怕Lag再大也只是一个时间问题。我当时处理的核心精力都放在“让消费线程别再阻塞”上而不是纠结怎么删消息。**止损方案必须提前设计。**如果消费逻辑里每个下游调用都有降级开关或fail-safe策略遇到事故时你只需要一个配置项就能止血没准备好的话现场改代码、发版本堆积早就把业务拖垮了。**监控告警的阈值要合理。**把告警设成“一有Lag就报警”团队很快会麻设成“持续增长且超过业务容忍阈值”才报警才能在真出事时有人响应。6.5 如果同样的问题再出现一次我会怎么做现在处理堆积问题的思路已经固化成一套固定动作分享出来供参考第一先看曲线定形态绝不盲目加机器。第二用jstack和日志快速定位是消费代码阻塞、队列不均还是Broker读瓶颈。第三先止血、再根治核心业务的堆积处理宁可快、宁可临时降级也不允许消息无限积压。第四处理完必须补监控和复盘把根因转成后续的防御手段。这套动作让我在处理类似问题时基本没有慌乱过。遇到堆积先别着急动手把“它是怎么来的”想清楚答案往往自己会浮现出来。
RELATED READING

延伸阅读

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