Skip to content

消费者扩容后再次堆积

这个知识点专门解释一个生产里很常见、也很容易误判的现象:

MQ 已经堆积了,临时把消费者实例加上去,Lag 或队列深度一开始下降,过一会儿又开始堆积。

这不是 MQ “不稳定”,也不一定是消费者“没有扩成功”。它通常说明:扩容消费者只提升了消费入口的并发,但整条链路真正的瓶颈不在消费者数量,或者扩容后把新的瓶颈打出来了。

学习目标

目标需要掌握的内容
知道现象明白为什么扩容初期有效,后面又重新堆积
知道原理理解 MQ 消费吞吐由分区/队列、消费者、线程池、下游、ACK、重试共同决定
会定位能通过 Lag 分布、消费 TPS、处理耗时、错误率、重试量、线程池、连接池、慢 SQL 判断瓶颈
会处理会选择限流、降并发、批处理、失败隔离、增加分区、优化下游、补偿重放等措施
会预防能提前设计监控、背压、幂等、死信、容量预案和故障演练

先用一句话理解

消费者扩容能不能解决堆积,取决于扩容后有效消费 TPS是否持续大于生产 TPS。

text
堆积是否下降 = 有效消费 TPS > 生产 TPS

但有效消费 TPS 不是只由消费者实例数决定,而是由整条链路的最短板决定:

text
有效消费 TPS =
min(
  Broker 投递能力,
  分区/队列并行度,
  消费者本机处理能力,
  下游数据库/接口/ES 能力,
  ACK 或 Offset 提交能力,
  重试与异常消息处理能力
)

所以扩容消费者只是提高了其中一项。如果最短板在数据库、远程接口、ES、Redis、分区数量、热点 key、重试风暴或 Broker,继续加消费者不会根治问题。

为什么会短暂有效

先看一个典型时间线。

mermaid
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、调用三方接口的并发也会变多。下游一旦开始排队、锁等待、限流、拒绝请求,单条消息处理耗时会升高。

例如原来:

text
10 个消费者,每条处理 20ms,消费 TPS 约 500

扩容后短时间:

text
30 个消费者,每条仍然 20ms,消费 TPS 约 1500

但数据库被打满后:

text
30 个消费者,每条处理变成 200ms,消费 TPS 约 150

这时消费者更多,反而更容易造成连接池等待、锁竞争、慢 SQL、超时重试,堆积会重新增长。

3. 并行度被分区或队列限制

MQ 并不是消费者越多吞吐越高。

产品并行度限制
Kafka同一个 Consumer Group 中,一个 Partition 同一时刻只能分配给一个消费者
RocketMQ一个 MessageQueue 同一时刻通常只会被同组一个消费者消费,顺序消息更明显
RabbitMQ同一个 Queue 可以多个消费者竞争,但会受 prefetch、ACK、Unacked、下游能力影响

如果 Kafka 只有 8 个分区,启动 20 个同组消费者,真正能并行消费的最多也就是 8 个分区,剩余消费者空闲。

mermaid
flowchart TD
    A["8 个 Partition"] --> B["最多 8 个活跃消费者"]
    B --> C{"启动 20 个消费者"}
    C -- "8 个有任务" --> D["参与消费"]
    C -- "12 个空闲" --> E["不能提升吞吐"]

本质:最短板决定吞吐

可以把 MQ 消费链路看成一条流水线。

mermaid
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、堆积量和最老消息延迟。

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

如果扩容后:

text
生产 TPS = 3000/s
消费 TPS = 4200/s

说明存量应该下降,只需要计算多久追平。

如果过一会儿变成:

text
生产 TPS = 3000/s
消费 TPS = 1800/s

说明新的瓶颈已经出现,继续堆积是必然结果。

扩容反弹的指标时间线

生产里最好不要只看一个时间点,要看扩容前、扩容刚完成、扩容 5 到 10 分钟后、扩容 30 分钟后的指标变化。

mermaid
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成功 TPSP95 耗时错误率Lag 趋势判断
T0 扩容前3000/s1800/s1700/s80ms1%上升消费追不上
T1 刚扩容3000/s5000/s4200/s90ms1%下降短暂有效
T2 5 分钟后3000/s5200/s2600/s350ms8%持平转升下游开始慢或失败增加
T3 15 分钟后3000/s5200/s1500/s900ms20%快速上升新瓶颈完全暴露

这张表的关键是区分:

指标变化说明
拉取 TPS 高,成功 TPS 低消息已经进消费者,但业务没真正处理完
P95/P99 升高下游、锁、线程池或慢消息导致长尾
错误率升高失败消息开始吞噬处理能力
Lag 下降但最老等待时间不降可能只处理了新消息,旧分区或慢消息仍卡住
消费者数增加但拉取 TPS 不涨分区/队列并行度、Broker 或消费者配置限制

