Skip to content

消息堆积与背压

消息堆积是 MQ 生产排障里非常高频的问题。

一句话理解:

消息堆积就是生产者写入消息的速度长期大于消费者处理消息的速度,Broker 中等待处理的消息越来越多。

背压可以理解为:

当下游处理不过来时,系统不要无限接收和堆内存排队,而是通过限流、暂停拉取、降低生产速度、延迟重试等方式,把压力向上游传递或吸收。

很多人遇到堆积只会说“加消费者”,但真实生产里这远远不够。因为堆积可能来自消费者慢、下游数据库慢、重试消息爆炸、顺序消息阻塞、分区太少、Broker 磁盘压力、ACK 配置错误、代码线程池无界排队等不同原因。

学习目标

目标你需要掌握什么
知道是什么理解堆积、积压、Lag、Ready、Unacked、重试堆积、背压的区别
知道为什么理解生产速度、消费速度、下游能力和 Broker 存储之间的关系
知道怎么工作会从生产、Broker、消费、下游、重试、顺序消息几个方向定位
知道不这样会怎样知道盲目扩容、盲目清消息、提前 ACK、无界线程池的风险
会排查能用 Kafka、RabbitMQ、RocketMQ 的指标判断堆积位置
会项目落地能设计削峰、限流、幂等、重试、死信、补偿和容量预案
会复盘能把“现象、原因、处理、预防”沉淀成稳定排障方案

堆积、积压、Lag、背压的区别

名词含义常见产品里的表现
消息堆积Broker 中待消费消息越来越多Queue depth 增大、Topic 消息滞留
消息积压和堆积基本同义,更强调未处理存量等待处理的消息量持续增长
Lag消费进度落后生产进度多少Kafka Consumer Lag、RocketMQ Diff
ReadyRabbitMQ 中已入队但还没投递给消费者的消息队列 Ready 数很高
UnackedRabbitMQ 中已投递但消费者未 ACK 的消息消费者拿走了但没确认
Retry 堆积消费失败后不断进入重试队列重试 Topic / 死信队列增长
背压下游慢时限制上游或暂停拉取,防止系统被压垮pause/resume、限流、熔断、降级

堆积和背压不是一个概念。堆积是现象,背压是治理手段。

如果面试或排障问到“Lag 分布是什么”,不要只看总 Lag。要继续看每个 Kafka Partition、RocketMQ MessageQueue 或 RabbitMQ Queue 分别积压多少,判断是整体消费慢还是少数通道热点、慢消息、顺序阻塞。详细看:Lag 分布与堆积定位

mermaid
flowchart TD
    A["生产者持续写入"] --> B["Broker 存储消息"]
    B --> C["消费者拉取或接收消息"]
    C --> D["执行业务处理"]
    D --> E["数据库 / 远程接口 / ES / 缓存"]
    E --> F{"下游是否处理得过来"}
    F -- "处理得过来" --> G["ACK 或提交 Offset"]
    F -- "处理不过来" --> H["消费速度下降"]
    H --> I["Broker 消息堆积"]
    H --> J["触发背压或限流"]

堆积的核心公式

先记住一个非常有用的公式:

text
堆积增长速度 = 生产 TPS - 消费 TPS

如果:

text
生产速度 = 5000 条/秒
消费速度 = 3000 条/秒

那么:

text
每秒新增堆积 = 2000 条

10 分钟后大约堆积:

text
2000 * 60 * 10 = 120 万条

如果后续生产恢复到 1000 条/秒,消费速度提升到 5000 条/秒,那么消化速度是:

text
净消化速度 = 5000 - 1000 = 4000 条/秒

消化 120 万条需要:

text
1200000 / 4000 = 300 秒

也就是 5 分钟左右。

这个公式能帮助你在生产排障时快速判断:不是看到堆积就慌,而是要判断生产速度、消费速度、存量、磁盘容量和预计消化时间。

容量评估模板

生产排障时不要只说“堆积了 100 万条”。100 万条日志消息和 100 万条支付回调消息,风险完全不同。真正要评估的是:还在不在增长、多久能追平、最老消息是否已经超过业务 SLA、Broker 是否撑得到消化完成。

必须收集的指标

指标说明为什么重要
当前堆积量Broker 中还没处理的消息数量判断存量规模
生产 TPS每秒新增多少消息判断压力源是否还在继续
成功消费 TPS每秒真正业务成功并 ACK 的消息数量判断系统真实处理能力
失败 TPS每秒失败、重试、进死信的消息数量判断是否有重试风暴
最老消息等待时间最早未处理消息已经等了多久比总 Lag 更能反映业务 SLA 风险
平均处理耗时 / P95 耗时单条消息业务处理耗时判断消费者是否被下游拖慢
分区 / 队列 Lag 分布每个并行通道分别堆积多少判断整体慢还是局部热点
Broker 磁盘水位消息文件还能存多久判断是否会先把 Broker 打满
下游容量DB、ES、Redis、HTTP 接口还能承受多少 QPS判断扩容消费者会不会把下游打挂

追平时间怎么算

text
净消化速度 = 成功消费 TPS - 当前生产 TPS
预计追平时间 = 当前堆积量 / 净消化速度

