Kafka 消费堆积:先判断是慢消费还是下游故障
Kafka 消费堆积:先判断是慢消费还是下游故障
一、Lag 上升只是症状
Kafka 消费 lag 上升时,很多人第一反应是加消费者实例。但 lag 上升可能来自消息突增、消费者处理慢、下游数据库故障、分区数不足、rebalance 频繁、单条消息卡住或业务逻辑异常。加实例只对部分场景有效,盲目扩容可能把下游打得更狠。
排查消费堆积,先判断是消费能力不足,还是下游不可用。如果消费者 CPU 很低、处理耗时高、下游错误多,问题不在 Kafka;如果消费者 CPU 打满、下游正常,才可能是消费能力不足。
二、排查链路:Lag、吞吐、耗时一起看
flowchart TD
A[消费 Lag 上升] --> B[看生产速率]
B --> C[看消费速率]
C --> D[看处理耗时]
D --> E[检查下游]
E --> F[扩容或降级]
Lag 要结合生产速率看。活动期间生产速率突然翻倍,短时间 lag 上升可能正常;生产速率恢复后能追上,就不是严重问题。若生产速率正常但 lag 持续上升,说明消费侧或下游有瓶颈。
还要看分区数。Kafka 同一个 consumer group 内,一个分区同一时间只能由一个消费者消费。消费者实例数超过分区数后,再扩容也没用。分区设计是吞吐上限的一部分,不能等堆积时才想起来。
三、监控指标:不要只盯一个 Lag
下面是一组建议指标。它们能帮助判断瓶颈位置。
kafka_consumer_metrics:
- records_lag_max
- records_consumed_rate
- records_processed_latency_p95
- poll_interval_ms
- rebalance_count
- downstream_error_rate
- commit_latency_ms
poll_interval_ms 过长可能触发 rebalance。处理逻辑太慢、单批消息太大或线程阻塞,都可能导致消费者没及时 poll。rebalance 频繁时,消费会反复暂停,lag 更难下降。
提交位点也要谨慎。业务处理成功后再提交,能避免丢消息;但如果单条消息一直失败,会阻塞后续消息。需要死信队列或跳过策略,避免坏消息卡住整个分区。
四、处理策略:扩容、限流和降级要配合
如果瓶颈在消费者 CPU,可以增加实例或优化处理逻辑;如果瓶颈在数据库,要限流或批量写入;如果瓶颈在外部接口,要降级、异步重试或进入死信。策略必须针对瓶颈,不要把所有堆积都当成消费者不够。
批量处理可以提升吞吐,但会增加单批失败成本。要控制 batch size,并记录每条消息处理结果。大批量写库时,还要注意事务时间和锁竞争。
最后,堆积恢复也要保护下游。lag 很大时,消费者追赶会形成高峰,可能把刚恢复的数据库再次打挂。可以限速追赶,优先处理高优先级 topic 或关键消息。恢复阶段也需要架构设计。
还要提前定义告警分级。短时间 lag 上升可以提醒,持续无法追平才需要升级;核心 topic 和低优先级 topic 的阈值也不一样。告警如果不分级,值班人员会被普通波动淹没。Kafka 的稳定性不只在 broker,也在消费侧的运营纪律里。
死信队列要有人看。把失败消息丢进 DLQ 后如果没人处理,只是把问题换了个位置。DLQ 应该有数量告警、重放工具和人工处理流程。
五、总结
Kafka 消费堆积要先判断慢消费、消息突增还是下游故障。Lag 只是症状,生产速率、消费速率、处理耗时、rebalance 和下游错误率才是证据。扩容有用,但不是唯一答案。