Appearance
Kafka 积压突然飙升时怎么排查
Kafka 积压是消息系统里最典型的故障征兆之一。
很多团队一看到 lag 飙升,第一反应就是:
- 扩消费者
- 重启消费者
- 调大分区数
这些动作有时有效,但如果不先分清积压根因,很容易把问题越处理越乱。
先说结论
- lag 飙升只是现象,根因通常在生产速度变快、消费速度变慢、消费暂停,或者 rebalance 抖动
- 第一目标是判断“还能不能追上”,而不是立刻做大改
- 第二目标是定位瓶颈在 broker、网络、消费者、下游依赖还是业务逻辑
- 很多积压问题本质不是 Kafka 性能,而是消费端依赖变慢或异常重试放大
更直接一点说,lag 问题最值得先回答三个问题:
- 是谁先变了,生产还是消费
- 现在只是追不上,还是已经基本停住
- 卡点在 Kafka 本身,还是卡在消息处理链路后半段
一个典型故障现场
某支付结果通知 topic 平时 lag 接近 0,突然 10 分钟内堆到几十万:
- 消费者实例数没变
- Broker CPU 正常
- 业务接口开始延迟升高
- 下游 MySQL 更新 RT 也在抖
这时如果只盯 Kafka,很容易错过真正瓶颈。
更常见的误判是:
- 看见 lag 涨了就扩容消费者
- 看见 consumer 在运行就以为它真的在消费
- 看见 Kafka broker 正常就以为消息系统没问题
实际上 lag 是一个非常典型的“放大器指标”,它会把消费链路后面的很多问题一起暴露出来。
第一阶段: 先判断积压类型
1. 生产突增型
生产流量突然暴增,但消费能力没跟上。
典型于大促、批处理任务、上游 bug 重复发送。
2. 消费变慢型
生产速率没明显变化,但消费者处理速度下降。
典型原因是:
- 下游数据库变慢
- 外部 RPC 超时
- 单条消息处理逻辑变重
3. 消费暂停型
消费者看似在线,但实际没在正常推进。
例如:
- 线程池卡死
- rebalance 频繁
- 位点提交异常
- 某个分区处理线程阻塞
先把积压归到这三类之一,排查会快很多。
4. 局部热点型
还有一种特别容易被忽略的情况:
- 总体消费速率还可以
- 但少数分区 lag 飙得特别高
这往往意味着热点 key、热点分区或局部消费实例异常。此时“全组扩容”往往效果有限。
第二阶段: 先做止血动作
1. 确认是否需要降级非核心消费
如果 topic 里有非关键业务,可以先降级非核心处理逻辑,优先把主链路 lag 压住。
2. 保护下游依赖
如果消费者因为 DB / Redis / RPC 变慢导致处理能力下降,盲目提并发只会把下游压垮。
3. 谨慎扩容消费者
扩容前先确认:
- 分区数是否允许并发提升
- 消费逻辑是否无共享瓶颈
- 下游是否还能承受更高并发
4. 必要时隔离失败消息和重试链路
如果是失败消息在反复立即重试,先把重试风暴切开,比继续压主消费链路更有效。
第三阶段: 看关键指标
排 Kafka 积压,我通常会重点看四组指标。
1. 生产速率 vs 消费速率
最重要的问题是:
- 当前消费速度有没有机会追平生产速度
如果生产每秒 5 万,消费每秒只有 2 万,不处理根因就不可能追平。
2. 分区维度 lag 分布
如果是所有分区一起涨,更像整体消费能力不足。
如果只有个别分区涨,往往是:
- 分区键倾斜
- 某个消费者实例异常
- 某个分区对应消息特别重
3. 消费端线程池和处理时长
很多积压其实卡在业务线程池:
- 队列打满
- 活动线程耗尽
- 单条消息处理耗时暴增
4. rebalance 和异常重试
消费者频繁 rebalance 会导致:
- 分区反复撤销和重新分配
- 实际有效消费时间下降
- 位点推进不稳定
5. 分区粒度的消费时长
很多时候“消费者在线、QPS 也有”,但某些分区上的单条消息特别重,实际推进速度极慢。
如果只看消费组总量,很容易把热点分区问题平均掉。
常见根因画像
1. 下游数据库变慢
最常见。
消费者处理本身不重,但更新数据库、查缓存、调 RPC 变慢,导致吞吐迅速下降。
2. 分区倾斜
某个热点 key 大量集中到单分区,导致:
- 整体看 lag 不算太夸张
- 但某个分区特别严重
3. 批处理策略不合理
单批次太大、提交位点过慢、串行处理链路过长,都会让消费端“看起来在线但推进很慢”。
4. 重试风暴
失败消息立即重试,反复打下游,正常消息也被拖慢,积压进一步滚雪球。
5. 发布、配置或线程池变更引发消费退化
很多 lag 故障并不是流量侧变化,而是最近一次发版让:
- 批次变大
- 线程池缩小
- 重试策略更激进
- 单条消息逻辑变重
结果消费者还活着,但有效吞吐明显下降。
一个实用排查顺序
线上 lag 飙升时,我通常按这个顺序看:
- 生产速率和消费速率谁先变了
- 是全部分区都涨,还是个别分区异常
- 消费者线程池、GC、CPU、网络有没有异常
- 下游 DB / Redis / RPC 是否同步变慢
- 最近是否发生 rebalance、发布、配置变更
这个顺序能帮你先判断 Kafka 是根因还是放大器。
如果需要更细一点,我会再加两步:
- 看失败消息比例和重试量有没有同步升高。
- 看是否有单个热点 key、热点分区、热点实例特别突出。
止血后的治理方向
1. 提升消费幂等和可回放能力
这样在需要临时并发补偿、批量重放时更安全。
2. 分区设计更合理
避免热点 key 长期集中在少数分区。
3. 重试链路隔离
失败消息不要和正常消息抢同一条主消费通道。
4. 下游保护
消费者吞吐上不去时,要优先考虑:
- 限流
- 批量化
- 异步落库
- 熔断降级
5. 给 lag 故障准备恢复策略
例如:
- 高峰期临时扩容模板
- 重试 Topic 限流方案
- 补偿消费和回放流程
- 非核心消费者一键降级开关
一个现场判断原则
如果 lag 在涨,但消费者“理论上”已经满负荷:
- 先不要急着怪 Kafka
- 先确认消费端是否真的把时间花在消息拉取上
很多时候,消费者只是把 Kafka 当作了“问题暴露点”,真正慢的是后面的业务链路。
总结
Kafka 积压故障的关键,不是看到 lag 就立刻扩容,而是先分清:
- 是谁变快了
- 是谁变慢了
- 还能不能追平
把“生产、消费、分区、失败重试、下游依赖、rebalance”这几条线捋清楚,积压问题通常就能很快定位到真正根因。真正稳的系统,不是不会积压,而是积压时能快速判断、快速止血、快速恢复。