如果 成功消费 TPS <= 当前生产 TPS,说明系统还在继续落后,理论上永远追不平。此时优先级不是“消化存量”,而是先止血:暂停补数据、限制生产速度、关闭非核心任务、隔离失败消息、保护下游。

例如:

text
当前堆积量 = 1800000
生产 TPS = 2000
成功消费 TPS = 5000
净消化速度 = 3000
预计追平时间 = 1800000 / 3000 = 600 秒

也就是大约 10 分钟。如果最老消息 SLA 是 30 分钟,而且 Broker 磁盘还能撑 2 小时,这通常是可控堆积。如果预计追平时间是 4 小时,但消息保留窗口只有 2 小时,或者最老消息已经超过支付、库存、通知业务 SLA,就必须升级处理。

mermaid
flowchart TD
    A["发现堆积"] --> B["收集生产 TPS"]
    B --> C["收集成功消费 TPS"]
    C --> D["计算净消化速度"]
    D --> E{"净消化速度是否大于 0"}
    E -- "否" --> F["先限流止血"]
    F --> G["隔离失败消息和保护下游"]
    E -- "是" --> H["估算追平时间"]
    H --> I{"是否超过业务 SLA 或存储窗口"}
    I -- "是" --> J["扩容、批处理、降级或补偿"]
    I -- "否" --> K["持续观察并准备预案"]

堆积处理决策表

消息类型能不能丢推荐处理错误做法
支付、退款、库存扣减不能丢限流止血、保证幂等、失败进补偿表、人工核对清队列、跳 offset
订单状态同步 ES通常不能直接丢,但可从 MySQL 重建暂停消费优化 ES、按订单更新时间补偿重建索引不记录范围直接跳过
短信、站内信、营销通知取决于业务时效超时可降级、延迟发送或按策略丢弃不经业务确认全部重发
日志、埋点、监控明细多数可按策略丢弃或采样降采样、缩短保留、离线补算核心报表为了日志把主业务拖垮
医疗采集原始数据通常不能丢本地落盘、断点续传、补采、对账只依赖内存队列
缓存刷新消息可以重建删除缓存或按主库批量刷新把缓存消息当核心交易处理

判断能不能跳过消息时,核心不是技术上能不能 reset offset,而是业务上有没有权威数据源可以恢复。如果消息是唯一事实来源,不能丢;如果 MySQL 主库、业务流水、binlog、采集原始文件能完整重建,可以在确认范围后跳过,再用补偿任务恢复。

正常堆积和异常堆积

不是所有堆积都是故障。MQ 本来就有削峰填谷能力。

类型特点是否问题
计划内堆积活动高峰期间短暂堆积,随后能按预期消化可接受
异常堆积堆积持续增长,消费速度追不上,预计消化时间越来越长需要处理
重试堆积同一批失败消息反复重试,占满消费能力高风险
顺序阻塞某个队列或分区被一条慢消息卡住高风险
Broker 存储告警磁盘水位高、文件保留压力大高风险

判断是否异常,关键看三点:

  1. 堆积量是否持续增长。
  2. 最老消息延迟是否超过业务 SLA。
  3. 当前消费能力是否能在可接受时间内追平。

堆积原因总览

mermaid
flowchart TD
    A["消息堆积"] --> B["生产速度突然升高"]
    A --> C["消费者实例少或并发低"]
    A --> D["下游数据库或接口慢"]
    A --> E["消费失败反复重试"]
    A --> F["顺序消息被单条消息阻塞"]
    A --> G["分区或队列数量不足"]
    A --> H["Broker 磁盘或网络压力"]
    A --> I["ACK 或 Offset 提交不合理"]
    A --> J["消费者代码阻塞或线程池打满"]

常见原因拆开看:

原因典型表现处理方向
生产突增活动、补数据、重放任务导致写入暴涨限流、削峰、批量写入、容量预案
消费者少消费 TPS 低,CPU/线程未充分利用增加实例、提高并发、批量消费
分区或队列少加消费者后吞吐不提升增加分区/队列,重新规划路由
下游慢消费线程卡在 DB、HTTP、ES优化 SQL、加索引、批处理、熔断降级
重试风暴消息失败后不断重投,正常消息也被拖慢失败隔离、延迟重试、死信队列
顺序阻塞某个队列一直不动,其他队列正常找到卡住消息,局部隔离或人工处理
ACK 错误RabbitMQ Unacked 很高,或消息重复很多手动 ACK、成功后确认、合理 prefetch
线程池无界JVM 内存上涨,队列越来越长有界队列、拒绝策略、背压
Broker 压力磁盘、网络、Page Cache、GC 异常扩容 Broker、清理历史、调整保留策略

消费全过程:堆积到底卡在哪里

很多人排查 MQ 堆积时只看“消费者数量”和“Lag 数字”,但不知道一条消息进入消费者之后经历了哪些步骤。真正定位时,要把消息消费链路拆开。

mermaid
flowchart TD
    A["Broker 中有待消费消息"] --> B["消费者拉取或接收消息"]
    B --> C["反序列化和参数校验"]
    C --> D["提交到本地线程池"]
    D --> E["执行业务逻辑"]
    E --> F["访问 DB / ES / Redis / HTTP"]
    F --> G["业务成功落库或写入下游"]
    G --> H["ACK 或提交 Offset"]
    H --> I["消费进度推进"]

每一步都可能变成堆积点。