所以面试里回答“扩容后又堆积”时,不要只说“下游瓶颈”。更完整的说法是:

我会先拉出扩容前后的生产 TPS、拉取 TPS、成功 TPS、P95/P99、错误率、重试量、Lag 分布和最老消息等待时间。拉取 TPS 上升但成功 TPS 下降,说明瓶颈在消费者内部或下游;消费者数增加但拉取 TPS 不涨,说明并行度或 Broker 投递受限;错误率和 Retry 增长,说明重试风暴;少数分区 Lag 高,说明热点、慢消息或局部消费者异常。

根因决策树

下面这棵树用于生产现场快速定位,不要跳步骤。

mermaid
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 数量变多”就结束,要确认消费者真的加入消费组或队列消费。

产品看什么
KafkaConsumer Group 成员数、分区分配结果、Rebalance 日志
RocketMQConsumerGroup 在线实例、MessageQueue 分配
RabbitMQQueue 的 Consumers 数、连接和 Channel 状态

常见假扩容:

  1. 新实例启动失败。
  2. 配错 ConsumerGroup。
  3. 订阅 Topic 或 tag 不一致。
  4. 网络不通,连不上 Broker。
  5. Kafka Rebalance 后新实例没有分到分区。

判断二:拉取 TPS 是否上升

如果消费者上线了,但拉取 TPS 没上升,说明消息没有更快进入消费者。

可能原因:

原因证据处理
分区/队列数不足活跃消费者数小于实例数增加分区/队列或提升单消费者能力
Broker 投递慢Broker 请求耗时、网络、磁盘高查 Broker 负载
消费者配置限制poll/batch/prefetch 太小调整批量和预取
Rebalance 抖动消费组频繁重新分配查心跳和处理时长

判断三:成功 TPS 是否上升

如果拉取 TPS 上升,但成功 TPS 没上升,说明消息卡在消费者内部或下游。

mermaid
flowchart TD
    A["拉取 TPS 上升"] --> B["消息进入消费者"]
    B --> C{"成功 TPS 不升"}
    C --> D["线程池排队"]
    C --> E["DB/ES/HTTP 慢"]
    C --> F["业务锁等待"]
    C --> G["失败重试"]
    C --> H["ACK/Offset 提交慢"]

这时最重要的是看:

  1. 消费者本地线程池队列。
  2. 单条处理耗时 P95/P99。
  3. DB 连接池 active 和等待时间。
  4. 慢 SQL、锁等待、ES rejected、HTTP 超时。
  5. 错误率、Retry、DLQ。

典型根因全过程

根因一:分区或队列并行度不足

过程:

mermaid
flowchart TD
    A["Topic 只有 8 个分区"] --> B["扩到 20 个消费者"]
    B --> C["最多 8 个消费者分到分区"]
    C --> D["12 个消费者空闲"]
    D --> E["成功 TPS 不明显提升"]
    E --> F["Lag 继续增长"]

判断证据:

证据说明
消费者实例数大于分区数多出来的实例没有分区
每个活跃消费者 CPU 不高不是机器算力不足
拉取 TPS 不随实例数增长并行通道限制

处理方式:

  1. 增加分区或队列,但要评估顺序语义。
  2. 提高单分区处理能力,例如批处理。
  3. 拆 Topic,把热点业务拆出去。
  4. 调整路由 key,避免数据过度集中。

根因二:下游数据库被打满

过程:

mermaid
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,不是算力不足

处理方式:

  1. 降低消费者并发,先保护 DB。
  2. 优化 SQL 和索引。
  3. 批量写入,减少单条网络往返。
  4. 缩小事务范围,避免长事务。
  5. 对热点数据做分片或串行化处理。

根因三:重试风暴

过程:

mermaid
flowchart TD
    A["某类消息持续失败"] --> B["进入快速重试"]
    B --> C["扩容后失败消息被更快取出"]
    C --> D["消费者大量时间处理失败消息"]
    D --> E["正常消息得不到处理"]
    E --> F["Lag 再次上升"]

判断证据:

证据说明
错误率升高不是纯容量问题
Retry Topic 增长失败消息在反复消耗资源
同一异常反复出现参数错误、数据错误、代码 bug
死信队列开始增长多次重试仍失败

处理方式:

  1. 区分短暂故障和永久故障。
  2. 永久故障不要无限重试,进入死信或异常表。
  3. 延迟重试,避免立即打爆消费者。
  4. 修复毒消息或数据问题后定向补偿。

根因四:热点 key 或慢消息

过程:

