消费者扩容后再次堆积
这个知识点专门解释一个生产里很常见、也很容易误判的现象:
MQ 已经堆积了,临时把消费者实例加上去,Lag 或队列深度一开始下降,过一会儿又开始堆积。
这不是 MQ “不稳定”,也不一定是消费者“没有扩成功”。它通常说明:扩容消费者只提升了消费入口的并发,但整条链路真正的瓶颈不在消费者数量,或者扩容后把新的瓶颈打出来了。
学习目标
| 目标 | 需要掌握的内容 |
|---|---|
| 知道现象 | 明白为什么扩容初期有效,后面又重新堆积 |
| 知道原理 | 理解 MQ 消费吞吐由分区/队列、消费者、线程池、下游、ACK、重试共同决定 |
| 会定位 | 能通过 Lag 分布、消费 TPS、处理耗时、错误率、重试量、线程池、连接池、慢 SQL 判断瓶颈 |
| 会处理 | 会选择限流、降并发、批处理、失败隔离、增加分区、优化下游、补偿重放等措施 |
| 会预防 | 能提前设计监控、背压、幂等、死信、容量预案和故障演练 |
先用一句话理解
消费者扩容能不能解决堆积,取决于扩容后有效消费 TPS是否持续大于生产 TPS。
堆积是否下降 = 有效消费 TPS > 生产 TPS但有效消费 TPS 不是只由消费者实例数决定,而是由整条链路的最短板决定:
有效消费 TPS =
min(
Broker 投递能力,
分区/队列并行度,
消费者本机处理能力,
下游数据库/接口/ES 能力,
ACK 或 Offset 提交能力,
重试与异常消息处理能力
)所以扩容消费者只是提高了其中一项。如果最短板在数据库、远程接口、ES、Redis、分区数量、热点 key、重试风暴或 Broker,继续加消费者不会根治问题。
为什么会短暂有效
先看一个典型时间线。
flowchart TD
A["消息开始堆积"] --> B["增加消费者实例"]
B --> C["入口并发上升"]
C --> D["短时间消费 TPS 提高"]
D --> E["Lag 或队列深度下降"]
E --> F["下游或分区瓶颈出现"]
F --> G["单条处理耗时变长"]
G --> H["有效消费 TPS 下降"]
H --> I["再次开始堆积"]短暂有效一般有三层原因。
1. 扩容确实提高了入口并发
刚扩容时,更多消费者可以同时拉取消息、反序列化、提交线程池、调用业务逻辑,所以消费 TPS 会短暂上升。只要这段时间消费 TPS 大于生产 TPS,堆积就会下降。
2. 真正瓶颈被更大并发打出来
消费者变多后,写数据库、写 ES、访问 Redis、调用三方接口的并发也会变多。下游一旦开始排队、锁等待、限流、拒绝请求,单条消息处理耗时会升高。
例如原来:
10 个消费者,每条处理 20ms,消费 TPS 约 500扩容后短时间:
30 个消费者,每条仍然 20ms,消费 TPS 约 1500但数据库被打满后:
30 个消费者,每条处理变成 200ms,消费 TPS 约 150这时消费者更多,反而更容易造成连接池等待、锁竞争、慢 SQL、超时重试,堆积会重新增长。
3. 并行度被分区或队列限制
MQ 并不是消费者越多吞吐越高。
| 产品 | 并行度限制 |
|---|---|
| Kafka | 同一个 Consumer Group 中,一个 Partition 同一时刻只能分配给一个消费者 |
| RocketMQ | 一个 MessageQueue 同一时刻通常只会被同组一个消费者消费,顺序消息更明显 |
| RabbitMQ | 同一个 Queue 可以多个消费者竞争,但会受 prefetch、ACK、Unacked、下游能力影响 |
如果 Kafka 只有 8 个分区,启动 20 个同组消费者,真正能并行消费的最多也就是 8 个分区,剩余消费者空闲。
flowchart TD
A["8 个 Partition"] --> B["最多 8 个活跃消费者"]
B --> C{"启动 20 个消费者"}
C -- "8 个有任务" --> D["参与消费"]
C -- "12 个空闲" --> E["不能提升吞吐"]本质:最短板决定吞吐
可以把 MQ 消费链路看成一条流水线。
flowchart TD
A["Broker 拉取/投递"] --> B["消费者反序列化"]
B --> C["线程池执行"]
C --> D["业务校验"]
D --> E["DB / ES / Redis / HTTP"]
E --> F["ACK 或提交 Offset"]任何一段慢下来,整条链路都会慢。
| 短板位置 | 典型表现 | 为什么扩容后又堆积 |
|---|---|---|
| 分区/队列少 | 消费者数量增加,但消费 TPS 不涨 | 并行通道数量固定,多出来的实例没有任务 |
| 热点 key | 某个分区 Lag 特别高 | 同一个热点业务 key 被路由到同一分区,只能串行或低并发处理 |
| 数据库慢 | 慢 SQL、锁等待、连接池满 | 扩容把更多写请求打到 DB,单条耗时升高 |
| ES 慢 | Bulk rejected、写入延迟高 | 扩容后 ES 写线程池或刷新压力变大 |
| 远程接口慢 | 超时、429、熔断、连接池满 | 下游限流后消费者线程大量等待 |
| 重试风暴 | Retry Topic 或死信增长 | 失败消息被更快取出并更快失败,吞噬消费能力 |
| ACK 过晚 | Unacked 高、重复投递多 | 消息已投递但确认慢,Broker 认为还没处理完 |
| ACK 过早 | 表面 Lag 下降,但业务缺数据 | 业务失败后消息已确认,只能靠补偿 |
| 本机线程池满 | 活跃线程满、队列增长 | 消费端内部开始排队,JVM 内存和延迟上升 |
| Broker 压力 | 拉取慢、网络高、磁盘高 | 更多消费者加重 Broker 协调和网络压力 |
排查第一步:先看 TPS 关系
不要先重启,也不要先继续加机器。先画出生产 TPS、消费 TPS、堆积量和最老消息延迟。
堆积增长速度 = 生产 TPS - 消费 TPS如果扩容后:
生产 TPS = 3000/s
消费 TPS = 4200/s说明存量应该下降,只需要计算多久追平。
如果过一会儿变成:
生产 TPS = 3000/s
消费 TPS = 1800/s说明新的瓶颈已经出现,继续堆积是必然结果。
扩容反弹的指标时间线
生产里最好不要只看一个时间点,要看扩容前、扩容刚完成、扩容 5 到 10 分钟后、扩容 30 分钟后的指标变化。
flowchart TD
A["T0 扩容前<br/>Lag 上升"] --> B["T1 刚扩容<br/>拉取 TPS 上升"]
B --> C["T2 短暂恢复<br/>Lag 下降"]
C --> D["T3 新瓶颈出现<br/>P95 耗时升高"]
D --> E["T4 有效 TPS 下降<br/>Lag 再次上升"]示例指标:
| 时间 | 生产 TPS | 拉取 TPS | 成功 TPS | P95 耗时 | 错误率 | Lag 趋势 | 判断 |
|---|---|---|---|---|---|---|---|
| T0 扩容前 | 3000/s | 1800/s | 1700/s | 80ms | 1% | 上升 | 消费追不上 |
| T1 刚扩容 | 3000/s | 5000/s | 4200/s | 90ms | 1% | 下降 | 短暂有效 |
| T2 5 分钟后 | 3000/s | 5200/s | 2600/s | 350ms | 8% | 持平转升 | 下游开始慢或失败增加 |
| T3 15 分钟后 | 3000/s | 5200/s | 1500/s | 900ms | 20% | 快速上升 | 新瓶颈完全暴露 |
这张表的关键是区分:
| 指标变化 | 说明 |
|---|---|
| 拉取 TPS 高,成功 TPS 低 | 消息已经进消费者,但业务没真正处理完 |
| P95/P99 升高 | 下游、锁、线程池或慢消息导致长尾 |
| 错误率升高 | 失败消息开始吞噬处理能力 |
| Lag 下降但最老等待时间不降 | 可能只处理了新消息,旧分区或慢消息仍卡住 |
| 消费者数增加但拉取 TPS 不涨 | 分区/队列并行度、Broker 或消费者配置限制 |
所以面试里回答“扩容后又堆积”时,不要只说“下游瓶颈”。更完整的说法是:
我会先拉出扩容前后的生产 TPS、拉取 TPS、成功 TPS、P95/P99、错误率、重试量、Lag 分布和最老消息等待时间。拉取 TPS 上升但成功 TPS 下降,说明瓶颈在消费者内部或下游;消费者数增加但拉取 TPS 不涨,说明并行度或 Broker 投递受限;错误率和 Retry 增长,说明重试风暴;少数分区 Lag 高,说明热点、慢消息或局部消费者异常。
根因决策树
下面这棵树用于生产现场快速定位,不要跳步骤。
flowchart TD
A["扩容后再次堆积"] --> B{"消费者实例是否真正上线"}
B -- "否" --> C["检查发布、注册、Group、appname、连接"]
B -- "是" --> D{"拉取 TPS 是否上升"}
D -- "否" --> E{"消费者数是否超过分区/队列数"}
E -- "是" --> F["并行度不足<br/>加实例无效"]
E -- "否" --> G["查 Broker、网络、消费者配置"]
D -- "是" --> H{"成功 TPS 是否持续大于生产 TPS"}
H -- "是" --> I["计算追平时间<br/>继续观察 SLA"]
H -- "否" --> J{"P95/P99 是否升高"}
J -- "是" --> K["查 DB、ES、HTTP、锁、连接池"]
J -- "否" --> L{"错误率或 Retry 是否升高"}
L -- "是" --> M["重试风暴或毒消息"]
L -- "否" --> N{"Lag 是否集中在少数分区"}
N -- "是" --> O["热点 key、慢消息、局部消费者异常"]
N -- "否" --> P["整体容量不足或 Broker 压力"]判断一:消费者是否真正上线
扩容不是“Pod 数量变多”就结束,要确认消费者真的加入消费组或队列消费。
| 产品 | 看什么 |
|---|---|
| Kafka | Consumer Group 成员数、分区分配结果、Rebalance 日志 |
| RocketMQ | ConsumerGroup 在线实例、MessageQueue 分配 |
| RabbitMQ | Queue 的 Consumers 数、连接和 Channel 状态 |
常见假扩容:
- 新实例启动失败。
- 配错 ConsumerGroup。
- 订阅 Topic 或 tag 不一致。
- 网络不通,连不上 Broker。
- Kafka Rebalance 后新实例没有分到分区。
判断二:拉取 TPS 是否上升
如果消费者上线了,但拉取 TPS 没上升,说明消息没有更快进入消费者。
可能原因:
| 原因 | 证据 | 处理 |
|---|---|---|
| 分区/队列数不足 | 活跃消费者数小于实例数 | 增加分区/队列或提升单消费者能力 |
| Broker 投递慢 | Broker 请求耗时、网络、磁盘高 | 查 Broker 负载 |
| 消费者配置限制 | poll/batch/prefetch 太小 | 调整批量和预取 |
| Rebalance 抖动 | 消费组频繁重新分配 | 查心跳和处理时长 |
判断三:成功 TPS 是否上升
如果拉取 TPS 上升,但成功 TPS 没上升,说明消息卡在消费者内部或下游。
flowchart TD
A["拉取 TPS 上升"] --> B["消息进入消费者"]
B --> C{"成功 TPS 不升"}
C --> D["线程池排队"]
C --> E["DB/ES/HTTP 慢"]
C --> F["业务锁等待"]
C --> G["失败重试"]
C --> H["ACK/Offset 提交慢"]这时最重要的是看:
- 消费者本地线程池队列。
- 单条处理耗时 P95/P99。
- DB 连接池 active 和等待时间。
- 慢 SQL、锁等待、ES rejected、HTTP 超时。
- 错误率、Retry、DLQ。
典型根因全过程
根因一:分区或队列并行度不足
过程:
flowchart TD
A["Topic 只有 8 个分区"] --> B["扩到 20 个消费者"]
B --> C["最多 8 个消费者分到分区"]
C --> D["12 个消费者空闲"]
D --> E["成功 TPS 不明显提升"]
E --> F["Lag 继续增长"]判断证据:
| 证据 | 说明 |
|---|---|
| 消费者实例数大于分区数 | 多出来的实例没有分区 |
| 每个活跃消费者 CPU 不高 | 不是机器算力不足 |
| 拉取 TPS 不随实例数增长 | 并行通道限制 |
处理方式:
- 增加分区或队列,但要评估顺序语义。
- 提高单分区处理能力,例如批处理。
- 拆 Topic,把热点业务拆出去。
- 调整路由 key,避免数据过度集中。
根因二:下游数据库被打满
过程:
flowchart TD
A["消费者扩容"] --> B["并发写 DB 增加"]
B --> C["连接池 active 打满"]
C --> D["SQL 排队和锁等待"]
D --> E["单条处理耗时升高"]
E --> F["成功 TPS 下降"]
F --> G["Lag 再次上升"]判断证据:
| 证据 | 说明 |
|---|---|
| DB 连接池 active 长期接近 max | 消费者线程在等连接 |
| 慢 SQL 增多 | SQL 或索引撑不住并发 |
| 锁等待增多 | 多消费者更新同一批数据 |
| 消费者 CPU 不高但线程阻塞 | 等 DB,不是算力不足 |
处理方式:
- 降低消费者并发,先保护 DB。
- 优化 SQL 和索引。
- 批量写入,减少单条网络往返。
- 缩小事务范围,避免长事务。
- 对热点数据做分片或串行化处理。
根因三:重试风暴
过程:
flowchart TD
A["某类消息持续失败"] --> B["进入快速重试"]
B --> C["扩容后失败消息被更快取出"]
C --> D["消费者大量时间处理失败消息"]
D --> E["正常消息得不到处理"]
E --> F["Lag 再次上升"]判断证据:
| 证据 | 说明 |
|---|---|
| 错误率升高 | 不是纯容量问题 |
| Retry Topic 增长 | 失败消息在反复消耗资源 |
| 同一异常反复出现 | 参数错误、数据错误、代码 bug |
| 死信队列开始增长 | 多次重试仍失败 |
处理方式:
- 区分短暂故障和永久故障。
- 永久故障不要无限重试,进入死信或异常表。
- 延迟重试,避免立即打爆消费者。
- 修复毒消息或数据问题后定向补偿。
根因四:热点 key 或慢消息
过程:
flowchart TD
A["热点订单/商户/设备 key"] --> B["持续路由到同一分区"]
B --> C["该分区消息远多于其他分区"]
C --> D["负责消费者长期处理同一通道"]
D --> E["少数分区 Lag 很高"]
E --> F["总 Lag 看起来也很高"]判断证据:
| 证据 | 说明 |
|---|---|
| Lag 集中在少数分区 | 不是整体消费慢 |
| 最大分区 key 分布集中 | 热点 key |
| 某条消息耗时特别长 | 慢消息 |
| 其他分区很快归零 | 局部阻塞 |
处理方式:
- 找最大 Lag 分区。
- 找负责该分区的消费者实例。
- 查该消费者日志和线程栈。
- 统计该分区消息 key 分布。
- 对热点 key 拆分、散列或单独 Topic。
- 对慢消息隔离处理,避免卡住后续消息。
根因五:Rebalance 抖动
过程:
flowchart TD
A["扩容消费者"] --> B["触发 Rebalance"]
B --> C["分区暂停并重新分配"]
C --> D["部分消费者处理时间过长"]
D --> E["心跳或 poll 超时"]
E --> F["再次 Rebalance"]
F --> G["消费吞吐上下波动"]
G --> H["Lag 反复下降又上升"]判断证据:
| 证据 | 说明 |
|---|---|
| 消费组频繁 Rebalance | 组不稳定 |
| 消费者日志有 poll 超时 | 单次处理太久 |
| Lag 周期性下降又暴涨 | 分区频繁暂停 |
| 新实例不断加入退出 | 部署或健康检查异常 |
处理方式:
- 减小单次 poll 的处理量。
- 把业务处理和 poll 心跳解耦。
- 调整
max.poll.interval.ms,但不要只靠调大掩盖慢处理。 - 扩容分批进行,避免一次性大量实例加入。
- 检查容器健康检查和 OOM。
排查流程
flowchart TD
A["扩容后再次堆积"] --> B["看生产 TPS 和消费 TPS"]
B --> C{"消费 TPS 是否持续大于生产 TPS"}
C -- "是" --> D["计算消化时间"]
C -- "否" --> E["定位消费瓶颈"]
E --> F["看分区/队列 Lag 分布"]
E --> G["看消费耗时 P95/P99"]
E --> H["看失败率和重试量"]
E --> I["看线程池和连接池"]
E --> J["看 DB/ES/HTTP 慢日志"]
F --> K["判断热点或并行度不足"]
G --> L["判断下游变慢"]
H --> M["判断重试风暴"]Lag 分布是什么
Lag 分布不是只看“总共堆积多少条”,而是看这些堆积分别落在哪些分区、队列或消费通道上。
这一节先讲在“扩容后再次堆积”场景里怎么用 Lag 分布判断方向。完整的 Lag 分布定义、Kafka/RocketMQ/RabbitMQ 指标、热点 key、慢消息和重试风暴定位流程,单独看:Lag 分布与堆积定位。
以 Kafka 为例:
单分区 Lag = LOG-END-OFFSET - CURRENT-OFFSET
消费组总 Lag = 所有分区 Lag 之和
Lag 分布 = 每个 Partition 分别落后多少例如总 Lag 都是 9000,下面两种情况的含义完全不同。
1. 均匀分布
Partition 0 Lag = 3000
Partition 1 Lag = 3000
Partition 2 Lag = 3000这说明多数分区都在落后,更像是整体消费能力不足,或者消费者共同依赖的下游数据库、ES、Redis、远程接口变慢。
2. 倾斜分布
Partition 0 Lag = 100
Partition 1 Lag = 8800
Partition 2 Lag = 100这说明堆积集中在少数分区,更像是热点 key、单分区慢消息、某个消费者异常、顺序消息阻塞或数据分布不均。
flowchart TD
A["查看 Lag 分布"] --> B{"Lag 是否集中在少数分区"}
B -- "否,整体都高" --> C["整体消费能力不足"]
C --> D["看消费者数量、线程池、DB、ES、HTTP"]
B -- "是,少数分区高" --> E["局部热点或局部阻塞"]
E --> F["看热点 key、慢消息、单消费者状态、顺序阻塞"]为什么不能只看总 Lag?因为总 Lag 只能告诉你“积压多不多”,Lag 分布才能告诉你“问题是全局慢还是局部卡住”。
| Lag 分布现象 | 常见原因 | 处理方向 |
|---|---|---|
| 所有分区 Lag 都高 | 消费者整体不足、下游整体慢、Broker 拉取慢 | 扩容消费者、优化下游、批处理、限流生产 |
| 少数分区 Lag 很高 | 热点 key、慢消息、单消费者异常 | 调整 key、拆热点、定位慢消息、重启异常实例 |
| 某个分区 Lag 不动 | 该分区消费者卡死、顺序消息阻塞 | 看线程栈、错误日志、单条消息处理耗时 |
| Lag 上下剧烈波动 | Rebalance、消费者频繁重启、下游抖动 | 查消费组稳定性、心跳、max.poll.interval.ms |
| Lag 不高但业务延迟高 | 单条消息处理慢,或消息时间跨度大 | 看最老消息等待时间,不只看数量 |
在 RocketMQ 里,类似地看每个 MessageQueue 的 Diff;在 RabbitMQ 里,虽然不是 Kafka 这种分区日志模型,也要看不同 Queue 的 Ready、Unacked、消费者数量和投递速率。核心思想一样:不要只看总量,要看堆积是不是集中在少数通道。
指标怎么看
| 指标 | 正常理解 | 异常信号 |
|---|---|---|
| 生产 TPS | 上游写入速度 | 活动、补数、重放任务导致持续高于预期 |
| 消费 TPS | 真正处理成功的速度 | 扩容后先升后降,说明被新瓶颈限制 |
| Lag / Queue Depth | 待处理存量 | 持续增长说明消费追不上生产 |
| Lag 分布 | 各分区、队列分别堆积多少 | 判断是整体慢还是局部热点、局部阻塞 |
| 最老消息等待时间 | 业务延迟 | 比消息数量更能反映 SLA 风险 |
| 单条处理耗时 P95/P99 | 消费端业务耗时 | 扩容后耗时升高,常见是下游被打慢 |
| 错误率 | 消费失败比例 | 错误率上升会触发重试风暴 |
| Retry / DLQ | 重试和死信量 | 失败消息正在吞噬正常消费能力 |
| 线程池队列长度 | 本机是否排队 | 队列增长说明消费者内部堆积 |
| DB 连接池 active | DB 连接是否够用 | active 长期打满说明大量线程等连接 |
| 慢 SQL / 锁等待 | DB 是否是短板 | 并发写入导致锁竞争或索引缺失 |
| ES rejected | ES 是否拒绝写入 | Bulk 太大、并发太高、分片压力过大 |
| Broker 磁盘和网络 | Broker 是否健康 | 磁盘高水位、网络高、拉取延迟高 |
Kafka 场景怎么判断
Kafka 最常见的误区是:消费者实例数超过分区数后,继续扩容没有意义。
1. 看每个分区 Lag
kafka-consumer-groups.sh \
--bootstrap-server 127.0.0.1:9092 \
--describe \
--group order-search-sync-group如果只有一个分区 Lag 特别高:
PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
0 1000 1200 200
1 2000 9000 7000
2 3000 3100 100这通常不是整体消费者不足,而是分区 1 有热点 key、慢消息、异常消费者或该分区下游处理慢。
2. 看消费者数量是否超过分区数
分区数 = 8
消费者实例数 = 20
有效并行消费者最多 = 8如果要继续提升吞吐,需要考虑:
- 增加 Topic 分区数,但要评估 key 路由和顺序语义。
- 提高单分区处理能力,例如批处理、减少单条耗时。
- 拆 Topic,把热点业务单独拆出去。
- 调整 key,避免热点集中到单个分区。
3. 看 Rebalance 是否频繁
Kafka 扩容、缩容、消费者心跳异常、max.poll.interval.ms 超时,都可能触发 Rebalance。Rebalance 期间分区重新分配,消费会暂停或抖动。
sequenceDiagram
participant C1 as 消费者1
participant C2 as 新消费者
participant G as Consumer Group
participant P as Partition
C2->>G: 加入消费组
G->>C1: 暂停并重新分配
G->>C2: 分配部分分区
C1->>P: 恢复消费
C2->>P: 开始消费频繁 Rebalance 会让扩容效果看起来忽高忽低。
RabbitMQ 场景怎么判断
RabbitMQ 要重点看 Ready 和 Unacked。
| 现象 | 重点判断 |
|---|---|
| Ready 很高,Unacked 不高 | 消费者没拿到消息,可能消费者少、连接异常、prefetch 太小 |
| Unacked 很高 | 消费者已经拿到消息但处理慢,可能业务卡住、ACK 慢、prefetch 太大 |
| Ready 和 Unacked 都高 | 生产太快,消费端和下游整体扛不住 |
| Consumers 下降 | 消费者掉线或部署异常 |
prefetch 控制每个消费者最多同时持有多少未 ACK 消息。
channel.basicQos(50);如果 prefetch 太大,单个慢消费者会拿走很多消息,导致其他消费者没有机会处理;如果 prefetch 太小,吞吐上不去。它不是越大越好,而是要和单条处理耗时、消费者数量、下游能力一起估算。
RocketMQ 场景怎么判断
RocketMQ 要关注 Topic 下每个 MessageQueue 的消费进度、ConsumerGroup Diff、重试队列和死信队列。
| 现象 | 可能原因 |
|---|---|
| 某几个队列 Diff 特别高 | 热点 key、顺序消息阻塞、单队列消费者慢 |
| Retry Topic 增长 | 消费失败反复重试 |
| DLQ 增长 | 多次失败后进入死信,需要人工或补偿处理 |
| 扩容后吞吐不涨 | MessageQueue 数量不足或下游瓶颈 |
| 顺序消息堆积 | 某条消息失败导致同队列后续消息不能继续处理 |
顺序消息尤其要小心。为了保证同一个业务 key 的顺序,消息会进入同一个队列。如果某条消息处理失败,后面的消息不能随便越过它,否则顺序语义就被破坏。
处理策略
1. 先止血
止血不是为了马上清空堆积,而是让堆积不再继续恶化。
| 场景 | 止血动作 |
|---|---|
| 上游生产过快 | 对生产接口限流,暂停补数据、重放、批导任务 |
| 下游 DB 慢 | 降低消费并发,减少事务范围,暂停非核心写入 |
| ES 写入慢 | 调整 bulk 大小,降低并发,查看 rejected,必要时暂停非核心同步 |
| 远程接口慢 | 设置超时,熔断降级,降低并发,转异步补偿 |
| 重试风暴 | 把不可恢复异常进入死信或异常表,不要立即无限重试 |
2. 再定位最短板
按照这个顺序会比较稳:
- 生产 TPS 是否异常升高。
- 消费 TPS 是否真的提升。
- 单条消费耗时是否升高。
- Lag 是否集中在少数分区/队列。
- 失败率、重试量、死信是否升高。
- DB、ES、Redis、HTTP、连接池、线程池是否打满。
- Broker 磁盘、网络、拉取延迟是否异常。
3. 最后消化存量
如果已经确认消费 TPS 可以稳定大于生产 TPS,再计算消化时间:
预计消化时间 = 当前堆积量 / (消费 TPS - 生产 TPS)如果消费 TPS 仍然小于生产 TPS,不要指望“慢慢就好了”。这时必须继续限流、优化下游、隔离失败消息或重新设计并行度。
Java Demo:用有界线程池形成背压
消费端最怕的是 Broker 没堆多少,本机 JVM 里堆了一堆任务。下面的代码用有界队列保护消费者。
ThreadPoolExecutor executor = new ThreadPoolExecutor(
16,
16,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1000),
new ThreadPoolExecutor.CallerRunsPolicy()
);CallerRunsPolicy 的作用是:线程池满了、队列也满了,就让提交任务的线程自己执行。拉取线程被迫变慢,消息更多留在 Broker,而不是无限进入 JVM 内存。
这背后的原则是:
- Broker 适合存消息,JVM 本地内存不适合无限存任务。
- 有界队列能让压力显性化,便于监控和限流。
- 拉取变慢是一种背压,不是单纯“性能差”。
Java Demo:Kafka 本地队列高水位暂停拉取
下面是一个简化思路:线程池队列快满时暂停拉取,恢复后继续拉取。
public class KafkaBackpressureLoop {
private final KafkaConsumer<String, String> consumer;
private final ThreadPoolExecutor executor;
public KafkaBackpressureLoop(KafkaConsumer<String, String> consumer,
ThreadPoolExecutor executor) {
this.consumer = consumer;
this.executor = executor;
}
public void run() {
while (true) {
boolean queueAlmostFull = executor.getQueue().remainingCapacity() < 100;
if (queueAlmostFull) {
consumer.pause(consumer.assignment());
} else {
consumer.resume(consumer.assignment());
}
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(300));
for (ConsumerRecord<String, String> record : records) {
executor.submit(() -> handle(record));
}
}
}
private void handle(ConsumerRecord<String, String> record) {
// 业务真正成功后,再记录可提交的 offset。
}
}真实项目不能在任务刚提交到线程池时就提交 offset。否则线程池里的任务失败了,Kafka 已经认为消息处理完成,会造成业务漏处理。
商业场景:订单同步 ES
链路:
flowchart TD
A["订单服务写 MySQL"] --> B["发送订单变更消息"]
B --> C["同步服务消费"]
C --> D["查询订单快照"]
D --> E["Bulk 写入 ES"]活动高峰时订单消息堆积,扩容同步服务后最初 Lag 下降,随后又增长。常见原因是 ES 写线程池、磁盘 IO、刷新、分片或 bulk 参数成为瓶颈。
处理方式:
- 看 ES bulk 请求耗时和 rejected 数量。
- 调整 bulk 大小,不要单条写,也不要一次过大。
- 降低消费者并发,避免把 ES 写线程池打爆。
- 用订单 ID 做幂等 upsert,允许消息重试。
- 用补偿任务从 MySQL 校验 ES,避免同步丢失。
- 必要时拆分索引、调整分片、优化 mapping 和 refresh 策略。
商业场景:采集任务异步入库
医疗、IoT、日志采集类系统经常把采集结果先写 MQ,再异步清洗、校验、落库。
扩容消费者后再次堆积,常见原因包括:
- 入库表没有合适索引,写入和去重查询变慢。
- 批量落库批次过大,事务太长。
- 同一设备或同一机构数据集中到热点分区。
- 清洗逻辑调用外部接口,外部接口限流。
- 异常数据重复重试,正常数据被挤压。
处理时不要只加消费者,而是要把采集链路拆成:
flowchart TD
A["采集消息"] --> B["格式校验"]
B --> C["清洗转换"]
C --> D["幂等去重"]
D --> E["批量落库"]
E --> F["失败隔离"]每一步都要有耗时指标。只有知道哪一步慢,才能判断该扩消费者、优化 SQL、加索引、拆分热点、还是降级外部接口。
不能做什么
| 错误做法 | 风险 |
|---|---|
| 无限加消费者 | 可能打爆 DB、ES、Redis、HTTP 下游 |
| 直接清空队列 | 业务数据丢失,订单、库存、索引、通知状态不一致 |
| 直接 reset offset 到 latest | 老消息被跳过,后续只能人工补数据 |
| 提前 ACK 或提交 offset | 表面堆积下降,实际业务可能漏处理 |
| 无界线程池排队 | JVM 内存不断上涨,最终 OOM |
| 失败消息立即无限重试 | 形成重试风暴,正常消息也处理不了 |
| 只看消息数量 | 忽略最老消息等待时间和业务 SLA |
预防设计
| 能力 | 为什么需要 |
|---|---|
| 生产 TPS 监控 | 知道入口压力是否异常 |
| 消费 TPS 监控 | 知道处理能力是否下降 |
| Lag 和最老消息延迟 | 判断堆积规模和业务影响 |
| 单条处理耗时 P95/P99 | 判断消费者是否变慢 |
| 下游耗时和连接池监控 | 判断瓶颈是否在 DB、ES、HTTP |
| 重试和死信监控 | 防止失败消息吞噬消费能力 |
| 有界线程池 | 防止 JVM 本机无限堆积 |
| 幂等消费 | 支撑重试、扩容、补偿和重放 |
| 限流和暂停拉取 | 下游慢时保护系统 |
| 容量压测 | 提前知道分区数、消费者数、DB、ES 的上限 |
关联知识点
| 知识点 | 为什么要继续学 |
|---|---|
| 消息堆积与背压 | 系统学习堆积、Lag、Ready、Unacked、背压、重试、死信和补偿 |
| Kafka 分区顺序与 Offset | 理解为什么 Kafka 受分区数量限制 |
| Kafka 可靠性幂等与事务 | 理解 offset、幂等、事务和重复消费 |
| RabbitMQ ACK 重试和死信 | 理解 Ready、Unacked、ACK、NACK、DLX |
| RocketMQ 消费重试 | 理解失败重试、死信和顺序消息阻塞 |
小结
消费者扩容后短暂有效又堆积,本质是“入口并发提高了,但有效消费 TPS 没有持续超过生产 TPS”。排查时不要停留在消费者数量,而要看分区/队列并行度、消费耗时、下游承载、重试风暴、线程池、连接池和 Broker 压力。成熟的处理方式是先止血、再定位最短板、最后消化存量,并用背压、幂等、死信、补偿和容量预案把问题变成可控问题。