阶段卡住的表现指标或证据常见处理
Broker 投递拉取慢、吞吐低Broker 网络、磁盘、请求耗时扩容 Broker、检查磁盘水位和网络
拉取到消费者消费者数量不少但 TPS 低poll 耗时、fetch 速率、连接状态检查消费者配置和网络
反序列化CPU 高、异常多反序列化异常日志、CPU profile修复消息格式、优化序列化
本地线程池拉得快但处理慢队列长度、活跃线程、拒绝次数有界队列、限流、pause/resume
业务逻辑单条耗时高P95/P99、方法耗时日志优化代码、批处理、减少锁
下游依赖线程阻塞等待慢 SQL、连接池满、HTTP 超时优化下游、降并发、熔断
ACK/Offset重复消费或 Unacked 高ACK 失败、commit 失败、Unacked成功后确认、合理提交策略

为什么“拉得快”不等于“处理得快”

有些消费者会先从 Broker 拉一批消息,再丢到本地线程池处理。如果本地线程池是无界队列,就可能出现一种假象:

  1. Broker Lag 短时间下降。
  2. 消费者 JVM 内部队列越来越长。
  3. 业务真正处理速度没变。
  4. 内存上涨、延迟上涨,最后 OOM 或任务超时。
mermaid
flowchart TD
    A["消费者快速拉取消息"] --> B["提交到无界队列"]
    B --> C["Broker Lag 表面下降"]
    B --> D["JVM 内存和本地队列上涨"]
    D --> E["业务处理仍然慢"]
    E --> F["最终 OOM 或超时"]

所以判断消费能力时,不能只看 Broker Lag 是否下降,还要看业务成功 TPS。真正有效的消费是:业务处理成功,并且 ACK 或 Offset 提交成功。

背压不是简单少拉一点

背压的目标不是“让消费者变慢”,而是让系统在下游变慢时仍然可控。

没有背压时,典型过程是:

mermaid
flowchart TD
    A["消费者持续从 Broker 拉消息"] --> B["提交到本地线程池"]
    B --> C["下游 DB / ES / HTTP 变慢"]
    C --> D["线程处理不过来"]
    D --> E["本地队列越来越长"]
    E --> F["内存上涨、GC 变频繁"]
    F --> G["处理更慢"]
    G --> H["超时、重试、堆积扩大"]

这是一种正反馈:越慢越堆,越堆越慢。背压要打断这个正反馈。

背压要控制哪几个闸门

闸门控制什么如果不控制会怎样
Broker 拉取量每次拿多少消息、是否暂停拉取消息从 Broker 转移到消费者内存
本地线程池队列本机最多排队多少任务无界队列导致 OOM 和长延迟
下游并发同时打 DB/ES/HTTP 多少请求把下游打满,引发全链路超时
失败重试速度失败消息多久重试重试风暴吞掉正常消费能力
上游生产速度是否暂停补数据、限流入口消费端永远追不上

一个可落地的背压策略

mermaid
flowchart TD
    A["消费者运行中"] --> B["采集本地队列长度、活跃线程、下游耗时"]
    B --> C{"是否超过高水位"}
    C -- "是" --> D["暂停拉取或降低 prefetch"]
    D --> E["继续处理已拿到的消息"]
    E --> F{"是否低于低水位"}
    F -- "否" --> E
    F -- "是" --> G["恢复拉取"]
    C -- "否" --> H["正常拉取和处理"]

为什么要有高水位和低水位两个阈值?因为如果只用一个阈值,消费者可能在“暂停、恢复、暂停、恢复”之间频繁抖动。高低水位可以形成缓冲区。

text
高水位:队列使用率超过 80%,暂停拉取
低水位:队列使用率低于 40%,恢复拉取

有界线程池为什么是必须的

消费端本地线程池不要使用无界队列。无界队列会让你误以为“没有拒绝,系统还能接”,实际上只是把消息堆在 JVM 里。

java
ThreadPoolExecutor executor = new ThreadPoolExecutor(
        16,
        16,
        60,
        TimeUnit.SECONDS,
        new ArrayBlockingQueue<>(2000),
        new ThreadFactory() {
            private final AtomicInteger counter = new AtomicInteger(1);

            @Override
            public Thread newThread(Runnable runnable) {
                return new Thread(runnable, "mq-consumer-" + counter.getAndIncrement());
            }
        },
        new ThreadPoolExecutor.AbortPolicy()
);

这个配置的含义:

参数含义为什么这样
corePoolSize = 16固定 16 个工作线程消费能力可预估,避免无限创建线程
maximumPoolSize = 16不临时放大线程下游慢时盲目加线程只会扩大压力
ArrayBlockingQueue<>(2000)最多本地排队 2000 条超过就触发背压,而不是无限占内存
AbortPolicy队列满时拒绝提交让上层暂停拉取或降速,而不是悄悄堆积

如果是 IO 型消费,可以适当提高线程数;如果下游 DB 已经慢,提高线程数通常只会增加连接池等待和锁竞争。

Offset 提交为什么比想象中难

Kafka 场景里,如果消费者把消息提交到线程池后立刻提交 offset,会产生“假消费成功”:

