首页 / 高并发 / 消息队列 / 消息积压
高并发 · 消息队列 · 故障排查

MQ 消息积压问题全解

从「什么是积压」到「为什么会积压」,再到「怎么止血、怎么根治、怎么预防」——一套能直接上手的方法论。

1 什么是 MQ 消息积压 / What is backlog

MQ(Message Queue,消息队列) 用来在系统之间异步传递消息。正常情况下,生产者发得快、消费者也消费得快,队列里基本不囤货。 所谓 消息积压(Backlog / Accumulation),指的是:生产者发送消息的速率持续大于消费者处理消息的速率,导致消息在 Broker(消息中间件)里越堆越多、长时间得不到消费。

一句话理解: 队列就像一个蓄水池,进水(生产)比出水(消费)快,水位(积压量)就会持续上涨。水位不降,就是积压。
生产者 发送快 ⚡ Broker 队列(堆积中) ↑ 水位持续上涨 消费者 处理慢 🐢
图 1:当「进水速度 > 出水速度」,队列中的消息就会不断堆积,这就是积压。

怎么判断「到底积压了没有」?看三个核心指标

指标含义怎么看
堆积量 / Lag队列里没被消费的消息条数(Kafka 叫 consumer lag)lag 持续 > 0 且在增长 = 正在积压;lag 回落 = 在恢复
消费延迟一条消息从生产到被消费经历的时长延迟从毫秒级涨到秒/分钟级,就是积压信号
生产 TPS vs 消费 TPS单位时间生产条数 vs 消费条数生产 TPS 长期 > 消费 TPS,必然积压

2 都有哪些原因 / Why it happens

积压的本质只有一句话:消费跟不上生产。但具体「为什么跟不上」,可以归纳成六大类。先定位属于哪一类,再对症下手。

最常见

① 消费端处理太慢

  • 业务逻辑重 / 计算量大
  • 慢 SQL、全表扫描、缺索引
  • 同步调用外部 RPC / HTTP 且对方慢
  • 锁竞争、死锁、串行化瓶颈
  • Full GC 频繁,线程被「冻住」
扩容

② 消费者实例 / 线程不足

  • 消费量上涨,但消费者没跟着扩容
  • 单线程消费,CPU 多核闲置
  • 线程池过小、被任务打满
  • 容器资源(CPU/内存)被限流
流量

③ 突发流量峰值

  • 秒杀、大促、抽奖瞬时洪峰
  • 批量任务 / 数据迁移集中灌入
  • 生产方重试、补偿逻辑疯狂补发
  • 没有削峰,洪峰直接打到队列
并行度

④ 队列 / 分区并行度不足

  • Kafka partition 数太少
  • RocketMQ queue 数太少
  • 消费者数 > 分区数,多出的消费者干等
  • 单分区成了全局瓶颈
雪崩

⑤ 消费失败与重试风暴

  • 消费抛异常 → 不断重试
  • 重试占用线程 / 拉满资源
  • 死信队列(DLQ)堆积
  • 重试放大流量,越堵越慢
基础设施

⑥ 中间件 / 基础设施问题

  • Broker 磁盘 IO 瓶颈、打满
  • 网络抖动、跨机房延迟
  • Broker 节点宕机 / 主从切换
  • 磁盘满、水位超阈值被限速

3 一般是什么原因导致的 / Top root causes

六类原因里,真实生产环境中 80% 的积压由下面三件事贡献。优先怀疑它们,能省掉大量排查时间。

实战经验:先查这三项,命中率最高。

根因 Top 1:消费逻辑慢(慢 SQL / 同步阻塞)

消费方法里有一条没走索引的 SQL、一次同步调用第三方接口超时、或一个串行循环。单条消息耗时从 5ms 变成 500ms,吞吐直接掉 100 倍。这是积压最普遍、也最容易被忽视的原因。

根因 Top 2:消费者没跟着业务量扩容

日常 1 个消费者刚好够,流量涨了 5 倍却还是 1 个实例、单线程消费。没有水平扩展 + 弹性伸缩,消费能力是硬上限,迟早被压垮。

根因 Top 3:突发流量没有削峰预案

大促 / 秒杀 / 批量任务一上来就是平时几十倍的瞬时流量,队列瞬间灌爆,而消费端没有任何限流、降级、临时扩容的预案,只能眼睁睁看着 lag 飙升。

一句话总结:积压 = 消费慢 × 没扩容 × 突发流量无预案。三件事同时占一个,就容易出事;三个都占,必出事。

