Kafka Consumer Lag 持续上升,说明消息写入速度超过消费处理能力,或消费者已经异常停止。积压本身不是根因,排查时必须同时观察生产速率、消费速率、消费者状态和下游依赖。
一、确认积压范围
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
--describe --group order-consumer重点查看每个分区的 CURRENT-OFFSET、LOG-END-OFFSET、LAG 和消费者实例。判断是全部分区同步积压,还是少数分区热点。
同时确认:
消费者实例是否在线;
Lag 是持续增长还是短时波动;
生产流量是否突增;
消费耗时是否增加;
是否频繁发生 rebalance;
下游数据库、缓存或接口是否变慢。
二、常见根因
消费逻辑变慢,例如数据库慢 SQL、外部接口超时。
消费者异常退出、线程阻塞或频繁 Full GC。
单条异常消息反复重试,阻塞整个分区。
分区数不足,消费者数量已经超过可并行分区数。
Key 分布不均,流量集中在少数热点分区。
max.poll.interval.ms、批量大小或提交策略不合理,引发 rebalance。Broker 磁盘、网络或副本同步异常,读写延迟升高。
三、应急处理
先修复崩溃或卡死的消费者,恢复基本消费能力。
在分区数允许时横向增加消费者实例。
对下游依赖限流、批量写入或临时降级非核心处理。
将无法处理的异常消息转入死信队列,避免无限重试阻塞。
评估消息保留时间,确保积压未超过可恢复窗口。
重置 Offset 会跳过或重复消费消息,不能把它当作普通清理命令。执行前必须确认业务幂等、补偿方案和数据影响。
四、判断扩容是否有效
消费者并行度上限通常受分区数限制。一个消费组内,同一分区在同一时刻只能分配给一个消费者。若 Topic 只有 6 个分区,增加到 20 个消费者并不会获得 20 倍并发。
扩容后应观察消费速率是否超过生产速率,并估算清空积压所需时间:
预计恢复时间 = 当前积压量 /(消费速率 - 生产速率)如果消费速率没有提升,应继续定位热点分区、下游瓶颈或 Broker 性能问题。
五、长期治理
监控 Lag 总量、增长速率、最老消息年龄和消费成功率。
为消费者建立重试、死信、幂等和超时机制。
根据峰值流量规划分区数,并避免不均衡的 Key 设计。
对消费耗时做分阶段指标,区分拉取、反序列化、业务处理和下游写入。
发布新版本时使用灰度和稳定的滚动策略,减少集体 rebalance。
定期演练消费者宕机、下游变慢和流量突增场景。
治理 Consumer Lag 的目标不是让指标永远为零,而是让系统能识别积压原因、快速恢复,并保证消息在重试和扩容过程中不丢失、不重复破坏业务。