mermaid
flowchart TD
    A["poll 拉到消息"] --> B["提交到线程池"]
    B --> C["立即 commit offset"]
    C --> D["业务线程稍后执行"]
    D --> E{"业务是否成功"}
    E -- "失败" --> F["Kafka 已认为成功<br/>消息不会自动重来"]

正确原则是:每个分区只提交已经连续成功处理到的位置

例如一个分区拉到 offset 100、101、102、103:

offset处理结果能否提交到这里
100成功可以提交 101
101成功可以提交 102
102失败不能越过 102
103成功仍不能提交 104

为什么 103 成功也不能提交 104?因为 offset 102 失败了。如果直接提交 104,102 就被跳过,造成消息丢失。并发消费时必须记录每个分区的成功连续位点,不能只记录最大完成 offset。

简化版思路:

java
public class PartitionOffsetTracker {
    private final Map<TopicPartition, Long> nextCommitOffset = new HashMap<>();
    private final Map<TopicPartition, SortedSet<Long>> finishedOffsets = new HashMap<>();

    public synchronized void markSuccess(TopicPartition tp, long offset) {
        finishedOffsets.computeIfAbsent(tp, k -> new TreeSet<>()).add(offset);
        long next = nextCommitOffset.getOrDefault(tp, offset);

        SortedSet<Long> finished = finishedOffsets.get(tp);
        while (finished.remove(next)) {
            next++;
        }
        nextCommitOffset.put(tp, next);
    }

    public synchronized Map<TopicPartition, OffsetAndMetadata> committableOffsets() {
        Map<TopicPartition, OffsetAndMetadata> result = new HashMap<>();
        for (Map.Entry<TopicPartition, Long> entry : nextCommitOffset.entrySet()) {
            result.put(entry.getKey(), new OffsetAndMetadata(entry.getValue()));
        }
        return result;
    }
}

这段代码只是演示原理。生产实现还要处理初始 offset、失败重试、分区 revoke、应用关闭、异常消息死信等细节。核心思想是:提交 offset 不是提交“已经拉到哪里”,而是提交“已经安全处理到哪里”。

毒消息为什么会拖垮整个消费组

毒消息是指某条消息因为数据格式、业务状态、外部依赖或代码 bug,导致每次消费都失败。它最危险的地方不是“这一条失败”,而是它会持续占用消费能力。

mermaid
flowchart TD
    A["消费者拿到毒消息"] --> B["业务处理失败"]
    B --> C["立即重试"]
    C --> D["再次失败"]
    D --> E["消费者反复处理同一批失败消息"]
    E --> F["正常消息排队等待"]
    F --> G["Lag 继续增长"]

处理原则:

做法说明
区分异常类型网络超时可重试,参数错误不要无限重试
延迟重试给下游恢复时间,避免快速打满消费者
最大重试次数超过次数进入死信或异常表
死信告警死信不是垃圾桶,必须有人看
修复后重放修复数据或代码后,从死信表按业务规则补偿

商业系统里,支付、库存、订单状态类消息进入死信后,不能只记录日志。要能按业务号查询事实源,确认是否需要补偿、撤销、人工审核或重放。

ACK 和 Offset 为什么影响堆积

不同 MQ 叫法不一样,但本质都是告诉 Broker:“这条消息我处理完了,可以推进进度。”

产品确认机制提前确认风险过晚确认风险
RabbitMQACK业务失败但消息已删除Unacked 高、重复投递、内存压力
KafkaCommit Offset业务失败但 offset 已推进Rebalance 后重复消费
RocketMQ返回消费状态业务失败但返回成功重试堆积、顺序阻塞

正确原则:

业务真正成功后再确认;确认失败或网络异常时,消费端要允许重复消息,并靠幂等保护业务。

生产 TPS、拉取 TPS、成功 TPS 的区别

排查堆积时要特别区分三个 TPS。

指标含义为什么重要
生产 TPS上游每秒写入 Broker 的消息数判断压力源
拉取 TPS消费者每秒从 Broker 取到的消息数判断 Broker 到消费者通不通
成功 TPS业务成功并 ACK/提交 Offset 的消息数判断真实消化能力

最容易误判的是拉取 TPS。比如消费者每秒拉 5000 条,但业务每秒只成功 1000 条,其余都在本地线程池排队或失败重试。此时真实消化能力是 1000,不是 5000。

text
有效消化速度 = 成功 TPS - 生产 TPS

如果成功 TPS 小于生产 TPS,即使消费者日志看起来“拉了很多”,堆积也一定会扩大。

堆积形成全过程

下面用订单同步 ES 举例,看堆积是怎么一步步形成的。

mermaid
flowchart TD
    A["订单服务持续发送订单变更消息"] --> B["Kafka Broker 写入消息"]
    B --> C["ES 同步消费者拉取消息"]
    C --> D["消费者构建搜索文档"]
    D --> E["调用 ES upsert"]
    E --> F{"ES 写入是否变慢"}
    F -- "否" --> G["提交 Offset"]
    F -- "是" --> H["单条耗时升高"]
    H --> I["消费者线程被占满"]
    I --> J["成功 TPS 下降"]
    J --> K["Lag 增长"]
    K --> L["最老消息等待时间升高"]

这里真正的瓶颈不是 Kafka,而是 ES 写入变慢。此时如果盲目扩容消费者,可能让更多请求打到 ES,导致 ES rejected、写入延迟更高,最终堆积更严重。

