java 技术随笔

消息积压排查与治理实战:从发现 consumer lag 到扩容消费与毒消息处理

凌晨两点告警电话响了:"消费积压 XX 万条!"你爬起来打开控制台,消费组 lag 疯狂上涨,下游接口还时不时超时。消息积压是 MQ 最常见的线上事故,但它不是随机事件,而是"消费速度 < 生产速度"的必然结果。本文先讲清积压是怎么形成的、怎么第一时间发现,再给出从应急到根治的完整打法。

先分清:积压、堆积、消息延迟是一回事吗

消费组 lag(落后条数)大于 0 就是积压。轻度积压(几百条、几秒追平)不用慌,真正危险的是持续增长的积压:消费能力持续低于生产,积压只会越滚越大,最终拖垮下游甚至触发告警停机。积压问题的本质是消费者处理不过来,不是 MQ 本身坏了。

消息为什么会积压:五大常见原因

1. 消费者实例太少或线程数不够。一个消费组只有 1 个实例,或每个实例消费线程数 = 1,而 Topic 有 32 个分区,等于 31 个分区没人消费——这是最常见的初级配置失误。

2. 单条消息处理太慢。消费逻辑里有同步调外部接口、查数据库、做复杂计算,单条耗时从 10ms 涨到 500ms,吞吐直接掉一个量级。

3. 下游依赖故障。消费里调的下游服务超时、数据库连接池打满、Redis 抖动,每条消息都在重试,整体处理速度趋近于 0。

4. 某条"毒消息"卡死线程。一条消息触发死循环或长时间阻塞,把消费线程占死,后面的消息全部排队。

5. 消费组频繁重平衡。Kafka 消费端处理超时被判定失活踢出组,触发 Rebalance,Rebalance 期间停止消费,恢复后可能又被踢出,形成恶性循环。

怎么第一时间发现积压

不要等用户投诉。生产环境必须建立三层观测:

第一层:指标监控。RocketMQ 控制台或 Kafka 的消费组 lag、消费 TPS、消费耗时是核心指标,配合 Prometheus + Grafana 采集,lag 超过阈值(比如 1 万)就告警。

第二层:日志。消费端记录每条消息的处理耗时,超过 P99 基线的打 warn 日志,快速发现"变慢的拐点"发生在哪次发版之后。

第三层:主动对账。核心业务(如支付回调)维护消息对账表,定期比对"我发了多少条 vs 我消费成功多少条",能发现监控没覆盖的静默丢失。

应急处理:积压已经发生怎么办

先止血,再排查,顺序很重要:

第一步:扩容消费者,而不是盲目重启。消费组扩容规则:Kafka/RocketMQ 的消费者数量最多到分区数,超过分区数再多实例也闲置。所以先确认当前消费者实例数与分区数,把实例数加到接近分区数。如果分区数太少(比如 Topic 只有 4 个分区),加机器没用,得考虑临时加分区。

第二步:定位瓶颈在"上游太快"还是"下游太慢"。看监控:生产 TPS 正常、消费 TPS 低,说明消费者自身慢;消费 TPS 高但 lag 还涨,说明真是生产洪峰,可能需要削峰限流而不是优化消费。

第三步:处理毒消息。如果某条消息一直消费失败,避免它阻塞队列:超过 N 次重试进死信队列(RocketMQ 默认死信),或跳过并记录,事后再人工补处理。查看消费日志里的异常堆栈,通常能立刻定位是哪条、哪个逻辑。

第四步:下游扛不住时的降级预案。如果是因为下游服务挂了导致消费全慢,应急期间先关闭对下游的调用(降级),消息落本地或先记录,等下游恢复后补偿,避免消费线程全部阻塞在下游超时上。

根治思路:让消费速度追得上生产速度

应急只能缓解,真正要解决的是把消费能力做上去:

提升单条处理速度。把消费逻辑里的同步外部调用改为异步化或本地缓存;批量拉取批量处理(一条 SQL 批量更新代替逐条 update);无关紧要的旁路逻辑(写日志、统计)从主链路剥离。

并行化消费。增加分区 + 增加消费者实例/线程,让处理能力横向扩展。注意 RocketMQ 顺序消息只能单线程消费一个队列,想并行就得按业务设计拆分队列。

消费线程池隔离。给消费逻辑配置独立的线程池和队列,避免消费任务互相挤占;不同重要级的 Topic 用不同线程池,防止一个慢 Topic 拖垮整个应用。

控制生产节奏。大促、定时任务集中触发导致的洪峰积压,要在生产端做削峰:业务侧限流、分批发送、错峰执行,而不是把压力全留给消费者。

终极兜底:积压到无法追平时怎么办

极端场景(积压几千万条,追平需要几天)时,别傻等消费完——旧消息可能已经失去时效性。成熟的方案是临时扩容队列 + 多组并行消费:新建一个临时 Topic(分区数是原来的 N 倍),写一个转发消费者把积压消息按原顺序灌到临时 Topic,再启动 N 倍消费者的临时消费组去消费,处理完删掉临时资源。这个"扩充队列 + 并行消费"的方案是各大厂积压治理的标准动作。

复盘:一次积压事故的标准排查时间线

收到告警 -> 控制台看 lag 与消费 TPS 确认积压程度 -> 看消费端异常日志找毒消息/下游超时 -> 应急扩容消费者 + 毒消息进死信 -> 下游恢复后观察 lag 回落 -> 复盘定位根因(发版引入的慢逻辑?下游变更?)并加监控阈值。

记住一句话:积压不可怕,可怕的是没有监控、没有预案、消费能力没有横向扩展的空间。把消费做成可扩容、可观测、可降级的,积压就永远只是"小插曲"而不是"大事故"。

标签
消息队列线上排查高并发运维