mermaid
flowchart TD
    A["热点订单/商户/设备 key"] --> B["持续路由到同一分区"]
    B --> C["该分区消息远多于其他分区"]
    C --> D["负责消费者长期处理同一通道"]
    D --> E["少数分区 Lag 很高"]
    E --> F["总 Lag 看起来也很高"]

判断证据:

证据说明
Lag 集中在少数分区不是整体消费慢
最大分区 key 分布集中热点 key
某条消息耗时特别长慢消息
其他分区很快归零局部阻塞

处理方式:

  1. 找最大 Lag 分区。
  2. 找负责该分区的消费者实例。
  3. 查该消费者日志和线程栈。
  4. 统计该分区消息 key 分布。
  5. 对热点 key 拆分、散列或单独 Topic。
  6. 对慢消息隔离处理,避免卡住后续消息。

根因五:Rebalance 抖动

过程:

mermaid
flowchart TD
    A["扩容消费者"] --> B["触发 Rebalance"]
    B --> C["分区暂停并重新分配"]
    C --> D["部分消费者处理时间过长"]
    D --> E["心跳或 poll 超时"]
    E --> F["再次 Rebalance"]
    F --> G["消费吞吐上下波动"]
    G --> H["Lag 反复下降又上升"]

判断证据:

证据说明
消费组频繁 Rebalance组不稳定
消费者日志有 poll 超时单次处理太久
Lag 周期性下降又暴涨分区频繁暂停
新实例不断加入退出部署或健康检查异常

处理方式:

  1. 减小单次 poll 的处理量。
  2. 把业务处理和 poll 心跳解耦。
  3. 调整 max.poll.interval.ms,但不要只靠调大掩盖慢处理。
  4. 扩容分批进行,避免一次性大量实例加入。
  5. 检查容器健康检查和 OOM。

排查流程

mermaid
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 为例:

text
单分区 Lag = LOG-END-OFFSET - CURRENT-OFFSET
消费组总 Lag = 所有分区 Lag 之和
Lag 分布 = 每个 Partition 分别落后多少

例如总 Lag 都是 9000,下面两种情况的含义完全不同。

1. 均匀分布

text
Partition 0 Lag = 3000
Partition 1 Lag = 3000
Partition 2 Lag = 3000

这说明多数分区都在落后,更像是整体消费能力不足,或者消费者共同依赖的下游数据库、ES、Redis、远程接口变慢。

2. 倾斜分布

text
Partition 0 Lag = 100
Partition 1 Lag = 8800
Partition 2 Lag = 100

这说明堆积集中在少数分区,更像是热点 key、单分区慢消息、某个消费者异常、顺序消息阻塞或数据分布不均

mermaid
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 连接池 activeDB 连接是否够用active 长期打满说明大量线程等连接
慢 SQL / 锁等待DB 是否是短板并发写入导致锁竞争或索引缺失
ES rejectedES 是否拒绝写入Bulk 太大、并发太高、分片压力过大
Broker 磁盘和网络Broker 是否健康磁盘高水位、网络高、拉取延迟高

Kafka 场景怎么判断

Kafka 最常见的误区是:消费者实例数超过分区数后,继续扩容没有意义。

1. 看每个分区 Lag

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

如果只有一个分区 Lag 特别高:

text
PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
0          1000            1200            200
1          2000            9000            7000
2          3000            3100            100

这通常不是整体消费者不足,而是分区 1 有热点 key、慢消息、异常消费者或该分区下游处理慢。

2. 看消费者数量是否超过分区数

text
分区数 = 8
消费者实例数 = 20
有效并行消费者最多 = 8

如果要继续提升吞吐,需要考虑:

  1. 增加 Topic 分区数,但要评估 key 路由和顺序语义。
  2. 提高单分区处理能力,例如批处理、减少单条耗时。
  3. 拆 Topic,把热点业务单独拆出去。
  4. 调整 key,避免热点集中到单个分区。

3. 看 Rebalance 是否频繁

Kafka 扩容、缩容、消费者心跳异常、max.poll.interval.ms 超时,都可能触发 Rebalance。Rebalance 期间分区重新分配,消费会暂停或抖动。

mermaid
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 要重点看 ReadyUnacked

现象重点判断
Ready 很高,Unacked 不高消费者没拿到消息,可能消费者少、连接异常、prefetch 太小
Unacked 很高消费者已经拿到消息但处理慢,可能业务卡住、ACK 慢、prefetch 太大
Ready 和 Unacked 都高生产太快,消费端和下游整体扛不住
Consumers 下降消费者掉线或部署异常

prefetch 控制每个消费者最多同时持有多少未 ACK 消息。

java
channel.basicQos(50);