堆积排查证据链

生产排查要形成证据链,而不是凭感觉说“应该是消费者慢”。

第一层:确认堆积规模

要看什么例子
Topic / Queueorder-event
ConsumerGrouporder-search-sync-group
总 Lag / Queue Depth120 万
最大单分区 LagPartition 3 有 80 万
最老消息等待时间42 分钟

第二层:确认速度关系

要看什么例子判断
生产 TPS3000/s压力是否还在
拉取 TPS5000/s是否能从 Broker 拿到
成功 TPS1200/s真实消化能力不足
失败 TPS800/s重试正在吞吐能力

如果拉取 TPS 高但成功 TPS 低,瓶颈一定在消费者内部或下游,而不是 Broker 没投递。

第三层:定位消费者内部

要看什么说明
消费线程活跃数是否线程全忙
本地队列长度是否内部排队
单条处理 P95/P99是否长尾严重
GC 和内存是否本地堆积导致内存上涨
错误日志是否大量失败消息反复重试

第四层:定位下游

下游证据
MySQL慢 SQL、锁等待、连接池 active 打满
ESbulk rejected、refresh/merge 压力、写入延迟
Redis慢命令、连接池耗尽、网络延迟
HTTP 接口超时、429、熔断、连接池等待

第五层:确认处理是否生效

处理后继续观察:

  1. 生产 TPS 是否下降或恢复正常。
  2. 成功 TPS 是否持续大于生产 TPS。
  3. 总 Lag 是否下降。
  4. 最大分区 Lag 是否下降。
  5. 最老消息等待时间是否下降。
  6. 失败率和重试量是否下降。
  7. Broker 磁盘水位是否稳定。

如果只看到总 Lag 下降,但最老消息等待时间不降,可能还有慢消息、顺序阻塞或局部分区卡住。

标准排查流程

排查堆积不要一上来就重启消费者。先回答五个问题:

  1. 哪个 Topic / Queue / ConsumerGroup 堆积?
  2. 当前生产 TPS 和消费 TPS 分别是多少?
  3. 是所有分区/队列都堆积,还是某几个特别严重?
  4. 消费失败率、重试量、死信量是否增加?
  5. 消费者卡在哪里:CPU、线程池、数据库、远程接口、锁等待、GC?
mermaid
flowchart TD
    A["发现消息堆积"] --> B["确认 Topic / Queue / Group"]
    B --> C["看生产 TPS 和消费 TPS"]
    C --> D{"堆积是否持续增长"}
    D -- "否" --> E["评估预计消化时间"]
    D -- "是" --> F["定位消费变慢原因"]
    F --> G["看消费者错误率和重试"]
    F --> H["看下游 DB / HTTP / ES"]
    F --> I["看线程池 / CPU / GC"]
    F --> J["看分区队列是否热点"]
    G --> K["隔离异常消息或进死信"]
    H --> L["优化下游或限流"]
    I --> M["扩容或调并发"]
    J --> N["调整 key / 分区 / 队列"]

Kafka 怎么看堆积

Kafka 中常用 Lag 表示消费落后量。

常用命令:

bash
kafka-consumer-groups.sh \
  --bootstrap-server 127.0.0.1:9092 \
  --describe \
  --group order-search-sync-group

典型输出里关注:

text
TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
order-event     0          1000            1500            500
order-event     1          2000            8000            6000
字段含义
CURRENT-OFFSET当前消费者提交到哪里
LOG-END-OFFSET分区最新消息位置
LAG还有多少消息没消费

如果只有某个分区 Lag 特别高,常见原因是:

  1. key 分布不均导致热点分区。
  2. 某条消息处理特别慢。
  3. 该分区对应消费者实例异常。
  4. 分区内消息顺序处理,无法被多个消费者并行处理。

Kafka 要特别注意:

同一个 Consumer Group 中,一个分区同一时刻只能被一个消费者消费。消费者实例数超过分区数,多出来的实例不会提升吞吐。

mermaid
flowchart TD
    A["Topic 有 3 个 Partition"] --> B["最多 3 个消费者并行消费"]
    B --> C{"启动 6 个消费者有用吗"}
    C -- "没有完全有用" --> D["只有 3 个消费者能分到分区"]
    C -- "想提升吞吐" --> E["增加分区或提高单消费者处理能力"]

RabbitMQ 怎么看堆积

RabbitMQ 重点看队列里的几个指标:

指标含义判断
Ready已在队列中,等待投递给消费者Ready 高说明消费者拿得慢或数量不足
Unacked已投递给消费者,但还没 ACKUnacked 高说明消费者处理慢或 ACK 卡住
Publish rate生产速率判断入口压力
Deliver / Ack rate投递和确认速率判断消费能力
Consumers消费者数量判断是否掉线或不足

典型情况:

现象可能原因
Ready 很高,Unacked 不高消费者数量不足、prefetch 太小、消费者没有正常拉取
Unacked 很高消费者拿了消息但处理慢,或者没 ACK
Ready 和 Unacked 都高生产过快,消费端整体处理不过来
消费者数为 0消费者服务挂了或连接失败

RabbitMQ 中 prefetch 很关键。它控制一个消费者最多同时拿多少未确认消息。

java
channel.basicQos(50);

