1. 先别急着加机器,把“堆积”这件事看透
RocketMQ 消息堆积,几乎是每个做电商、做交易、做日志收集的团队都绕不开的坎。面试官问这个问题,表面上考的是“你会不会扩容 Consumer”,实际上想听的是一整套排查链路:堆积到底卡在哪个环节?是 Broker 写不进去,还是 Consumer 拉不动?是消费逻辑本身太慢,还是下游依赖把线程池拖死了?答不出这层,背再多“增加 Consumer 数量”的八股都没用。
先明确一个概念:RocketMQ 的堆积,本质上就是Consumer 的消费速率长期低于 Producer 的生产速率,导致消息在 Broker 的 CommitLog 里越攒越多,消费位点(ConsumerOffset)和最大位点(MaxOffset)之间的差值不断变大。这个差值就是我们常说的“堆积数量”,在控制台里能看到曲线一路飙升。
但这只是表象。真正要命的是,堆积往往不是单一原因,而是多因素叠加。我遇到过最典型的一个案例:业务方说“订单消息堆积了 10 万条”,我上去一看,Consumer 的线程池配的是 20,消费单条消息要调远程接口,接口 P99 是 300ms,单线程每秒只能吃 3 条,20 个线程满打满算也就 60 TPS。而 Producer 那边峰值每秒写入 500 条,堆积自然只会越来越多。这种情况你哪怕把 Consumer 扩到 100 个实例,如果线程数不变、下游耗时不变,照样白搭。
所以,处理堆积的第一步永远是先定位,而不是先动手。这里我建议按三个层面去排查:
- Broker 层:看 Broker 的写入性能是否正常,有没有频繁 GC、磁盘 IO 打满、PageCache 压力过大。
- Consumer 层:看消费线程数、消费逻辑耗时、拉取消息的阻塞情况、线程池活跃度。
- 下游依赖层:看消费时调用的数据库、Redis、远程 RPC 接口是否存在慢查询、连接池耗尽、超时重试等问题。
这三层里,Consumer 层是 80% 堆积问题的根源。为什么?因为 Broker 的写入路径是顺序写盘 + PageCache 加速,只要磁盘没坏,写入吞吐基本不是瓶颈;真正拉垮的往往是消费端那种“单条消息处理太重”的设计。比如你在消费回调里做了同步的 HTTP 调用,每次等 2 秒超时,那线程池瞬间就被占满,堆积就是必然的。
还有一个容易被忽略的点:RocketMQ 的消费并发模型。如果你用的是 DefaultMQPushConsumer,它内部是一个拉取线程不断向 Broker 拉消息,然后丢进一个线程池去执行消费回调。线程池默认的核心线程数、队列大小你都可能没调过,结果就是拉取线程拉得飞快,线程池里的任务却处理不过来,消息全堆在本地队列里,Consumer 进程本身就成了堆积点。这种情况连堆积曲线都看不出异常,因为 Broker 端没积压,但消费延迟照样高。
我自己处理这类问题有个习惯:先看 Consumer 的线程池活跃度,再看消费单条消息的平均耗时,最后才看堆积总量。把这三个数字列出来,问题基本就水落石出了。下面把每个环节的具体套路拆开讲。
2. 定位堆积根源:三步排查法,把瓶颈钉死
2.1 第一步:区分“生产堆积”还是“消费堆积”
很多人一看到堆积就急着去扩容 Consumer,这是典型的“头痛医头”。实际上,RocketMQ 的堆积可能发生在生产端,也可能发生在消费端,两者的处理方式完全不同。
生产堆积的表现是:Producer 发送消息的耗时持续走高,重试频繁,Broker 的写入 TPS 明显下降。这种情况通常是 Broker 磁盘性能不行,或者 PageCache 被写脏数据挤占得太厉害。你可以用mqadmin topicStatus命令查看 Topic 的写入情况,或者直接监控 Broker 的写入延迟指标。
消费堆积的表现则恰恰相反:Broker 写入 TPS 正常,但 ConsumeQueue 的消费进度落后越来越多,Consumer 实例的消费 TPS 远低于生产 TPS。这是最经典的堆积场景,后面所有步骤都围绕它展开。
怎么快速区分?很简单,看两个指标:Broker 端的消息生产速率,和Consumer 端的消息消费速率。前者大于后者,堆积必然发生;两者持平,堆积数量会维持不变;前者小于后者,堆积会被慢慢消化。把这组数据拉出来一对比,方向就有了。
2.2 第二步:Consumer 侧细看,线程数、耗时、阻塞三件套
确定了是消费堆积,接下来就要把 Consumer 解剖开。
第一看线程数。如果你用的是DefaultMQPushConsumer,setConsumeThreadMin和setConsumeThreadMax这两个参数决定了消费线程的上下限。很多时候大家都只设置了 Consumer 的实例个数,忽略了单实例内部的线程数。假设你有 10 个消费实例,每个实例 20 个线程,那就是 200 个并发;如果每个实例只默认了 5 个线程,那就是 50 个并发,差距四倍。所以扩容 Consumer 之前,先确认单实例的线程池是否拉满。
第二看消费耗时。给消费回调里的每个方法调用打上耗时日志,或者用 Arthas 的trace命令看消费链路里哪一步最慢。我见过最夸张的案例:消费逻辑里同步调了三个下游服务,每个都设置了 5 秒超时,三个串行下来最坏要 15 秒才处理完一条消息。你想想,就算给你 100 个线程,100 条消息同时卡在远程调用上,整个消费集群基本等于瘫痪。这种问题靠加机器是治标不治本,必须改异步化、批量处理或者走降级方案。
第三看阻塞。RocketMQ 消费线程池默认用的是LinkedBlockingQueue,队列长度默认是 1000。如果线程数不够,任务就会在本地队列里排队,排队越多,消费延迟越高。这个本地队列的积压,Broker 端是看不到的,只有通过 JMX 查看线程池的队列大小才能发现。我建议把消费线程池的监控指标挂到监控平台,重点盯queueSize、activeCount、completedTaskCount三个值。
2.3 第三步:下游依赖的体检,数据库和远程调用不能漏
消费逻辑一般都要读写数据库、调 RPC、操作 Redis,这些下游服务的健康状况直接决定消费速率。
数据库方面,最常见的坑是慢 SQL和连接池打满。每条消息都要执行一次 INSERT 或 UPDATE,本来应该走索引秒回,结果因为没加索引、或者查询条件用了函数导致全表扫描,一次就要 500ms。200ms 的延迟可以忍,500ms 以上就会肉眼可见地拉低消费速率。用慢查询日志把这类 SQL 揪出来,加上索引、优化 SQL,通常能立竿见影。
Redis 方面,坑在大 Value和热点 Key。如果你在消费逻辑里对同一个热点 Key 做频繁读写,Redis 单线程模型下所有请求都会被串行化排队,延迟从 1ms 飙到 100ms 都是常见的。这种情况要么把热点 Key 拆散,要么做本地缓存,要么把聚合操作放到消费完成后异步处理。
远程 RPC 方面,最怕的是超时重试风暴。消费线程调 RPC 接口,超时设为 3 秒,失败后框架自动重试 3 次,那一条消息在最坏情况下要 12 秒才能确定失败。更重要的是,RPC 超时后线程不会立即释放,它要等够超时时间才走重试逻辑,这段时间里线程池资源被死占。我见过一个团队,消费逻辑里调了个成功率只有 60% 的接口,堆积曲线直接飙升到 200 万条,最后把接口修好,堆积才慢慢消化掉。所以,消费逻辑里对远程调用的超时时间、重试次数一定要做严格限制,该降级就降级。
3. 应对堆积的核心手段:扩容、重平衡、重置位点三把刀
定位完原因,接下来就是动手处理。根据堆积的程度和业务容忍度,我一般把手段分成三级:轻量级扩容、主动重平衡、重置消费位点。逐个说。
3.1 轻量级扩容:先调线程池,再考虑加实例
如果堆积数量不大,比如几万条,而且下游依赖一切正常,这时候最经济的手段是先调 Consumer 的消费线程数。
假设现在是 10 个消费实例,每个实例 10 个线程,一共 100 并发。你把每个实例的consumeThreadMax从 10 调到 40,并发数就变成了 400,消费速率理论上可以翻四倍。但注意一个前提:下游凭依要能扛住这四倍的压力。数据库连接池够不够?远程接口的 TPS 上限是多少?如果下游扛不住,调线程数只会让消费从“堆积”变成“大量超时失败”,情况反而更糟。
如果线程数已经是合理的(比如每个实例跑 20-30 个线程),堆积量仍然在涨,那就需要加消费实例了。加实例要注意 RocketMQ 的负载均衡机制:同一个 ConsumerGroup 下的实例会均摊 Topic 下的队列(MessageQueue),默认的AllocateMessageQueueAveragely策略会把队列平均分配给每个实例。
假设 Topic 有 16 个队列,当前有 4 个消费实例,那每个实例分到 4 个队列。你加到 8 个实例,每个实例分到 2 个队列,单实例负载减半,消费速率自然提升。但这里有个脑袋坑:如果队列数本身少于实例数,加实例就没用了。16 个队列最多只能被 16 个实例平分,加到 32 个实例也只会出现一部分实例没有队列可分,白白浪费资源。所以扩容之前,先确认 Topic 的队列数是否 ≥ 实例数。如果队列数不够,还得先加队列,但队列数只支持在创建 Topic 时指定最大数量,运行期可以改但影响面较大,建议提前规划好队列数。
3.2 主动重平衡:让新实例尽快接活儿的技巧
默认情况下,新实例启动后要等 RocketMQ 的 Rebalance 机制自动触发,这个时间默认为 20 秒(RebalanceInterval)。如果你在业务高峰期扩容,这几秒甚至几十秒的延迟都会让堆积进一步加剧。
想要尽快让新实例参与消费,可以调整两个参数:pollNameServerInterval和heartbeatBrokerInterval。前者控制客户端轮询 NameServer 获取路由信息的时间间隔,默认 30 秒;后者控制向 Broker 发送心跳的频率,默认 30 秒。在集群规模不大、不会给 NameServer 造成压力的情况下,可以把这两个值调小到 10 秒甚至 5 秒,新实例能在更短的时间内获取到最新路由并触发重平衡。
另外还有一个容易忽略的骚操作:直接调用consumer.rebalance()方法。RocketMQ 客户端 API 里并没有一个公开的rebalance()接口让你手动触发,实际的 Rebalance 逻辑在RebalanceImpl里,普通 API 拿不到。所以实操上更靠谱的做法是:把新实例的启动时间错开,先启动一个,观察它的消费情况,确认正常后再启动下一个,避免多个实例同时触发 Rebalance 造成队列分配的抖动。
这里面还有个真实遇到过的情况:新实例启动后,明明同一个 ConsumerGroup 下已经有多个实例存在,但新实例的消费 TPS 很低,甚至一直没有消息进来。原因通常是重平衡还没触发,或者是订阅关系不一致。比如新实例的subExpression写的是TagA,而老实例写的是TagB,两者订阅的 Tag 不同,重平衡时就会把队列分给不同 Tag 的实例,导致看起来“新实例没消息”。这种问题排查起来很反直觉,后来我们养成了习惯:新实例的订阅配置一律从配置中心复制,绝不手填。
3.3 重置消费位点:堆积太多时的兜底大杀器
当堆积量已经达到几百万甚至上千万条,靠调线程、加实例慢慢消费可能需要好几个小时,业务等不起。这时候有两个选择:跳过堆积消息或者从最新位点开始消费。
跳过的做法是重置消费位点。在 RocketMQ 控制台的 Consumer 管理页面里,可以对指定 ConsumerGroup 的消费位点进行重置。重置有两种模式:
- 从最新位点开始消费:也就是跳过所有堆积消息,只消费新消息。适合堆积消息已经失去业务价值、或者后续消息依赖后续状态才能处理的场景。
- 从指定时间点开始消费:比如你确定 13:00 之后的消息是完整的,可以把消费位点重置到 13:00 左右,丢弃之前的脏数据或无效数据。
重置位点是高风险操作,执行前必须确认三个问题:
- 被跳过的消息是否真的可以丢弃?如果下游对数据完整性有强要求,比如交易流水、订单状态变更,跳过之后会产生数据不一致,这种场景不能用。
- 是否有多个消费实例在跑?重置位点是针对 ConsumerGroup 级别的,如果多个实例负载不均衡,重置后可能有的实例消费进度被调整,有的没有,导致重复消费或者漏消费。
- 业务方是否已经确认?任何重置操作都应该走变更流程,至少要有业务负责人的口头确认和邮件记录。
我遇到过一次紧急情况:线上堆积了 2000 万条某业务 Topic 的消息,消息内容是一个小时内拍下的商品链接快照,但实际上这些快照在数据库里都有实时状态,堆积消息里的数据已经过期。当时经过业务确认后,直接把消费位点重置到了最新,堆积归零,整个消费集群恢复了健康。这个例子很极端,但足以说明:重置位点是处理大规模堆积时最有效的兜底手段,没有之一。
4. 治本之道:消费逻辑优化和链路治理的实操细节
堆积极少是“运气差”造成的,绝大多数是设计和代码的债。跟在后面一次次应急,不如把消费逻辑本身变稳。分享几个我自己压箱底的做法。
4.1 消费逻辑里少做重活儿,能异步就异步
RocketMQ 消费回调里,你做的事越少,消费速率越快,这是铁律。常见的高危操作包括:
- 在回调里同步调用第三方 HTTP 接口
- 在回调里做复杂计算、大文件读取
- 在回调里批量执行数据库写操作但不分批
- 在回调里依赖另一个 MQ 的消费结果
这里最典型的反面案例就是“消费回调里发短信/发邮件”。很多人觉得发短信很快,实际上短信渠道的 HTTP 接口延迟经常在 500ms 到 2 秒之间,而且短信服务在高峰期还会限流。如果把发短信写在消费回调里,消费速率瞬间被拉低。我的做法是:消费回调只负责把消息内容持久化到本地表,再发一个内部标记,真正的短信发送由定时任务或单独的异步线程池去处理。这样消费回调全程耗时控制在几毫秒以内,堆积永远不会发生。
还有一类重活儿是数据库写入。如果消费逻辑要批量插入数据,强烈建议用batch方式,一次插入 100 条甚至 500 条,而不是一条一条 insert。假设单条 insert 是 2ms,100 条就是 200ms;改成 batch 后,100 条可能只需要 10ms,性能提升 20 倍。如果你是用 MyBatis,写一个<foreach>批量 insert 的 SQL 并不难,收益却非常夸张。
4.2 重试要有上限,死信队列是最后的收尸人
消费失败导致的重试,如果不加控制,会无限占用消费线程,把堆积越拖越重。RocketMQ 默认的重试机制是:消费失败后,会按照RETRY_TIMES重新投递,默认最多重试 16 次,每次间隔时间逐步拉长。这本来是保护机制,但如果你在消费回调里一直抛异常,又不理阴界限,16 次重试会消耗大量消费线程和网络资源。
我的建议是:业务消费失败要区分“可重试”和“不可重试”。对于订单状态没同步之类的不一致数据,可以重试;对于消息本身格式错误、参数非法之类的数据,重试 100 次都是白搭。这类消息应该直接捕获异常,记录日志,然后走 RocketMQ 的死信队列(DLQ),让专门的运维人员定期处理。
死信队列的名字一般是%DLQ%ConsumerGroupName,里面的消息有完整的原始内容和消费失败原因。我通常会写一个死信队列的消费程序,把这些消息按照业务类型自动分类,自动恢复的就重新投递到原 Topic,无法自动恢复的钉钉告警,让研发手工介入。这套机制跑通之后,堆积问题至少减少一半。
4.3 消费链路里做“水位预警”,别等堆起来才动手
我在每个消费集群上都设了四道预警线,四道线的经验值是:
| 预警级别 | 触发条件 | 响应动作 |
|---|---|---|
| 提示 | 堆积数量 > 5000 | 检查消费速率是否正常 |
| 警告 | 堆积数量 > 50000 | 评估是否需要扩容 Consumer |
| 严重 | 堆积数量 > 100000 | 拉业务方和运维一起排查链路 |
| 紧急 | 堆积数量 > 上百万且持续增长 | 执行重置位点或临时停 Producer 限流 |
监控体系是处理堆积的第一道防线,等招聘的面试官都开始问“堆了怎么办”的时候,你肯定不希望是自己在生产环境里慌忙救火。
水位预警怎么做?最简单的方式是创建消费者来获取消费位点,再通过 RocketMQ 的mqadmin consumerProgress命令定期采集,计算每个 ConsumerGroup 的堆积差值,写入监控系统。更现代的做法是接入 RocketMQ Dashboard,它自带消费进度监控和告警功能,界面直接能看到某个 ConsumerGroup 的消费延迟。如果你是自研的监控系统,那核心逻辑就是:埋点采集ConsumerOffset和MaxOffset,两者之差超过阈值就触发告警。
四道水位的预警级别里,第三级开始就必须有明确的处理 owner,第4级需要触发紧急预案流程。这个过程一定要固化到文档里,不然每次都临时拉人开会,效率极低。
5. 从一次真实的线上事故看完整处理流程
理论讲了这么多,还是用一个完整案例把整个过程串起来,读起来更直观。
某天下午 3 点左右,监控告警弹出:交易中心某个 Topic 的 ConsumerGroup 堆积量在 20 分钟内从 2 万涨到了 33 万。这个 Topic 承载的是订单状态变更消息,消费逻辑里会同步调用户中心的 RPC 接口获取用户信息,然后写库。
我当时的排查流程是这样的:
第一步,打开监控看生产速率和消费速率。生产速率大约每秒 800 条,消费速率只有每秒 120 条,生产远大于消费,堆积在预期内,问题出在消费侧。
第二步,看消费线程池指标。发现 30 个消费线程里面有 28 个处于RUNNABLE状态,但是从线程栈看,几乎全部卡在HttpURLConnection.getInputStream上,也就是远程 RPC 调用处。这不难判断:用户中心接口的耗时在飙升。
第三步,查看用户中心接口的监控,发现 P99 从平时的 80ms 涨到了 3 秒,明显是下游出问题了。联系用户中心负责人后发现,该接口所在的集群下午发布了一个版本,带了一个性能回退的 bug,导致接口处理变慢。
我当时没有直接重启消费进程,而是先做了两件事:第一,把消费逻辑里的远程调用超时时间从 3 秒改到 1 秒,超时直接走降级方案(本地缓存用户信息);第二,临时把消费线程数从 30 调到 60,尽量提高短时消费速率。这样改完之后,消费速率从每秒 120 条慢慢回升到了 400 条左右,堆积曲线开始掉头向下。用户中心那边修复后又过了一个小时,消费速率恢复到每秒 750 条左右,堆积彻底清完。
这中间有一个细节:改消费线程数可以在运行期内通过updateConsumerThread向 Broker 提交新的线程数配置吗?不能,这个参数是 Consumer 启动时就固定的,你要改它必须重启消费实例。所以我当时是起了一个新的 ConsumerGroup 临时实例,或者直接改了配置重启。不管哪条路,都要接受短时间内可能重复消费几万到几十万条消息的风险。
针对这种运行期调整需求,我们后来做了一套通用方案:消费线程数、消费超时时间、远程调用超时时间全部做成配置中心动态配置,调整时不用重启消费进程,直接推送新配置,线程池和客户端参数实时变更。这个方案上线后,处理堆积的响应时间从半小时缩到了一分钟。
这个案例给我们的启发:堆积问题的根子往往不在 RocketMQ 本身,而在离它最近的消费代码和下游依赖上。你能控制的是消费线程数、超时时间、批量大小、异步化程度,但下游能不能扛住压力,取决于你们团队的链路治理水平。这也是为什么很多团队用上了 K8s 自动扩容,堆积照样还是会发生。
6. 面试官想听什么?从应答框架到避坑细节
既然题目是“面试官问消息堆积怎么处理”,我猜读者里至少有半数是想把这个问题答到面试加分水平的。那就按面试场景拆一下,哪些话该说、哪些坑别踩。
6.1 回答框架:先定位,再扩容,再治本
面试时听到这个问题,千万别一开口就“加机器”。我建议按下面这套逻辑递进回答,既有思路、又有细节、还能体现经验:
- 先描述堆积的本质:生产速率 > 消费速率,导致未消费消息在 Broker 上积压。
- 再给出定位链路:先看生产速率和消费速率的差值,再查消费线程池活跃度、消费单条耗时、下游依赖耗时,最后判断是哪一个环节成了瓶颈。
- 然后说应急处理:如果是临时问题(比如下游抖动),可以调消费线程数、加消费实例、用死信队列接收失败消息,甚至重置消费位点跳过无用消息。
- 最后说长效治理:代码层面做异步化、批量写、失败快速降级;监控层面部署消费水位告警;压测层面定期模拟全链路压力测试。
这套框架能充分展示你不光会背 API,而是真的能对整个消息链路做诊断。
6.2 加分项:主动聊“顺序消费”和“幂等”
如果你能在答案里自然地带出 RocketMQ 的两个进阶主题,面试官会有明显的好感。
一个是顺序消费。顺序消费和堆积是天然矛盾的:它要求同一个 Queue 里的消息只能被同一个消费线程串行消费,所以哪怕你有 100 个消费线程,能用的也只是一个。面试时你可以说:“如果业务要求全局顺序消费,那堆积的应对手段会受限,这时候更多要考虑把 Topic 分区设计得合理,把不需要全局顺序的业务拆出来。” 这段能体现你想得比“贪多求快”更深一层。
另一个是幂等消费。堆积处理过程中,重置位点、重复投递、失败重试都会导致消息被消费多次,如果消费逻辑不是幂等的,就会产生重复数据。我建议你在消费回调里顺手实现一组幂等校验,比如用唯一业务 ID 查 Redis 是否存在,存在就直接返回消费成功,不在才执行真正的逻辑。这样无论消息被重复投递多少回,都不会产生脏数据。面试时提到这一点,基本就是“有真实生产经验”的明证。
6.3 避坑点:千万别踩的三个雷
面试中常有候选人聊得兴起,最后死在两个细节上。我也列出来提醒自己团队的人,免得赔了技术又折了 Offer。
- 雷区一:张口就说“日志消费端堆积,加 Kafka partitions”。这是概念混淆,RocketMQ 里对应的是 Queue,不是 Partition。虽然两者都是分区概念,但说错会显得你底子不牢。
- 雷区二:把“消息堆积”说成“消息丢失”。堆积的根本特征是消息还在 Broker 上,只是来不及消费,数据没有丢。如果回答里把这两者混作一团,说明你压根没理解 RocketMQ 的存储模型。
- 雷区三:重提“延迟消息”时无脑消耗。如果业务用到了延迟消息,堆积排查时要把延迟队列的调度情况单独拉出来看,不能一体化运维。RocketMQ 的延迟消息是存在单独的 Schedule 队列里的,它在到达执行时间之前不会进入真正的消费队列,所以它产生的“堆积”是预期行为,不需要处理。懂得区分“预期堆积”和“异常堆积”,也是区分菜鸟和老鸟的信号之一。
7. 实操中积累的监控脚本与习惯
最后分享几段我日常排查堆积时必用的命令和脚本,都是可以直接抄作业的东西。
7.1 用 mqadmin 快速看堆积
RocketMQ 安装包自带的mqadmin命令是排查堆积的第一个工具。按下面的流程快速获取关键信息:
# 列出所有消费者组的消费进度 mqadmin consumerProgress -n {namesrvAddr} # 查看指定消费者组每个队列的消费积压情况 mqadmin consumerStatus -g {consumerGroup} -n {namesrvAddr} # 查看某个Topic的队列分布和写队列数量 mqadmin topicStatus -t {topic} -n {namesrvAddr}consumerStatus会显示每个队列的当前 Offset 和最大 Offset,两者差值就是队列维度上的堆积量。如果你发现某个队列的堆积量明显大于其他队列,那就是队列负载不均,原因可能是生产端的消息 key 全都路由到了同一队列,或者消费端的重平衡把某个队列分给了性能差的实例。
7.2 用 JMX 看消费线程池内部状态
如果想让问题定位更快,给 RocketMQ Client 进程开启 JMX 端口(启动参数加-Dcom.sun.management.jmxremote.port=9999),然后用 jconsole 或者脚本读取以下 MBean:
org.apache.rocketmq.client:type=Consumer,group={group}里的QueueSize、ActiveCount、PoolSizeInvokeGetConsumerStatus可以获取所有队列的消费进度信息
通过这些指标能直接看出线程池有没有被打满、任务在本地队列里积压了多少。我见过的情况是:ActiveCount一直等于PoolSize,说明线程全在干活,基本可以判断是消费逻辑耗时造成的瓶颈。
7.3 写消费告警脚本的通用模板
如果你用的监控平台只能配简单规则,没法直接读 RocketMQ 指标,可以写一个定时执行的脚本去采集堆积量,再对接企业微信/钉钉告警。下面是一个简化的 Python 模板,换掉连接信息就能跑:
import subprocess import json import requests NAMESRV = "127.0.0.1:9876" GROUP = "consumer_group_a" THRESHOLD = 50000 def get_accumulation(): cmd = f"mqadmin consumerStatus -g {GROUP} -n {NAMESRV}" result = subprocess.run(cmd, shell=True, capture_output=True, text=True) print(result.stdout) if __name__ == "__main__": get_accumulation()这个脚本只是骨架,生产环境建议直接走 RocketMQ Dashboard 自带的监控告警,或者用 Prometheus 抓取相关指标。但如果你手头只有裸集群,临时用脚本顶上是完全可行的。
7.4 一个容易被忽略的习惯:关闭消费线程池自动扩展
很多人在DefaultMQPushConsumer里把consumeThreadMin和consumeThreadMax设置成相同值,比如都设成 30。为什么呢?因为 RocketMQ 的线程池在consumeThreadMin和consumeThreadMax之间会按照负载自动调整线程数,但调整频率和策略比较保守,往往跟不上堆积的突增。
把两个值设为相同,等于告诉线程池不要自动伸缩,直接固定线程数。这样线程池的行为可预期,排查问题时少一个变量。实际生产经验也证明,固定线程数比自动伸缩更稳定。
最后分享一个我自己的处理习惯
单独强调一点经验。遇到堆积,我很少先发群消息“消费堆积啦,大家谁的锅”。我的顺序永远是:
- 先把生产速率、消费速率、线程池活跃度、单条消费耗时、下游依赖耗时这五个数凑齐;
- 再判断问题属于“代码慢”、“下游慢”,还是“容量不够”;
- 然后决定是调线程数、加实例、走降级,还是直接重置位点。
这套流程看起来简单,但真正执行过的人会明白,难的不是操作,是冷静地凑齐数据再下判断。很多人一看到堆积就慌,又是重启、又是清消息,最后把问题越搞越大。只要你能按这个节奏来,消息堆积不是大事,它更像是一个信号:你的消费链路哪里设计得不够稳,该优化了。