如果 prefetch 太大,单个慢消费者会拿走很多消息,导致其他消费者没有机会处理;如果 prefetch 太小,吞吐上不去。它不是越大越好,而是要和单条处理耗时、消费者数量、下游能力一起估算。

RocketMQ 场景怎么判断

RocketMQ 要关注 Topic 下每个 MessageQueue 的消费进度、ConsumerGroup Diff、重试队列和死信队列。

现象可能原因
某几个队列 Diff 特别高热点 key、顺序消息阻塞、单队列消费者慢
Retry Topic 增长消费失败反复重试
DLQ 增长多次失败后进入死信,需要人工或补偿处理
扩容后吞吐不涨MessageQueue 数量不足或下游瓶颈
顺序消息堆积某条消息失败导致同队列后续消息不能继续处理

顺序消息尤其要小心。为了保证同一个业务 key 的顺序,消息会进入同一个队列。如果某条消息处理失败,后面的消息不能随便越过它,否则顺序语义就被破坏。

处理策略

1. 先止血

止血不是为了马上清空堆积,而是让堆积不再继续恶化。

场景止血动作
上游生产过快对生产接口限流,暂停补数据、重放、批导任务
下游 DB 慢降低消费并发,减少事务范围,暂停非核心写入
ES 写入慢调整 bulk 大小,降低并发,查看 rejected,必要时暂停非核心同步
远程接口慢设置超时,熔断降级,降低并发,转异步补偿
重试风暴把不可恢复异常进入死信或异常表,不要立即无限重试

2. 再定位最短板

按照这个顺序会比较稳:

  1. 生产 TPS 是否异常升高。
  2. 消费 TPS 是否真的提升。
  3. 单条消费耗时是否升高。
  4. Lag 是否集中在少数分区/队列。
  5. 失败率、重试量、死信是否升高。
  6. DB、ES、Redis、HTTP、连接池、线程池是否打满。
  7. Broker 磁盘、网络、拉取延迟是否异常。

3. 最后消化存量

如果已经确认消费 TPS 可以稳定大于生产 TPS,再计算消化时间:

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

如果消费 TPS 仍然小于生产 TPS,不要指望“慢慢就好了”。这时必须继续限流、优化下游、隔离失败消息或重新设计并行度。

Java Demo:用有界线程池形成背压

消费端最怕的是 Broker 没堆多少,本机 JVM 里堆了一堆任务。下面的代码用有界队列保护消费者。

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

CallerRunsPolicy 的作用是:线程池满了、队列也满了,就让提交任务的线程自己执行。拉取线程被迫变慢,消息更多留在 Broker,而不是无限进入 JVM 内存。

这背后的原则是:

  1. Broker 适合存消息,JVM 本地内存不适合无限存任务。
  2. 有界队列能让压力显性化,便于监控和限流。
  3. 拉取变慢是一种背压,不是单纯“性能差”。

Java Demo:Kafka 本地队列高水位暂停拉取

下面是一个简化思路:线程池队列快满时暂停拉取,恢复后继续拉取。

java
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

链路:

mermaid
flowchart TD
    A["订单服务写 MySQL"] --> B["发送订单变更消息"]
    B --> C["同步服务消费"]
    C --> D["查询订单快照"]
    D --> E["Bulk 写入 ES"]

活动高峰时订单消息堆积,扩容同步服务后最初 Lag 下降,随后又增长。常见原因是 ES 写线程池、磁盘 IO、刷新、分片或 bulk 参数成为瓶颈。

处理方式:

  1. 看 ES bulk 请求耗时和 rejected 数量。
  2. 调整 bulk 大小,不要单条写,也不要一次过大。
  3. 降低消费者并发,避免把 ES 写线程池打爆。
  4. 用订单 ID 做幂等 upsert,允许消息重试。
  5. 用补偿任务从 MySQL 校验 ES,避免同步丢失。
  6. 必要时拆分索引、调整分片、优化 mapping 和 refresh 策略。

商业场景:采集任务异步入库

医疗、IoT、日志采集类系统经常把采集结果先写 MQ,再异步清洗、校验、落库。

扩容消费者后再次堆积,常见原因包括:

  1. 入库表没有合适索引,写入和去重查询变慢。
  2. 批量落库批次过大,事务太长。
  3. 同一设备或同一机构数据集中到热点分区。
  4. 清洗逻辑调用外部接口,外部接口限流。
  5. 异常数据重复重试,正常数据被挤压。

处理时不要只加消费者,而是要把采集链路拆成:

mermaid
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 压力。成熟的处理方式是先止血、再定位最短板、最后消化存量,并用背压、幂等、死信、补偿和容量预案把问题变成可控问题。