如果 prefetch 太大,消费者会一次拿走很多消息,导致 Unacked 很高,其他消费者分不到消息;如果太小,吞吐可能上不去。

RocketMQ 怎么看堆积

RocketMQ 常看 ConsumerGroup 的消费进度和 Diff。

常用排查方向:

指标含义
Topic 队列堆积某个 Topic 下消息等待消费
ConsumerGroup Diff消费进度落后多少
Retry Topic消费失败进入重试的消息
DLQ多次失败后的死信消息
消费 TPS消费者实际吞吐
Broker 磁盘CommitLog、ConsumeQueue 存储压力

RocketMQ 顺序消息要特别注意:如果某条消息一直处理失败,它所在队列的后续消息可能被阻塞,导致局部堆积。

mermaid
flowchart TD
    A["队列 Q0"] --> B["消息 1 成功"]
    B --> C["消息 2 失败或超慢"]
    C --> D["消息 3 等待"]
    D --> E["消息 4 等待"]
    C --> F["Q0 局部堆积"]

处理顺序消息堆积时,不能简单跳过消息。要看业务是否允许:

  1. 修复数据后重试。
  2. 把异常消息转人工处理。
  3. 将失败原因落库,后续补偿。
  4. 调整业务 key,减少热点队列。

堆积时先做什么

生产故障中,建议按这个顺序处理:

1. 先止血

目标是让堆积不再继续恶化。

常见动作:

  1. 临时关闭非核心生产入口。
  2. 对上游接口限流。
  3. 暂停补数据、重放、批量导入任务。
  4. 降级非核心消费者逻辑。
  5. 对明显失败的消息进入死信或异常表,不要无限快速重试。

2. 再定位瓶颈

判断消费者慢在哪里:

瓶颈现象处理
CPU 高计算、序列化、压缩、复杂规则慢优化代码或扩容
DB 慢慢 SQL、锁等待、连接池满加索引、批处理、扩容、限流
远程接口慢HTTP 调用耗时高超时、熔断、异步补偿
线程池满活跃线程满,队列增长有界队列、调并发、背压
单分区热点某分区 Lag 特别高调整 key、增加分区、拆 Topic
重试过多失败消息反复消费延迟重试、死信、人工处理

3. 最后消化存量

消化堆积时要计算预计时间:

text
预计消化时间 = 当前堆积量 / (当前消费 TPS - 当前生产 TPS)

如果消费 TPS 小于生产 TPS,说明永远追不上,必须先限流或扩容。

消化手段:

  1. 临时增加消费者实例。
  2. 提高单消费者批处理能力。
  3. 临时放宽非核心业务校验。
  4. 将异常消息隔离到死信队列。
  5. 对历史低价值消息按业务规则降级处理。
  6. 增加分区或队列,但要评估顺序和路由影响。

扩容消费者短暂有效,后面又堆积

这一节先给出排查总览。完整知识点拆到了独立页面:消费者扩容后再次堆积

现象通常是:

  1. 刚增加消费者实例后,Lag 或队列深度开始下降。
  2. 过一会儿下降速度变慢。
  3. 随后堆积又开始增长。

这说明:扩容消费者只短暂提高了消费入口能力,但真正瓶颈没有消失,或者新的瓶颈被打出来了。

mermaid
flowchart TD
    A["消息堆积"] --> B["扩容消费者"]
    B --> C["短时间消费 TPS 上升"]
    C --> D{"是否超过生产 TPS 且下游扛得住"}
    D -- "是" --> E["堆积逐步消化"]
    D -- "否" --> F["堆积再次增长"]
    F --> G["下游 DB / ES / HTTP 变慢"]
    F --> H["分区或队列限制并行度"]
    F --> I["重试消息占满消费能力"]
    F --> J["线程池 / 连接池 / Broker 到达瓶颈"]

常见原因:

原因为什么会短暂有效后又堆积
下游数据库扛不住消费者变多后并发写 DB 增加,DB 慢 SQL、锁等待、连接池打满,单条消费耗时变长
远程接口限流或变慢刚开始并发提升,随后下游接口触发限流、熔断或排队
分区/队列数不足Kafka 分区数、RocketMQ 队列数、RabbitMQ 队列并行度限制了真正并发,消费者加多后边际收益很低
热点分区或热点队列某个 key 的消息集中到一个分区/队列,整体扩容无法解决单分区串行瓶颈
重试风暴扩容后失败消息被更快取出并更快重试,正常消息反而被失败消息挤占
消费者线程池打满实例增加后本机或下游线程池、连接池、内存、CPU 被打满,处理耗时上升
Broker 压力变大更多消费者同时拉取、提交 offset、ACK、网络传输,Broker 磁盘或网络成为瓶颈
生产速度也在上涨扩容期间上游流量继续增大,消费 TPS 仍低于生产 TPS
ACK 或 Offset 策略不当过早提交导致失败补偿压力,过晚提交导致重复消费和重平衡后重放
Rebalance 频繁Kafka 扩容消费者会触发 Rebalance,短时间内分区重新分配,消费会暂停或抖动

怎么判断是哪种原因

不要只看消费者数量,要看扩容前后这些指标是否变化:

