RocketMQ的消息堆积问题,几乎是我面试中被问到概率最高的实战题。说实话,这个问题能看出一个人是只会用MQ还是真在生产环境扛过事。这次不聊那些"提高消费能力""扩容"之类的官话,我直接把这几年在处理堆积问题时踩过的坑、用过的招、背后的原理都摊开讲,从排查链路到应对方案,再到事后复盘,一套完整的实战思路给你。
1. 先把"堆积"这个事定义清楚:什么情况才算真的堆积
很多刚接触RocketMQ的同学有个误区:看到监控面板上Consumer Lag大于0就慌得不行,赶紧到处问怎么办。其实堆积本身不可怕,可怕的是堆积的速度和处理速度之间的关系。我一般这样判断堆积是不是需要介入:
堆积速率 = 生产端写入速率 - 消费端处理速率
如果消费速率只是短暂低于生产速率,比如峰值流量持续几分钟,积压几百条甚至几千条消息,但很快消费速率又追上来了,这种属于正常波动,不需要做任何处理。
真正需要警惕的是:生产速率长期高于消费速率,消费位点(Lag)持续单调递增,而且增长趋势没有放缓的迹象。这时候堆积就不是"一会儿就好"的问题,而是系统性的消费瓶颈。
判断方法也很简单,登录RocketMQ控制台或者用mqadmin命令行工具查消费进度:
mqadmin consumerProgress -g group_name -n 127.0.0.1:9876重点关注Diff列(剩余未消费消息数)和LastTimestamp列(最后一次消费时间)。如果Diff持续上涨,或者LastTimestamp和当前时间差距越来越大,那就不是虚惊,是真堆积了。
生产环境我一般设两个阈值:堆积量超过50万条,或者堆积时间超过10分钟,就触发告警。为什么是这两个值?50万条基本上说明消费端有持续性的性能问题,不是瞬时流量能解释的;10分钟以上则说明消费受阻趋势明显,再拖可能影响业务数据的及时性。当然这两个值得根据业务调整,有的业务要求秒级消费,有的允许分钟级延迟,不能一刀切。
还有个很多人忽略的信号:消息消费的最大重试次数触顶。RocketMQ默认每条消息最多重试16次,超过之后消息进入死信队列。如果你在监控里发现死信队列的Topic(%DLQ%消费组名)开始有消息进来,说明堆积背后可能不仅是慢,而是有消息压根消费不了,一直在原地重试,这比单纯慢更麻烦。
2. 堆积为什么可怕:表面是延迟,实际是连环雪崩
处理堆积之前,我建议你先想清楚一个事:堆积到底会造成什么后果?理解了后果,你才知道为什么有些解决方案是饮鸩止渴,有些是正道。
2.1 业务层面的"新鲜度"被破坏
消息队列的核心价值之一,是解耦和削峰填谷,但隐含的前提是消息能在合理时间内被处理。假设你的业务是订单超时自动关单,消费者每5秒扫一次延迟消息,正常情况下用户下单后30分钟未支付就会关单。结果消息堆积了1个小时,那关单动作就延后了1个小时——用户可能在第31分钟付了款,系统却依然把订单关了,因为系统以为用户30分钟没付。这种业务逻辑错误比消息丢失还隐蔽,你很难排查出来。
再比如你对接的下游接口有签名时效,消息里带的token有效期是15分钟,堆积半小时后消费者拿到的消息已经失效,怎么调下游都是失败,失败又触发重试,重试再堆积,整个链路彻底卡死。
2.2 Broker存储层面的连锁反应
消息堆积不只是消费侧的事,它是真实占用Broker磁盘的。RocketMQ默认单个CommitLog文件是1GB,消息堆积越久,CommitLog文件越多。我遇到过最夸张的一次,某业务消费组一周没消费,Topic的数据量从每天几GB涨到几十GB,直接把Broker磁盘干到90%以上,最后触发了Broker的自我保护机制,连其他正常消费的Topic都被拖慢了。
还有一个容易踩的坑:消息堆积导致磁盘写入变慢,进而引发PageCache压力剧增,整个Broker的读写性能都会下降,最终影响的不只是堆积的这个Topic,而是所有共用同一Broker的业务。这就是为什么我常说,堆积问题如果不能快速止血,大概率会演变成全集群事故。
2.3 消费端本身的资源耗尽
消费者线程池默认20个线程,如果每条消息处理耗时从50ms涨到500ms,消费能力直接掉一个数量级。更可怕的是,如果你的消费逻辑里有数据库连接池、HTTP连接池这类有限资源,消息一堆积,并发处理的请求变多,连接池被打满,后续消息都在等连接释放,消费速度进一步下降,形成正反馈死循环。
3. 正规排查链路:从"知道堆积"到"定位根因"的完整路径
不少人在堆积发生时第一反应是"赶紧扩机器",其实这是本末倒置。你不先搞清楚堆积在哪一环,扩再多机器可能都没用。我自己的排查顺序是这样的:
3.1 先看消费组整体的Lag分布
用控制台或者命令行,先看一眼堆积的Topic在哪个消费组下,Lag是所有队列均匀分布,还是集中在某几个队列。
这个信息非常关键。如果是均匀分布,说明消费者实例整体处理能力不足;如果集中分布,说明某几个队列的消息处理特别慢,或者绑定的消费者实例出问题了。
RocketMQ的消费模式是:一个消费组里多个消费者实例,每个实例负责一部分队列。如果某台消费机器CPU飙高、网络抖动、或者Full GC频繁,它负责的那几个队列消息就会越积越多。
3.2 再查消费线程状态
找到可疑的消费者实例后,直接去看JVM线程栈。我一般用jstack抓线程快照:
jstack <pid> > jstack_output.txt重点找ConsumeMessageThread开头的线程,看它们卡在哪个方法上。最常见的几种情况:
- 卡在
java.net.SocketInputStream.read上——说明在等下游接口响应,大概率是下游变慢了 - 卡在
java.sql.Connection相关的地方——说明在等数据库连接或者慢SQL执行 - 卡在
LockSupport.park——说明线程在等锁,可能有资源竞争
上次排查一个堆积问题,jstack抓出来所有消费线程全部卡在一个外部API的HttpClient调用上,那个API的TP99本来30ms,当时那个接口因为对方线上故障直接HTTP超时120秒不回,消费线程全堵在Socket读取上。不抓线程栈你根本想不到问题出在消费逻辑调用链的下游。有时候是下游接口性能突然劣化,比如从几百毫秒变成几秒,消费速率自然就掉下来了。
3.3 逐层拆解消费逻辑耗时
线程栈只能定性,要定量还得在代码里做链路追踪。我建议每个消费逻辑都加入埋点,统计这几段耗时:
- 从Broker拉取消息到拿到消息体的耗时(网络IO)
- 消息反序列化耗时
- 业务逻辑处理耗时(查库、调接口、写缓存)
- 提交消费位点的耗时
正常情况下一批消息(默认32条)从拉取到消费完成,总耗时应该在几百毫秒到1秒之间。如果某一段耗时异常,比如反序列化就要几秒,十有八九是消息体太大或者序列化方式低效(比如用了JDK原生序列化处理大对象)。
3.4 确认是不是消费位点提交的问题
还有一种经常被忽视的堆积原因:消费成功了,但位点没提交成功。RocketMQ有同步提交和异步提交两种模式。如果业务逻辑里把consumeMessage的返回结果异常处理了,或者CONSUME_SUCCESS没有正确返回,Broker会认为消息还没消费成功,继续往消费端投递,同时Lag一直不降。但消息总是重复投递确实会导致消费速度叠加变慢,因为大量时间和资源都耗费在处理重复消息上。
我遇到过一个小伙伴,消费逻辑里用了try-catch吞掉了所有异常,然后不管什么情况都返回CONSUME_SUCCESS。表面上看没什么问题,但如果业务处理实际失败了,消息也算消费成功了,不重试——这其实是丢消息,不是堆积。反过来,如果你是先处理业务、再提交位点,而业务处理抛异常后没有正确返回RECONSUME_LATER,也可能导致消息一直重投。
4. 核心处理手段:按紧急程度分级,不要一上来就动架构
处理堆积方案很多,但要讲究顺序和取舍。我的原则是:先止血,再治本。先让堆积不再恶化,保住业务,再去想架构上的修改。
4.1 最直接的手段:水平扩容消费者实例
如果确认是消费端整体处理能力不足,最简单有效的方式就是增加消费者实例。但这里有个容易踩的坑:消费者实例数量不能超过队列数量。
RocketMQ的分配策略决定了,一个队列同一时间只能被一个消费者实例消费。假设Topic的读队列是16个,你起了20个消费者实例,其实只有16个实例在干活,剩下4个空闲。所以扩容之前先看队列数,实例数压到队列数以内。
如果队列数本身就少,比如只有8个,那你扩容也没用,因为并发上限就8。这种情况要考虑增加队列数,但要注意:动态增加队列数后,同一个消费组里的消费位点会重新分配,可能引发少量消息重复消费,需要在消费逻辑里做好幂等。
从实际经验看,扩容消费者实例的效果是最立竿见影的。之前一次促销活动,一个消费组处理订单消息,原本6台实例,Lag持续增长。我直接临时扩到10台(配合扩容Topic队列到16),十分钟内Lag从几十万降到了几千。前提是消费逻辑没有瓶颈(比如数据库连接池不够),否则扩再多实例都是白搭,因为它们都在抢同一个连接池。
4.2 优化单条消息消费耗时
如果消费逻辑本身有性能问题,扩容只能是暂时的。比如我们之前有个消费逻辑很重,每条消息要调用三个下游系统,同步串行执行,总耗时800ms。后来改成并行调用三个下游,总耗时降到250ms,消费能力直接提升了3倍多,堆积压力瞬间小了不少。
再比如批量消费。RocketMQ的ConsumeMessageService一次最多拉32条消息,如果你的消费逻辑是一条条处理,建议改成批量处理——拿到List<MessageExt>后批量查数据库、批量写缓存、批量调下游,很多时候性能提升非常明显。
还有一种优化思路:消费逻辑里面不要做太重的写入操作。消息消费路径上的数据库写入,尽量削峰。比如把高频的写操作合并成批量写,或者用异步线程池去处理非核心链路,保证消费主线程能快速处理完一批消息后马上拉下一批。
4.3 临时应急方案:跳过堆积的旧消息
如果堆积量非常大,比如几百万条,按正常消费速度要跑几个小时才能消化,而业务方已经明确说"这些旧消息已经过期了,处理了也没意义",这时候你就得考虑跳过堆积消息了。
RocketMQ提供了重置消费位点的功能,可以把消费位点直接推进到最新位置,跳过中间堆积的消息。命令是这样:
mqadmin resetOffsetByTime -g group_name -t topic_name -s timestamp-s参数填一个时间戳,比如你想让消费者从当前时间开始消费,就把时间戳设为当前时间。这样消费者会跳过所有旧消息,直接从最新消息开始拉取。
但这个方法要特别谨慎,只能在你确定旧消息不需要处理时使用。如果不确定,千万别乱重置位点,否则消息源丢失的锅就甩你头上了。我见过一个事故,运维误操作把位点重置到了未来时间,所有消息都投递不了了,排查了半天最后才发现是位点被重置了。
4.4 降级策略:给非核心消费让路
还有一种场景:堆积的核心原因是生产端流量暴增,而且这种暴增会持续一段时间。这时候如果消费端怎么扩都跟不上,就得考虑降级了。
我常用的方案是优先级分离:给消费端设置一个可以配置的"丢弃策略",对于非核心业务的消息,如果消费积压超过一定阈值(比如消息的存活时间超过5分钟),直接丢弃不再处理。核心业务消息照常处理。这样消费端的压力会集中在最核心的数据上,保证最有价值的数据不丢。
这个方案听起来有点"暴力",但在大促高峰期确实比死扛有效。毕竟非核心消息晚处理一分钟和丢掉没本质区别,不如把资源让给核心链路。
5. 从根上防堆积:架构和设计层面的自我救赎
应急处理只是治标,真正要解决堆积问题,得在设计层面就把"堆积的可能性"降到最低。分享几个我一直坚持的实践。
5.1 消费幂等是底线,不是选项
前面提到扩容、重置位点、动态调整队列都会引发重复消费,所以消费逻辑的幂等性是绝对不能妥协的。我见过太多系统上线时不做幂等,等到堆积发生需要扩容时,一扩容就出现大量重复数据。
推荐做法:消费逻辑开始时先查一下唯一键是否已经处理过,处理过就直接返回成功。这个唯一键可以是业务订单号、流水号,甚至可以用消息体的哈希值,但一定要结合具体业务设计。
5.2 给消费端加"背压"闭环
RocketMQ本身没有提供消费速率限制的API,但你可以在消费端自己控制。我目前的做法是:消费者拧开线程池的Semaphore,限制最大并发消息数。如果下游系统已经出现性能问题,并发数自动调低,宁可堆积,也不要把下游打死。
维护一个动态参数配置,每分钟检查一次下游系统的健康状态(比如接口错误率、响应时间),如果异常就自动降低消费并发,等下游恢复再调回正常值。这套背压机制帮我避免了好几次"下游抖动引发消费端雪崩"的事故。
提示:RocketMQ 5.x版本的SimpleConsumer支持
suspend和resume操作,比4.x的DefaultPushConsumer更容易实现流控,建议新项目优先用5.x的消费模型。
5.3 消息粒度的取舍:宁可小不可大
一条消息包含的数据量太大,是堆积的隐形帮凶。消息体大,传输耗时增加,反序列化耗时增加,存储占用增加,消费处理也可能变慢。更关键的是,如果一条超大消息在处理中频繁触发重试,重试带来的网络和CPU开销是成倍增长的。
所以我一般建议:消息体里只放关键标识和必要参数,具体数据由消费者按需去查。比如订单消息,只放订单ID和变更类型,消费者收到后按ID去订单服务查询全量数据。这样消息体从几KB降到几百字节,消费效率提升非常明显。唯一要注意的是,消费时突然查不到数据(比如数据被删了)的处理逻辑。
5.4 延迟消息别滥用,容易"无声堆积"
RocketMQ的定时/延迟消息好用是好用,但用多了会有一个坑:延迟消息在Topic里积压是"不可见"的,它们不进入消费队列,直到延迟时间到了才被投递。
如果你有大量精度要求不高的延迟消息(比如定时通知、定时任务),建议评估一下是否真的需要延迟消息。之前我们有一个场景,每天有几百万条延迟30分钟的消息,高峰期Broker里积压的延迟消息占存储的30%。后来改成业务数据库里记录延迟时间,由调度任务轮询扫表,Broker压力小了很多。
6. 面试时怎么把"消息堆积"讲出层次感
既然标题是"面试被问到",最后分享一点面试表达的心得。消息堆积这个问题,面试官想听的绝对不是一个"扩容"答案,而是你的思路链路。
我建议按这个层次来回答:
第一层:怎么发现堆积。监控体系、消费位点Lag、告警阈值怎么设,体现你有线上运维经验。
第二层:怎么定位根因。从消费组Lag分布看到消费者实例,从jstack看线程状态定位卡点,从埋点数据定位耗时慢的环节,体现你会排查问题而不是瞎猜。
第三层:怎么应急处理。扩容、优化消费逻辑、跳过旧消息、降级策略,每种方案的适用场景和风险是什么,体现你懂取舍、有全局观。
第四层:怎么长期治理。幂等设计、背压机制、消息粒度优化、延迟消息治理,体现你对系统稳定性的思考深度。
这样一套回答下来,面试官基本上能确定你是在生产环境真刀真枪处理过堆积问题的人,而不是背了一堆面试八股。
我个人这几年的体会是,消息堆积这个问题,一半是技术问题,一半是管理问题。技术问题在于你能不能快速定位、快速止血;管理问题在于事前有没有做好监控、预案、容量评估。处理堆积最顺利的一次,我们从头到尾没做任何架构调整,完全靠监控告警发现得早、扩容预案执行得快、消费逻辑的瓶颈定位得准,半小时内就恢复了。反而是最有"技术含量"的大规模架构改造,往往是在堆积反复发生之后才被逼着做出来的。希望看完这篇文章的你,下次遇到堆积,能比当时的我更从容一点。