MQ消息积压排查实战:从定位消费卡点到应急止损与优化
发布时间:2026/9/16 7:24:06 作者:尧图编辑部 阅读量:1,286

1. 一次真实的积压事故从告警刷屏到业务受损凌晨两点值班手机开始连续震动。打开告警平台一看MQ 的消费延迟已经从几百条涨到了几十万条消费端最慢的消费者 lag 超过 3 个小时。这意味着用户下单后一直收不到状态变更通知积分发放、短信通知、对账文件生成这些下游业务全部停摆。更麻烦的是生产者那边并没有异常发送速率一直平稳问题明显出在消费端但消费端的 CPU 和内存看起来都不高日志里也没有明显的报错——这是消息积压排查里最让人头疼的一种情况系统没有报错只是处理不过来了。那段时间我前后处理过好几次类似的积压事故从 RocketMQ 到 Kafka 都遇到过。这类问题的排查思路其实高度相似先搞清消费卡在哪一环再定位是能力不足还是逻辑被阻塞最后才是调参数、改代码、上应急手段。这篇文章就把我实际排查积压问题的完整链路、用过的命令、踩过的坑以及事后怎么防止再次积压一次性讲清楚。适合正在维护 MQ 生产环境的开发、运维同学尤其是那种消费者偶尔卡一下、积压几分钟又恢复但高峰期积压持续上涨的中间态场景。2. 消费卡顿定位三大突破口日志、线程状态与外部依赖2.1 第一步永远是看消费日志但要看处理耗时而不是有没有报错很多人排查积压的第一反应是翻报错日志这其实是个误区。消息积压的场景里消费端往往没有 Exception只是每条消息的处理时间从正常的 50ms 慢慢变成了 2s、5s甚至更久。表面上看系统没问题实际吞吐已经崩了。我当时做的第一件事是把消费端每个消息从拉取到处理完成的时间打点拉出来按耗时分布统计。正常的处理耗时曲线应该是集中在某个低值区间的比如 P99 在 200ms 以内。当积压发生时P99 大概率已经飙升到秒级。这个数据能帮你快速判断到底是个别消息特别慢还是所有消息都变慢了。如果是个别消息慢重点看那几条消息携带的业务参数——比如某个订单号特别大的订单、某个还没上线的新渠道来源顺着业务特征去查对应数据如果是所有消息普遍变慢那问题通常在消费端的公共依赖上比如数据库连接池、Redis、下游 RPC 接口。这一步看起来简单但很多人跳过耗时分析直接去查代码容易白忙一场。2.2 线程堆栈用 jstack 看消费者线程到底卡在哪儿日志只能告诉你慢告诉不了你卡在哪个方法里。这时候就需要抓线程堆栈。以 Java 技术栈为例消费者一般跑在线程池里。找出消费线程池的线程 ID然后多次执行 jstack 抓取堆栈间隔几秒抓一次连续抓五六次。重点看线程状态RUNNABLE 状态且栈顶是 Socket 读写或 JDBC 调用说明卡在等待外部响应BLOCKED 状态说明线程在等锁可能和某个资源竞争有关WAITING / TIMED_WAITING 状态看它在等哪个条件变量常见的是等线程池队列、等连接池归还连接。我遇到过一个典型案例消费者线程大量阻塞在数据库连接池的等待队列上。业务方查了慢 SQL发现有一条 UPDATE 语句在高峰期执行时间涨到 4 秒连接被占满后所有消费线程都拿不到连接整条消费链路就僵住了。jstack 的堆栈里清楚显示着HikariCP.getConnection阻塞顺着这条线索很快定位到慢 SQL。这里有个技巧不要只看一次堆栈因为线程状态是瞬时的。连续抓几次看状态是否一直在同一个地方。如果每次都卡在同一个方法基本可以断定这就是瓶颈点。2.3 外部依赖排查慢 SQL、Redis 超时、下游接口响应线程堆栈定位到卡在外部依赖之后下一步就是把所有外部依赖逐一排查。我的习惯是按耗时占比从高到低排数据库看慢查询日志、看连接池活跃连接数、看锁等待情况。尤其是那种高峰期突然出现的锁竞争比如大量消息在更新同一张表的同一行行锁会直接把消费线程串行化。Redis看慢日志、看网络耗时。有些操作比如KEYS或大 Key 的SMEMBERS在数据量大时非常致命。下游 RPC看调用方超时时间设置。很多消费逻辑里同步调用了外部服务下游服务一旦抖动消费速度立刻被打崩。一个很容易被忽略的细节是消费线程池的线程数是有限的假设核心线程数 20每个线程阻塞等下游 3 秒整体吞吐就是 20/3 ≈ 6.7 TPS。哪怕消息本身处理逻辑再简单吞吐也被死死限制住。所以外部依赖的响应时间对消费吞吐的影响是乘数级的这一点在估算消费能力时一定要算进去。2.4 从控制台/管理页面确认积压现状再说说热词里提到的MQ 怎么在页面查看消息。不同 MQ 的管理端虽然不一样但核心信息都是类似的积压量和消费位点。以常见的管理页面为例RocketMQ Dashboard在消费详情页面可以看到 Consumer Group 的实时积压数量按 Topic 划分还能看到消费位点和最新位点的差值。Kafka 的 Kafka Tool 或 Confluent Control Center直接看 Consumer Group 的 Current Offset 和 Log End Offset两者差值就是 lag。云厂商的 MQ 控制台一般都有消息堆积监控图能看到每个消费组的历史积压曲线。有两点经验供参考。第一看积压要看曲线不要只看当前值。如果积压量在持续上涨说明消费速度已经彻底跟不上生产速度如果积压量在缓慢下降说明消费能力还是够的只是暂时积压了一批可以不用太慌。第二管理页面能看到的积压量是结果是 lag 的汇总值要找到哪台机器上的哪个消费者实例处理最慢还得看消费端监控。把 Dashboard 和消费端日志、监控配合起来才能定位到具体实例。3. 堆积根因的常见分类并发不足、重试风暴与 Rebalance3.1 并发消费配置不合理你以为的并发不是真正的并发很多人有个误解只要消费线程数配得多吞吐就一定高。这句话只对了一半。消费线程数确实影响吞吐但前提是每个线程都在高效工作而不是在排队等锁、等 IO。以 Kafka 为例消费并发其实受分区数限制。一个 Consumer Group 里同一个分区的消息只会被组内的一个消费者实例消费。如果你的 Topic 只有 5 个分区那你起 20 个消费者实例也只有 5 个在干活。我见过不少团队在 Kafka 上盲目加消费者实例结果积压一点没缓解反而因为 Rebalance 频繁把消费搞得更乱。RocketMQ 的逻辑不太一样它的消费队列MessageQueue是按 Consumer 实例数量平均分配的同一个消费者实例内部还有线程池并发消费。但线程池默认的消费线程数如果设置过小比如只有 10而单条消息处理耗时又长吞吐自然上不去。反过来线程数设置过大也可能出问题大量线程同时去查数据库、调下游把下游打挂引发更严重的连锁故障。3.2 单条消息抛异常导致的重试风暴积压问题里最阴险的一种情况消费逻辑里没有对单条消息做异常隔离一条坏消息触发了无限重试。我之前排查过一个线上事故某条消息的业务数据里有个字段是 null消费端在转换时抛了 NullPointerException。消费框架将这条消息放入重试队列重试了十几次还是失败每次重试都占用一个消费线程好几百毫秒。基础消息量本身不大但这条坏消息就像一颗老鼠屎把整个队列的处理节奏拖慢最终导致积压量上涨。这种问题的可怕之处在于重试消息夹杂在正常消息里反复消费消费者日志里会周期性出现同一条报错但因为不是持续报错很多人会忽略。排查时注意看日志里是否有规律性的重复异常一旦发现要立刻把这条消息单独拎出来看看是数据问题还是代码 bug。处理方案有两个层面代码上给消费逻辑加上异常分类业务异常和系统异常分开处理重试次数超过阈值就直接进入死信队列运维上遇到已经积压的死信可以临时隔离或者跳过先把正常消息消费完。后面第 5 部分会讲具体的应急操作。3.3 客户端 GC 停顿与 CPU 抖动导致消费暂停还有一种不太容易发现的积压原因消费端进程本身没有问题但 JVM 在频繁 GC尤其是 Full GC。我处理过一起事故消费端老年代不断增长Full GC 每次持续好几秒。GC 期间整个进程的线程都暂停了消费者当然也在暂停积压量就一波一波地涨GC 前消费跟上GC 积压一些再跟上再积压。这种积压的特征是曲线呈锯齿状周期性上涨又下降但整体趋势是缓慢上升的。排查思路是看 GC 日志和内存快照。GC 日志里如果 Full GC 频繁触发大概率是内存里存在大对象或对象无法及时回收。比如消费逻辑里把大批量的消息数据加载到内存做聚合或者第三方 SDK 持有大量缓存。解决方向是优化内存占用、调整堆大小和 GC 参数而不是调消费线程数。这里要提醒一句别一看到积压就慌着加资源先看 GC 和 CPU很多时候根因在 JVM 层面。3.4 Rebalance 引发的消费暂停与分区分配不均Kafka 场景里Rebalance 是积压的另一个重要诱因。消费者实例增加、减少、或者心跳超时都会触发 Rebalance。Rebalance 期间整个 Consumer Group 会短暂停止消费如果 Rebalance 频繁发生比如每几分钟一次那消费有效工作时间会大打折扣。一种常见情况是消费者实例处理消息太慢超过了max.poll.interval.ms的默认设置5 分钟broker 认为该消费者已经死亡踢出分组并触发 Rebalance。但消费者实际还活着只是卡在一条消息上结果就被踢了。重新加入分组、重新分配分区、再消费、再卡住形成一个恶性循环。处理方式调大max.poll.interval.ms同时配合调大max.poll.records因为调大间隔意味着单次拉取的消息要相应增加否则消费能力反而下降并在代码里保证 poll 循环的单次耗时可控。另外检查会话超时时间和心跳线程设置避免因为 GC 停顿导致心跳发送延迟、被误判为死节点。4. 消费速度优化的实践清单参数调优与代码改造4.1 消费参数调优的具体方向与合理范围积压发生时先做短平快的参数调整往往比改代码更快见效。但参数不是拍脑袋调要有依据。各 MQ 的常用参数对照如下基于常见版本默认值具体以你的版本为准参数项作用经验值/调整思路消费线程数RocketMQ consumeThreadMin/Max决定消费者实例内并发处理消息的线程数从默认值逐步调大观察 CPU 和下游负载一般从 20 调至 40~60 需谨慎压测每次拉取条数Kafka max.poll.records单次 poll 返回的最大消息数默认 500如果单条处理耗时长建议调小到 100~200避免单次 poll 处理超时消费超时时间Kafka max.poll.interval.ms两次 poll 的最大间隔默认 3000005 分钟如果消息处理耗时可能超过 5 分钟需要调大并配合调大 max.poll.records拉取字节数fetch.max.bytes单次拉取的数据量上限消息体较大时调大减少网络往返次数消费端心跳间隔heartbeat.interval.ms心跳发送频率确保小于 session.timeout.ms 的 1/3避免被误判为宕机参数调整的核心逻辑是让并发数与单条耗时匹配如果单条消息处理 100ms20 个线程的理论吞吐是 200 TPS如果单条消息处理 1s20 个线程的理论吞吐只有 20 TPS。你要先算出现状再决定是调并发还是先优化单条耗时。有一个很容易踩的坑只调大消费线程数不关注下游依赖的承受能力。线程数翻倍意味着下游数据库和 RPC 的 QPS 也翻倍如果下游本来就在高负载运行可能直接把下游打挂。所以调参之后要盯下游监控随时准备回退。4.2 代码层面的消费耗时改造串行改并行、同步改异步参数调优是治标代码优化才是治本。消费逻辑的耗时改造我通常会按优先级做三件事。第一件事把消费逻辑里的串行外部调用改成并行。举个例子一条消息可能需要同时调用用户服务查用户信息、调用订单服务查订单详情、调用库存服务查库存。这三个调用之间没有依赖关系如果按顺序执行总耗时是三次调用之和改成并行调用总耗时接近最慢的那一次。用 CompletableFuture 或者简单的线程池就能实现改造量不大吞吐提升却可能是成倍的。第二件事非核心逻辑异步化。比如消息处理完成后要发一条审计日志、推送一次数据同步这类操作不影响主流程完全可以从消费线程里剥离出去丢进独立的异步线程池。这样消费线程只做最核心的业务逻辑处理完立刻返回吞吐立刻提升。第三件事减少重复计算。比较典型的是消费逻辑里频繁查询配置、频繁查询基础数据。这些数据如果变化频率低完全可以本地缓存。我之前优化过一个消费链路代码里每次消费都查一次数据库配置表改成本地缓存、每 5 分钟刷新一次之后单条消息耗时从 45ms 降到 12ms积压问题直接消失。4.3 消息体瘦身积压时最容易被忽略的加速手段消息体积对消费速度的影响往往被低估。消息体越大网络传输耗时越长序列化/反序列化的 CPU 开销也越大。在高峰期消息体从 2KB 涨到 200KB消费吞吐可能下降 30% 以上。有一类积压问题的根源特别隐蔽生产端往消息里塞了太多不必要的数据。比如订单变更消息本该只传订单 ID 和变更类型结果把整个订单对象、关联的商品信息、用户信息全塞进去了。消费端明明只需要 3 个字段却要把 200KB 的 JSON 反序列化一遍。遇到这种情况我通常建议生产端做消息瘦身只传关键业务 ID消费端按需查询数据。虽然会增加一次查询开销但网络和序列化的开销下降幅度更明显。如果生产端一时改不了消费端也可以做优化不要每次都对完整消息体做全量反序列化很多框架支持按需解析需要的字段也能减少部分开销。4.4 积压场景下的临时加速方案扩容、跳过、批量补偿当积压已经发生且业务受损时等代码优化上线来不及得先上应急手段。我按操作风险从低到高排序给出三个常用方案方案一临时扩容消费者实例。这对 Kafka 这类支持横向扩展的 MQ 特别有效——但要先确认分区数够不够。如果分区只有 5 个加再多消费者也只有 5 个在消费需要同时把 Topic 的分区数扩大。RocketMQ 同样可以通过加消费者实例来提升消费并发度。方案二临时调整消费参数。加大消费线程数、调大每次拉取的消息数量、缩短单条消息处理中的等待时间。这些操作无需发版通过配置中心或 JMX 就能做适合快速止血。方案三跳过/丢弃部分积压消息。这个操作要非常谨慎最好在业务方确认消息可丢的情况下再做。比如某些实时性要求不高的统计类消息积压 3 小时后就算消费了也没有意义那就可以直接跳过。具体做法是把积压 Topic 的消费者指向一个空逻辑或者把指定消息标记为已消费。要注意跳过消息前一定先做备份万一后续需要补数还有数据可用。5. 积压事故的应急止损与常态化防御5.1 止损优先级先恢复业务再查根因处理积压事故最忌一条路走到黑非要先找出根因才动手。我的原则是先恢复业务再定位根因。哪怕只是临时让积压量停止上涨也为后续排查争取了时间。止损的操作思路按顺序来第一步如果是消费逻辑抛异常导致的重试风暴先把异常消息隔离停止无谓的重试第二步适当调大消费并发和拉取量看看消费速度是否立竿见影地提升第三步如果积压量还在涨果断加消费者实例或扩容分区第四步同步排查根因但不能因为排查而耽误前面的止损动作。这个顺序背后的道理是积压每多一分钟业务损失就多一分下游依赖的超时、超库存、超时未发货等问题都会连锁爆发。先把水位降下来再去做第 2 章、第 3 章讲的那些定位分析。5.2 重置消费位点的风险与正确姿势有些场景下积压消息已经失去时效性或者消费端代码已经修复、旧消息再消费也没意义这时候可以考虑重置消费位点Offset让消费者从最新位置开始消费跳过积压的历史消息。但重置位点是高危操作务必注意三点确认业务可接受消息丢失。重置位点意味着把 lag 直接清零积压的那部分消息不会再被当前消费者处理。在低峰期操作。重置位点期间消费者可能需要 Rebalance如果业务正在依赖消费结果短时间的消费暂停会造成新的问题。先记录原始位点。操作前把消费组的当前 offset、topic、partition 信息完整记录万一后续需要回滚或补数还能恢复。Kafka 里用kafka-consumer-groups.sh --reset-offsets可以精确控制把消费位点重置到指定时间或指定偏移量。RocketMQ 的控制台也支持按时间重置位点。我个人的经验是能按时间重置就按时间重置比如把位点重置到积压开始前的时刻这样既跳过了坏消息又不会丢掉积压开始前未消费完的数据。5.3 监控告警体系把积压消灭在爆发之前积压事故处理得多了你会发现大部分积压本来是可以提前发现的。问题在于很多团队的监控告警设得太粗比如只在积压量超过 10 万条时才告警但等积压到 10 万条业务早就受到影响多时了。合理的告警策略是把积压量和消费能力水位结合起来积压绝对值针对每个 Topic 和消费组设置积压阈值。核心业务可以设 1 万条非核心业务可以设 10 万条分级处理。积压变化率告警不能只看当前值要看趋势。积压量 10 分钟内翻倍比积压量 1000 条但持续缓慢上涨的信号价值更高。消费延迟时间以最早一条未消费消息的停留时间作为指标比积压条数更直观。比如延迟超过 5 分钟告警超过 30 分钟电话通知。有了监控还要有值班响应预案。我的习惯是准备一份积压应急操作手册里面写明如何判断积压类型、第一步止损操作是什么、需要联系哪些下游业务方、哪些 Topic 可以跳过消息。手册越详细线上出问题时越不慌。5.4 压测与容量评估提前知道你的消费极限在哪里最后聊一个很多团队容易忽略的环节消费端压测。大多数积压问题在平时不会暴露因为日常消息量远低于消费能力上限只有到业务峰值比如大促、活动秒杀时消息量突然翻好几倍消费端瓶颈才彻底暴露。建议在业务上线前或大促前做一次消费链路压测。压测目标很简单搞清楚在当前代码、当前参数、当前下游负载的情况下消费端最大吞吐是多少 TPS。方法是用压测工具向 Topic 里灌消息观察消费端的处理速度逐步加压到消费端出现 lag 持续增长的点那个点就是消费能力的上限。有了这个数据你就能回答一个很实际的问题如果大促消息量翻 5 倍现有消费端扛得住吗需要加多少消费者实例要不要提前优化代码。压测中发现的瓶颈往往和线上积压高度一致——可能是下游接口慢、可能是消费线程池配置不合理、可能是数据库锁竞争。提前解决掉这些比等线上积压了再熬夜排查要省心得多。从我这几年处理 MQ 积压问题的经验来看积压本身不可怕可怕的是积压发生后手足无措、乱调参数、盲目扩容。只要排查链路清晰——先看耗时分布再抓线程堆栈接着确认外部依赖然后分类根因、对症下药大部分积压问题都能在一个小时内控制住。而真正让团队省心的还是平时的参数规划、压测验证和监控告警。把这几件事做在前面凌晨两点的告警自然就少了。