指标如果出现这种变化说明
消费 TPS 先升后降扩容初期有效,后面被新瓶颈限制看 DB、接口、线程池、Broker
单条消费耗时升高消费者并发把下游打慢了下游是瓶颈
DB 连接池 active 打满每个消费者都在等连接连接池或数据库瓶颈
慢 SQL、锁等待增加并发写入造成数据库竞争优化 SQL、索引、批处理、锁范围
Kafka 某个分区 Lag 高热点分区或单分区消费者慢调 key、增分区、拆 Topic
RabbitMQ Unacked 高消息已投递但消费者处理慢看业务耗时、prefetch、ACK
重试队列增长失败消息在吞噬消费能力死信隔离、修复异常数据
CPU 不高但 Lag 高可能卡在 IO、锁、外部接口看线程栈和下游耗时
CPU 很高序列化、压缩、规则计算或日志过多优化代码或扩容 CPU

排查流程:

mermaid
flowchart TD
    A["扩容后再次堆积"] --> B["比较生产 TPS 和消费 TPS"]
    B --> C{"消费 TPS 是否仍低于生产 TPS"}
    C -- "是" --> D["继续定位消费瓶颈"]
    C -- "否" --> E["计算预计消化时间"]
    D --> F["看单条消费耗时"]
    F --> G{"耗时是否变高"}
    G -- "是" --> H["查 DB / HTTP / ES / 连接池"]
    G -- "否" --> I["查分区队列并行度和热点"]
    I --> J["查重试队列和死信"]
    H --> K["限流、批处理、优化下游"]
    J --> L["隔离异常消息或调整路由"]

应该怎么处理

处理思路不是“继续无限加消费者”,而是找到真正瓶颈。

真正原因处理方式
下游 DB 慢优化 SQL、加索引、批量写入、减少事务范围、限流消费者
下游接口慢设置超时、熔断、降级、异步补偿、降低并发
分区/队列不足增加分区或队列,重新规划 key 和路由
热点 key拆分热点 key、加随机后缀、按更细粒度路由,但要评估顺序要求
重试风暴区分可恢复和不可恢复异常,延迟重试,失败进入死信或异常表
本机线程池堆积使用有界队列、CallerRunsPolicy、暂停拉取,形成背压
Broker 压力扩容 Broker、调整保留策略、检查磁盘水位和网络
上游持续过快对生产者限流、削峰、暂停补数据任务

不能做什么

这些操作很危险:

错误操作后果
直接清空队列可能丢业务数据,后续无法补偿
Kafka 直接 reset offset 到 latest老消息全部跳过,可能造成订单、库存、索引不同步
消费端提前 ACK业务失败后消息已删除,只能人工补
盲目无限加线程打爆数据库、Redis、ES 或远程服务
无界线程池排队JVM 内存上涨,最终 OOM
失败消息立即无限重试重试风暴,正常消息也处理不了
分片/队列随便加很多协调成本、存储成本、顺序语义都可能出问题

生产处理时必须记住:不能一上来清消息或重置 offset,必须先确认业务是否允许丢弃、是否有补偿来源、是否能从主库重建。

背压怎么设计

背压的核心思想是:消费者处理不过来时,不要把压力无限堆到本机内存里。

mermaid
flowchart TD
    A["消费者拉取消息"] --> B{"本地处理队列是否接近满"}
    B -- "否" --> C["继续拉取并处理"]
    B -- "是" --> D["暂停拉取或降低拉取量"]
    D --> E["让消息留在 Broker"]
    E --> F["等待本地处理能力恢复"]
    F --> C

常见背压手段:

层级手段
生产者限流、削峰、降级、批量发送、熔断非核心入口
Broker队列容量、磁盘水位、保留策略、延迟重试、死信
消费者拉取Kafka pause/resume、RabbitMQ prefetch、RocketMQ 消费线程配置
消费者执行有界线程池、信号量、批处理、超时控制
下游资源DB 连接池保护、慢 SQL 优化、接口熔断、限流

Java Demo:有界线程池防止本机堆积

错误做法是使用无界队列:

java
ExecutorService executor = Executors.newFixedThreadPool(20);

newFixedThreadPool 底层使用无界队列。消费速度跟不上时,消息任务会不断堆到 JVM 内存里,最后可能 OOM。

更推荐使用有界队列:

java
ThreadPoolExecutor executor = new ThreadPoolExecutor(
    20,
    20,
    60,
    TimeUnit.SECONDS,
    new ArrayBlockingQueue<>(1000),
    new ThreadPoolExecutor.CallerRunsPolicy()
);

CallerRunsPolicy 的含义是:线程池和队列都满了,就让提交任务的线程自己执行任务。这样拉取线程会变慢,形成一种简单背压,不会无限把任务塞进内存。

Kafka Demo:处理不过来时暂停拉取

Kafka 消费者可以在本地队列满时暂停拉取,处理能力恢复后再继续。

java
public class BackpressureKafkaConsumer {
    private final KafkaConsumer<String, String> consumer;
    private final ThreadPoolExecutor executor;

    public BackpressureKafkaConsumer(KafkaConsumer<String, String> consumer,
                                     ThreadPoolExecutor executor) {
        this.consumer = consumer;
        this.executor = executor;
    }

    public void pollLoop() {
        while (true) {
            if (executor.getQueue().remainingCapacity() < 100) {
                consumer.pause(consumer.assignment());
            } else {
                consumer.resume(consumer.assignment());
            }

            ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
            for (ConsumerRecord<String, String> record : records) {
                executor.submit(() -> handle(record));
            }
        }
    }