4 针对积压如何解决 / Two-layer fix

解决分两层:第一层「紧急止血」是救火,先让 lag 不再涨、开始回落;第二层「根治优化」是治本,避免下次再积压。

🚨 紧急止血(先救火)

  • 临时扩容消费者:增加 consumer 实例数(K8s 直接 scale up)
  • 提高消费并发:多线程 / 线程池消费单条消息
  • 提高并行度:增加 partition / queue 数,让更多消费者并行
  • 降级非核心消息:跳过 / 丢弃可补偿的非关键消息,先保核心链路
  • 限流生产者:暂停或节流发送,给消费端喘息空间(削峰)
  • 建临时 Topic 加速:新开 N 个 partition + N 个 consumer 疯狂消费旧队列,消费完再回写正常链路
注意: 止血手段多是「临时加机器 / 丢消息」,会牺牲成本或一致性,不能长期依赖

✅ 根治优化(治本)

  • 优化消费逻辑:异步化、批处理、消除慢 SQL、加缓存、去掉锁
  • 消费端水平扩展 + 弹性伸缩:HPA 按 lag 自动扩缩容
  • 合理规划分区 / 队列数:容量预留,避免单分区瓶颈
  • 消息分级:核心 / 非核心隔离,互不影响
  • 限流 + 削峰填谷:令牌桶 / 漏桶,洪峰平滑后再入队
  • 幂等消费:保证重复消费不出错,放心重试 / 重放
  • 批量消费 + 批量落库:一次拉一批、一次写一批,吞吐翻倍
目标: 让「消费能力」长期 > 「峰值生产能力」,积压从架构上不可能发生。

5 解决思路与方案(五步法) / Methodology

遇到积压不要慌,按下面五步走,每一步都有明确产出。这是一套可以直接抄作业的排查与处理流程。

① 监控发现 ② 定位原因 ③ 紧急止血 ④ 根治优化 ⑤ 预防演练 lag 告警 / 延迟突增 看耗时·吞吐·外部依赖 扩容·降级·限流 优化·弹性·隔离 压测·容量规划·预案
图 2:MQ 积压处理五步法——发现 → 定位 → 止血 → 根治 → 预防。

① 监控发现:怎么知道积压了?

看 consumer lag、消费延迟、生产/消费 TPS 曲线。设置告警阈值(如 lag > 1 万且持续增长自动报警),别等用户投诉才发现。

② 定位原因:到底卡在哪?

看消费单条耗时(是不是慢 SQL / 外部调用);看消费 TPS(是不是实例不够);看分区数(是不是并行度受限);看错误日志(是不是重试风暴);看 Broker 指标(是不是磁盘 / 网络)。先量再猜

③ 紧急止血:先让 lag 回落

按第 4 节「止血」操作:临时扩容消费者 + 多线程消费 + 提分区数;非核心消息降级跳过;必要时限流生产者。目标是先止住上涨、开始下降

④ 根治优化:把根因拔掉

优化消费逻辑(异步 / 批处理 / 去慢 SQL);消费端接弹性伸缩(按 lag 自动扩缩);核心 / 非核心消息隔离;接入限流削峰;保证幂等。让消费能力长期 > 峰值生产。

⑤ 预防演练:让下次不再发生

做容量规划与压测,知道系统上限在哪;准备大促 / 秒杀的扩容与降级预案;把「积压应急手册」固化成 runbook,定期演练。

6 监控、告警与预防 / Prevent

最好的处理是「让它不发生」。把下面几件事做在前面,积压基本不会成为故障。

预防手段具体做法解决哪类问题
实时监控 lagKafka exporter / RocketMQ dashboard 接 Prometheus + Grafana,盯 consumer lag发现慢、定位快
分级告警lag 超阈值(如 1 万)warning,超更大值(如 10 万)page 值班早发现早处理
弹性伸缩消费端按 lag 自动扩缩容(K8s HPA / 自定义指标)消费者不足
容量规划 + 压测提前测出单消费者吞吐上限,按比例预留实例与分区突发流量
限流削峰生产端令牌桶 / 漏桶,洪峰平滑后入队;非核心消息走降级通道突发流量
消息分级隔离核心 / 非核心分 topic、分队列,互不影响重试风暴 / 局部堵塞
幂等 + 死信队列重复消费安全;失败消息进 DLQ 不阻塞主链路消费失败
口诀收尾: 积压不可怕,可怕的是「没监控、没预案、没扩容」。记住——发现靠监控 止血靠扩容降级 根治靠优化+弹性 预防靠压测+预案