消息堆积与背压
消息堆积是 MQ 生产排障里非常高频的问题。
一句话理解:
消息堆积就是生产者写入消息的速度长期大于消费者处理消息的速度,Broker 中等待处理的消息越来越多。
背压可以理解为:
当下游处理不过来时,系统不要无限接收和堆内存排队,而是通过限流、暂停拉取、降低生产速度、延迟重试等方式,把压力向上游传递或吸收。
很多人遇到堆积只会说“加消费者”,但真实生产里这远远不够。因为堆积可能来自消费者慢、下游数据库慢、重试消息爆炸、顺序消息阻塞、分区太少、Broker 磁盘压力、ACK 配置错误、代码线程池无界排队等不同原因。
学习目标
| 目标 | 你需要掌握什么 |
|---|---|
| 知道是什么 | 理解堆积、积压、Lag、Ready、Unacked、重试堆积、背压的区别 |
| 知道为什么 | 理解生产速度、消费速度、下游能力和 Broker 存储之间的关系 |
| 知道怎么工作 | 会从生产、Broker、消费、下游、重试、顺序消息几个方向定位 |
| 知道不这样会怎样 | 知道盲目扩容、盲目清消息、提前 ACK、无界线程池的风险 |
| 会排查 | 能用 Kafka、RabbitMQ、RocketMQ 的指标判断堆积位置 |
| 会项目落地 | 能设计削峰、限流、幂等、重试、死信、补偿和容量预案 |
| 会复盘 | 能把“现象、原因、处理、预防”沉淀成稳定排障方案 |
堆积、积压、Lag、背压的区别
| 名词 | 含义 | 常见产品里的表现 |
|---|---|---|
| 消息堆积 | Broker 中待消费消息越来越多 | Queue depth 增大、Topic 消息滞留 |
| 消息积压 | 和堆积基本同义,更强调未处理存量 | 等待处理的消息量持续增长 |
| Lag | 消费进度落后生产进度多少 | Kafka Consumer Lag、RocketMQ Diff |
| Ready | RabbitMQ 中已入队但还没投递给消费者的消息 | 队列 Ready 数很高 |
| Unacked | RabbitMQ 中已投递但消费者未 ACK 的消息 | 消费者拿走了但没确认 |
| Retry 堆积 | 消费失败后不断进入重试队列 | 重试 Topic / 死信队列增长 |
| 背压 | 下游慢时限制上游或暂停拉取,防止系统被压垮 | pause/resume、限流、熔断、降级 |
堆积和背压不是一个概念。堆积是现象,背压是治理手段。
如果面试或排障问到“Lag 分布是什么”,不要只看总 Lag。要继续看每个 Kafka Partition、RocketMQ MessageQueue 或 RabbitMQ Queue 分别积压多少,判断是整体消费慢还是少数通道热点、慢消息、顺序阻塞。详细看:Lag 分布与堆积定位。
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["触发背压或限流"]堆积的核心公式
先记住一个非常有用的公式:
堆积增长速度 = 生产 TPS - 消费 TPS如果:
生产速度 = 5000 条/秒
消费速度 = 3000 条/秒那么:
每秒新增堆积 = 2000 条10 分钟后大约堆积:
2000 * 60 * 10 = 120 万条如果后续生产恢复到 1000 条/秒,消费速度提升到 5000 条/秒,那么消化速度是:
净消化速度 = 5000 - 1000 = 4000 条/秒消化 120 万条需要:
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 | 判断扩容消费者会不会把下游打挂 |
追平时间怎么算
净消化速度 = 成功消费 TPS - 当前生产 TPS
预计追平时间 = 当前堆积量 / 净消化速度如果 成功消费 TPS <= 当前生产 TPS,说明系统还在继续落后,理论上永远追不平。此时优先级不是“消化存量”,而是先止血:暂停补数据、限制生产速度、关闭非核心任务、隔离失败消息、保护下游。
例如:
当前堆积量 = 1800000
生产 TPS = 2000
成功消费 TPS = 5000
净消化速度 = 3000
预计追平时间 = 1800000 / 3000 = 600 秒也就是大约 10 分钟。如果最老消息 SLA 是 30 分钟,而且 Broker 磁盘还能撑 2 小时,这通常是可控堆积。如果预计追平时间是 4 小时,但消息保留窗口只有 2 小时,或者最老消息已经超过支付、库存、通知业务 SLA,就必须升级处理。
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 存储告警 | 磁盘水位高、文件保留压力大 | 高风险 |
判断是否异常,关键看三点:
- 堆积量是否持续增长。
- 最老消息延迟是否超过业务 SLA。
- 当前消费能力是否能在可接受时间内追平。
堆积原因总览
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 数字”,但不知道一条消息进入消费者之后经历了哪些步骤。真正定位时,要把消息消费链路拆开。
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 拉一批消息,再丢到本地线程池处理。如果本地线程池是无界队列,就可能出现一种假象:
- Broker Lag 短时间下降。
- 消费者 JVM 内部队列越来越长。
- 业务真正处理速度没变。
- 内存上涨、延迟上涨,最后 OOM 或任务超时。
flowchart TD
A["消费者快速拉取消息"] --> B["提交到无界队列"]
B --> C["Broker Lag 表面下降"]
B --> D["JVM 内存和本地队列上涨"]
D --> E["业务处理仍然慢"]
E --> F["最终 OOM 或超时"]所以判断消费能力时,不能只看 Broker Lag 是否下降,还要看业务成功 TPS。真正有效的消费是:业务处理成功,并且 ACK 或 Offset 提交成功。
背压不是简单少拉一点
背压的目标不是“让消费者变慢”,而是让系统在下游变慢时仍然可控。
没有背压时,典型过程是:
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 多少请求 | 把下游打满,引发全链路超时 |
| 失败重试速度 | 失败消息多久重试 | 重试风暴吞掉正常消费能力 |
| 上游生产速度 | 是否暂停补数据、限流入口 | 消费端永远追不上 |
一个可落地的背压策略
flowchart TD
A["消费者运行中"] --> B["采集本地队列长度、活跃线程、下游耗时"]
B --> C{"是否超过高水位"}
C -- "是" --> D["暂停拉取或降低 prefetch"]
D --> E["继续处理已拿到的消息"]
E --> F{"是否低于低水位"}
F -- "否" --> E
F -- "是" --> G["恢复拉取"]
C -- "否" --> H["正常拉取和处理"]为什么要有高水位和低水位两个阈值?因为如果只用一个阈值,消费者可能在“暂停、恢复、暂停、恢复”之间频繁抖动。高低水位可以形成缓冲区。
高水位:队列使用率超过 80%,暂停拉取
低水位:队列使用率低于 40%,恢复拉取有界线程池为什么是必须的
消费端本地线程池不要使用无界队列。无界队列会让你误以为“没有拒绝,系统还能接”,实际上只是把消息堆在 JVM 里。
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,会产生“假消费成功”:
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。
简化版思路:
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,导致每次消费都失败。它最危险的地方不是“这一条失败”,而是它会持续占用消费能力。
flowchart TD
A["消费者拿到毒消息"] --> B["业务处理失败"]
B --> C["立即重试"]
C --> D["再次失败"]
D --> E["消费者反复处理同一批失败消息"]
E --> F["正常消息排队等待"]
F --> G["Lag 继续增长"]处理原则:
| 做法 | 说明 |
|---|---|
| 区分异常类型 | 网络超时可重试,参数错误不要无限重试 |
| 延迟重试 | 给下游恢复时间,避免快速打满消费者 |
| 最大重试次数 | 超过次数进入死信或异常表 |
| 死信告警 | 死信不是垃圾桶,必须有人看 |
| 修复后重放 | 修复数据或代码后,从死信表按业务规则补偿 |
商业系统里,支付、库存、订单状态类消息进入死信后,不能只记录日志。要能按业务号查询事实源,确认是否需要补偿、撤销、人工审核或重放。
ACK 和 Offset 为什么影响堆积
不同 MQ 叫法不一样,但本质都是告诉 Broker:“这条消息我处理完了,可以推进进度。”
| 产品 | 确认机制 | 提前确认风险 | 过晚确认风险 |
|---|---|---|---|
| RabbitMQ | ACK | 业务失败但消息已删除 | Unacked 高、重复投递、内存压力 |
| Kafka | Commit Offset | 业务失败但 offset 已推进 | Rebalance 后重复消费 |
| RocketMQ | 返回消费状态 | 业务失败但返回成功 | 重试堆积、顺序阻塞 |
正确原则:
业务真正成功后再确认;确认失败或网络异常时,消费端要允许重复消息,并靠幂等保护业务。
生产 TPS、拉取 TPS、成功 TPS 的区别
排查堆积时要特别区分三个 TPS。
| 指标 | 含义 | 为什么重要 |
|---|---|---|
| 生产 TPS | 上游每秒写入 Broker 的消息数 | 判断压力源 |
| 拉取 TPS | 消费者每秒从 Broker 取到的消息数 | 判断 Broker 到消费者通不通 |
| 成功 TPS | 业务成功并 ACK/提交 Offset 的消息数 | 判断真实消化能力 |
最容易误判的是拉取 TPS。比如消费者每秒拉 5000 条,但业务每秒只成功 1000 条,其余都在本地线程池排队或失败重试。此时真实消化能力是 1000,不是 5000。
有效消化速度 = 成功 TPS - 生产 TPS如果成功 TPS 小于生产 TPS,即使消费者日志看起来“拉了很多”,堆积也一定会扩大。
堆积形成全过程
下面用订单同步 ES 举例,看堆积是怎么一步步形成的。
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 / Queue | order-event |
| ConsumerGroup | order-search-sync-group |
| 总 Lag / Queue Depth | 120 万 |
| 最大单分区 Lag | Partition 3 有 80 万 |
| 最老消息等待时间 | 42 分钟 |
第二层:确认速度关系
| 要看什么 | 例子 | 判断 |
|---|---|---|
| 生产 TPS | 3000/s | 压力是否还在 |
| 拉取 TPS | 5000/s | 是否能从 Broker 拿到 |
| 成功 TPS | 1200/s | 真实消化能力不足 |
| 失败 TPS | 800/s | 重试正在吞吐能力 |
如果拉取 TPS 高但成功 TPS 低,瓶颈一定在消费者内部或下游,而不是 Broker 没投递。
第三层:定位消费者内部
| 要看什么 | 说明 |
|---|---|
| 消费线程活跃数 | 是否线程全忙 |
| 本地队列长度 | 是否内部排队 |
| 单条处理 P95/P99 | 是否长尾严重 |
| GC 和内存 | 是否本地堆积导致内存上涨 |
| 错误日志 | 是否大量失败消息反复重试 |
第四层:定位下游
| 下游 | 证据 |
|---|---|
| MySQL | 慢 SQL、锁等待、连接池 active 打满 |
| ES | bulk rejected、refresh/merge 压力、写入延迟 |
| Redis | 慢命令、连接池耗尽、网络延迟 |
| HTTP 接口 | 超时、429、熔断、连接池等待 |
第五层:确认处理是否生效
处理后继续观察:
- 生产 TPS 是否下降或恢复正常。
- 成功 TPS 是否持续大于生产 TPS。
- 总 Lag 是否下降。
- 最大分区 Lag 是否下降。
- 最老消息等待时间是否下降。
- 失败率和重试量是否下降。
- Broker 磁盘水位是否稳定。
如果只看到总 Lag 下降,但最老消息等待时间不降,可能还有慢消息、顺序阻塞或局部分区卡住。
标准排查流程
排查堆积不要一上来就重启消费者。先回答五个问题:
- 哪个 Topic / Queue / ConsumerGroup 堆积?
- 当前生产 TPS 和消费 TPS 分别是多少?
- 是所有分区/队列都堆积,还是某几个特别严重?
- 消费失败率、重试量、死信量是否增加?
- 消费者卡在哪里:CPU、线程池、数据库、远程接口、锁等待、GC?
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 表示消费落后量。
常用命令:
kafka-consumer-groups.sh \
--bootstrap-server 127.0.0.1:9092 \
--describe \
--group order-search-sync-group典型输出里关注:
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 特别高,常见原因是:
- key 分布不均导致热点分区。
- 某条消息处理特别慢。
- 该分区对应消费者实例异常。
- 分区内消息顺序处理,无法被多个消费者并行处理。
Kafka 要特别注意:
同一个 Consumer Group 中,一个分区同一时刻只能被一个消费者消费。消费者实例数超过分区数,多出来的实例不会提升吞吐。
flowchart TD
A["Topic 有 3 个 Partition"] --> B["最多 3 个消费者并行消费"]
B --> C{"启动 6 个消费者有用吗"}
C -- "没有完全有用" --> D["只有 3 个消费者能分到分区"]
C -- "想提升吞吐" --> E["增加分区或提高单消费者处理能力"]RabbitMQ 怎么看堆积
RabbitMQ 重点看队列里的几个指标:
| 指标 | 含义 | 判断 |
|---|---|---|
| Ready | 已在队列中,等待投递给消费者 | Ready 高说明消费者拿得慢或数量不足 |
| Unacked | 已投递给消费者,但还没 ACK | Unacked 高说明消费者处理慢或 ACK 卡住 |
| Publish rate | 生产速率 | 判断入口压力 |
| Deliver / Ack rate | 投递和确认速率 | 判断消费能力 |
| Consumers | 消费者数量 | 判断是否掉线或不足 |
典型情况:
| 现象 | 可能原因 |
|---|---|
| Ready 很高,Unacked 不高 | 消费者数量不足、prefetch 太小、消费者没有正常拉取 |
| Unacked 很高 | 消费者拿了消息但处理慢,或者没 ACK |
| Ready 和 Unacked 都高 | 生产过快,消费端整体处理不过来 |
| 消费者数为 0 | 消费者服务挂了或连接失败 |
RabbitMQ 中 prefetch 很关键。它控制一个消费者最多同时拿多少未确认消息。
channel.basicQos(50);如果 prefetch 太大,消费者会一次拿走很多消息,导致 Unacked 很高,其他消费者分不到消息;如果太小,吞吐可能上不去。
RocketMQ 怎么看堆积
RocketMQ 常看 ConsumerGroup 的消费进度和 Diff。
常用排查方向:
| 指标 | 含义 |
|---|---|
| Topic 队列堆积 | 某个 Topic 下消息等待消费 |
| ConsumerGroup Diff | 消费进度落后多少 |
| Retry Topic | 消费失败进入重试的消息 |
| DLQ | 多次失败后的死信消息 |
| 消费 TPS | 消费者实际吞吐 |
| Broker 磁盘 | CommitLog、ConsumeQueue 存储压力 |
RocketMQ 顺序消息要特别注意:如果某条消息一直处理失败,它所在队列的后续消息可能被阻塞,导致局部堆积。
flowchart TD
A["队列 Q0"] --> B["消息 1 成功"]
B --> C["消息 2 失败或超慢"]
C --> D["消息 3 等待"]
D --> E["消息 4 等待"]
C --> F["Q0 局部堆积"]处理顺序消息堆积时,不能简单跳过消息。要看业务是否允许:
- 修复数据后重试。
- 把异常消息转人工处理。
- 将失败原因落库,后续补偿。
- 调整业务 key,减少热点队列。
堆积时先做什么
生产故障中,建议按这个顺序处理:
1. 先止血
目标是让堆积不再继续恶化。
常见动作:
- 临时关闭非核心生产入口。
- 对上游接口限流。
- 暂停补数据、重放、批量导入任务。
- 降级非核心消费者逻辑。
- 对明显失败的消息进入死信或异常表,不要无限快速重试。
2. 再定位瓶颈
判断消费者慢在哪里:
| 瓶颈 | 现象 | 处理 |
|---|---|---|
| CPU 高 | 计算、序列化、压缩、复杂规则慢 | 优化代码或扩容 |
| DB 慢 | 慢 SQL、锁等待、连接池满 | 加索引、批处理、扩容、限流 |
| 远程接口慢 | HTTP 调用耗时高 | 超时、熔断、异步补偿 |
| 线程池满 | 活跃线程满,队列增长 | 有界队列、调并发、背压 |
| 单分区热点 | 某分区 Lag 特别高 | 调整 key、增加分区、拆 Topic |
| 重试过多 | 失败消息反复消费 | 延迟重试、死信、人工处理 |
3. 最后消化存量
消化堆积时要计算预计时间:
预计消化时间 = 当前堆积量 / (当前消费 TPS - 当前生产 TPS)如果消费 TPS 小于生产 TPS,说明永远追不上,必须先限流或扩容。
消化手段:
- 临时增加消费者实例。
- 提高单消费者批处理能力。
- 临时放宽非核心业务校验。
- 将异常消息隔离到死信队列。
- 对历史低价值消息按业务规则降级处理。
- 增加分区或队列,但要评估顺序和路由影响。
扩容消费者短暂有效,后面又堆积
这一节先给出排查总览。完整知识点拆到了独立页面:消费者扩容后再次堆积。
现象通常是:
- 刚增加消费者实例后,Lag 或队列深度开始下降。
- 过一会儿下降速度变慢。
- 随后堆积又开始增长。
这说明:扩容消费者只短暂提高了消费入口能力,但真正瓶颈没有消失,或者新的瓶颈被打出来了。
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 |
排查流程:
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,必须先确认业务是否允许丢弃、是否有补偿来源、是否能从主库重建。
背压怎么设计
背压的核心思想是:消费者处理不过来时,不要把压力无限堆到本机内存里。
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:有界线程池防止本机堆积
错误做法是使用无界队列:
ExecutorService executor = Executors.newFixedThreadPool(20);newFixedThreadPool 底层使用无界队列。消费速度跟不上时,消息任务会不断堆到 JVM 内存里,最后可能 OOM。
更推荐使用有界队列:
ThreadPoolExecutor executor = new ThreadPoolExecutor(
20,
20,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1000),
new ThreadPoolExecutor.CallerRunsPolicy()
);CallerRunsPolicy 的含义是:线程池和队列都满了,就让提交任务的线程自己执行任务。这样拉取线程会变慢,形成一种简单背压,不会无限把任务塞进内存。
Kafka Demo:处理不过来时暂停拉取
Kafka 消费者可以在本地队列满时暂停拉取,处理能力恢复后再继续。
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 已经认为消息处理完了。
更稳的做法是:
- 控制每次 poll 的数量。
- 业务真正成功后再提交 offset。
- 批量处理时记录每个分区成功处理到的最大 offset。
- 失败消息进入延迟重试或死信,不要阻塞整个消费组。
RabbitMQ Demo:手动 ACK
RabbitMQ 消费端应该在业务成功后再 ACK。
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,否则它会反复投递,造成重试风暴。
幂等是消化堆积的前提
堆积后常常需要重试、扩容、重放、补偿。只要有这些动作,就可能重复消费。
消费端必须幂等。
常见设计:
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)
);消费时:
@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 堆积
场景:
- 订单服务写 MySQL。
- 事务成功后发送订单变更消息。
- 搜索同步服务消费消息,查询 MySQL 最新订单快照,写入 ES。
- 活动高峰订单量暴涨,ES 写入慢,Kafka Lag 上升。
处理思路:
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"]注意:
- 订单详情页仍然查 MySQL,不应该依赖 ES 强一致。
- ES 同步消息可以堆积,但必须有延迟监控和补偿。
- 消费端要做幂等 upsert。
- 不要因为 Lag 高就跳过 offset,否则搜索索引会漏数据。
商业场景:短信通知堆积
短信通知和订单状态不同。它通常允许降级和限流。
如果短信通道慢导致堆积:
- 限制非核心短信。
- 优先发送支付、登录、安全类短信。
- 营销短信进入延迟队列或丢弃,但要有业务确认。
- 对同一用户同类短信做合并或去重。
- 通道恢复后按优先级补发。
这说明不同业务的堆积处理策略不同。订单、库存、支付不能随便丢;营销、通知、日志可以按规则降级。
生产监控指标
| 指标 | 为什么重要 |
|---|---|
| 生产 TPS | 判断入口流量 |
| 消费 TPS | 判断处理能力 |
| Lag / Queue Depth | 判断堆积规模 |
| 最老消息等待时间 | 判断业务延迟是否超 SLA |
| 消费失败率 | 判断是否有异常消息 |
| 重试消息量 | 判断是否重试风暴 |
| 死信队列数量 | 判断是否有不可处理消息 |
| 消费耗时 P95/P99 | 判断慢在哪里 |
| 下游 DB/接口耗时 | 判断瓶颈是否在下游 |
| Broker 磁盘水位 | 判断是否有存储风险 |
| 消费者线程池队列长度 | 判断本机是否堆积 |
告警不要只看消息数量,还要看最老消息延迟。例如堆积 1 万条日志可能没事,但堆积 100 条支付消息超过 5 分钟就是严重问题。
关联跳转
详细原理、流程图、排查命令、Demo 和商业场景沉淀在本页。刷题复习入口看 MQ面试题。
小结
消息堆积的本质是生产速度长期大于消费速度。排查时要看生产、Broker、消费者、下游和重试链路,处理时先止血、再定位、最后消化存量。背压是为了让系统在下游变慢时保持可控,不让压力无限堆到内存、线程池、数据库或 Broker 磁盘。真正成熟的 MQ 方案必须同时具备监控告警、幂等消费、失败重试、死信隔离、限流背压和补偿能力。