    private void handle(ConsumerRecord<String, String> record) {
        // 业务处理:写数据库、同步 ES、调用下游等
    }
}

真实生产还要注意 offset 提交时机。不要任务刚提交到线程池就提交 offset,否则线程池里的业务失败时,Kafka 已经认为消息处理完了。

更稳的做法是:

  1. 控制每次 poll 的数量。
  2. 业务真正成功后再提交 offset。
  3. 批量处理时记录每个分区成功处理到的最大 offset。
  4. 失败消息进入延迟重试或死信,不要阻塞整个消费组。

RabbitMQ Demo:手动 ACK

RabbitMQ 消费端应该在业务成功后再 ACK。

java
public void handle(Message message, Channel channel) throws IOException {
    long deliveryTag = message.getMessageProperties().getDeliveryTag();
    try {
        OrderEvent event = parse(message);
        orderService.handle(event);

        channel.basicAck(deliveryTag, false);
    } catch (RecoverableException e) {
        channel.basicNack(deliveryTag, false, true);
    } catch (Exception e) {
        channel.basicNack(deliveryTag, false, false);
    }
}

含义:

操作含义
basicAck确认成功,Broker 可以删除消息
basicNack(..., true)失败并重新入队,适合短暂故障
basicNack(..., false)不重新入队,可进入死信队列

不要对永久失败的消息无限 requeue,否则它会反复投递,造成重试风暴。

幂等是消化堆积的前提

堆积后常常需要重试、扩容、重放、补偿。只要有这些动作,就可能重复消费。

消费端必须幂等。

常见设计:

sql
create table mq_consume_log (
    id bigint primary key,
    message_id varchar(128) not null,
    consumer_group varchar(128) not null,
    status varchar(32) not null,
    created_at datetime not null,
    unique key uk_msg_group (message_id, consumer_group)
);

消费时:

java
@Transactional(rollbackFor = Exception.class)
public void consume(OrderCreatedEvent event) {
    if (!consumeLogRepository.tryInsert(event.messageId(), "order-search-sync")) {
        return;
    }

    searchIndexRepository.upsertOrder(event.orderId());
    consumeLogRepository.markSuccess(event.messageId(), "order-search-sync");
}

这样即使消息重复投递,也不会重复执行业务副作用。

商业场景:订单同步 ES 堆积

场景:

  1. 订单服务写 MySQL。
  2. 事务成功后发送订单变更消息。
  3. 搜索同步服务消费消息,查询 MySQL 最新订单快照,写入 ES。
  4. 活动高峰订单量暴涨,ES 写入慢,Kafka Lag 上升。

处理思路:

mermaid
flowchart TD
    A["订单消息 Lag 上升"] --> B["确认 ES 写入耗时"]
    B --> C["检查 bulk 大小和错误率"]
    C --> D["临时增加同步消费者"]
    D --> E{"分区是否足够"}
    E -- "足够" --> F["扩容消费者追赶"]
    E -- "不足" --> G["提高单实例批量能力"]
    F --> H["监控 Lag 下降速度"]
    G --> H
    H --> I["补偿任务校验 MySQL 和 ES"]

注意:

  1. 订单详情页仍然查 MySQL,不应该依赖 ES 强一致。
  2. ES 同步消息可以堆积,但必须有延迟监控和补偿。
  3. 消费端要做幂等 upsert。
  4. 不要因为 Lag 高就跳过 offset,否则搜索索引会漏数据。

商业场景:短信通知堆积

短信通知和订单状态不同。它通常允许降级和限流。

如果短信通道慢导致堆积:

  1. 限制非核心短信。
  2. 优先发送支付、登录、安全类短信。
  3. 营销短信进入延迟队列或丢弃,但要有业务确认。
  4. 对同一用户同类短信做合并或去重。
  5. 通道恢复后按优先级补发。

这说明不同业务的堆积处理策略不同。订单、库存、支付不能随便丢;营销、通知、日志可以按规则降级。

生产监控指标

指标为什么重要
生产 TPS判断入口流量
消费 TPS判断处理能力
Lag / Queue Depth判断堆积规模
最老消息等待时间判断业务延迟是否超 SLA
消费失败率判断是否有异常消息
重试消息量判断是否重试风暴
死信队列数量判断是否有不可处理消息
消费耗时 P95/P99判断慢在哪里
下游 DB/接口耗时判断瓶颈是否在下游
Broker 磁盘水位判断是否有存储风险
消费者线程池队列长度判断本机是否堆积

告警不要只看消息数量,还要看最老消息延迟。例如堆积 1 万条日志可能没事,但堆积 100 条支付消息超过 5 分钟就是严重问题。

关联跳转

详细原理、流程图、排查命令、Demo 和商业场景沉淀在本页。刷题复习入口看 MQ面试题

小结

消息堆积的本质是生产速度长期大于消费速度。排查时要看生产、Broker、消费者、下游和重试链路,处理时先止血、再定位、最后消化存量。背压是为了让系统在下游变慢时保持可控,不让压力无限堆到内存、线程池、数据库或 Broker 磁盘。真正成熟的 MQ 方案必须同时具备监控告警、幂等消费、失败重试、死信隔离、限流背压和补偿能力。