Kafka 消费堆积排查

mq 📚 learning kafka-consumer-lag-troubleshoot · kafka · consumer-lag · troubleshooting · rebalance · kafka-consumer-groups · operations

TL;DR(30 秒扫完)

  • 4 步排查:确认现象 → 分场景定位 → 根因分类 → 分级方案
  • 第一步:kafka-consumer-groups.sh --describe --group 看 Lag 分布
  • 3 种场景:均匀增长 / 集中某分区 / 从 0 突增
  • 3 类根因:Consumer 侧(下游依赖/JVM/业务/配置/网络)+ Producer 侧(突发/重试风暴)+ Broker 集群侧(磁盘/副本/Controller)
  • 3 级方案:P0(不能丢)扩容+限流+下游优化 / P1(可延迟)降级消费 / P2(可丢弃)跳过消息
  • 3 个坑:消费者数不能超分区数 / Rebalance 期间堆积是正常的 / 扩容分区会改变消息路由

关键结论

结论 A排查堆积第一步是看 Lag 分布,不是盲目扩容
结论 B均匀堆积 → 整体消费慢,重点查下游依赖和 JVM;集中堆积 → 特定分区问题,重点查该分区
结论 C扩容消费者前先确认分区数上限,消费者数不能超分区数
结论 D堆积恢复时可能重复消费,消费端必须做幂等

完整讲解(费曼四步)

STEP 1 · 概念

消费堆积是消费者消费速度跟不上生产速度,导致 Lag(Log-End-Offset - Current-Offset)持续增长的现象。

STEP 2 · 大白话

医院排队比喻:
  • Producer = 来挂号的病人
  • Consumer = 医生
  • Lag = 排队人数
堆积 100 万 = 排队 100 万人。可能是:
  • 病人突然增多(Producer 突发)
  • 医生处理变慢(Consumer 慢:下游依赖/JVM/业务逻辑)
  • 医生排班出问题(Rebalance 期间整个组暂停)
  • 医生数量不足(消费者数不够,但不能超过分区数)

STEP 3 · 底层

4 步排查流程

第一步:确认现象(10 分钟内)
# 看 Lag 分布
kafka-consumer-groups.sh \
  --bootstrap-server <broker> \
  --describe \
  --group <group>

# 关键输出
# CURRENT-OFFSET | LOG-END-OFFSET | LAG
观察两点:
  • Lag 分布:均匀 vs 集中在某分区
  • 堆积速度:线性 vs 指数(指数是正反馈问题)
第二步:分场景定位
场景表现排查方向
均匀增长所有分区 Lag 都涨消费者整体处理慢
集中某分区某分区 Lag 暴涨该分区消息/消费逻辑问题
从 0 突增突然从 0 到 100 万Producer 突发 / Consumer 崩溃 / Rebalance
第三步:根因分类 Consumer 侧(最常见):
类别具体原因
下游依赖DB 慢查询、外部 API 超时、Redis 阻塞
JVM 问题FullGC、OOM、CPU 100%
业务逻辑死循环、幂等去重表锁竞争、大事务
配置问题max.poll.records 太大、max.poll.interval.ms 超时
网络问题与 Broker 连接不稳定、DNS 解析失败
Producer 侧:
  • 突发大量消息(大促、批量任务、下游积压补偿)
  • 重试风暴(Consumer 消费失败反复重试)
Kafka 集群侧:
  • Broker 磁盘满、副本同步慢
  • Controller 频繁切换、Leader 频繁迁移
  • ZooKeeper/KRaft 异常
第四步:分级方案
等级场景方案
P0数据不能丢扩容消费者(注意分区数上限)+ 限流生产者 + 下游扩容/优化 + 消费幂等
P1可以延迟降级消费逻辑、临时跳过非核心消息
P2可以丢弃跳过堆积消息(--to-latest 或修改 offset)

3 个关键陷阱

  • 消费者数 > 分区数:多出的消费者空转,扩容无效
  • 解决:先确认分区数,不够再扩分区
  • Rebalance 期间:整个组暂停消费,堆积是正常的(几秒到几十秒)
  • 解决:用 CooperativeStickyAssignor 减少停顿
  • 扩容分区会改变消息路由:影响顺序性
  • 解决:评估业务是否依赖顺序,必要时新建 Topic 迁移

常用排查命令

# 1. 看消费组状态
kafka-consumer-groups.sh --bootstrap-server <broker> \
  --describe --group <group>

# 2. 看消费者客户端信息
kafka-consumer-groups.sh --bootstrap-server <broker> \
  --describe --group <group> --members --verbose

# 3. 看分区详情
kafka-topics.sh --bootstrap-server <broker> \
  --describe --topic <topic>

# 4. 看 Consumer Lag 监控(Prometheus)
kafka_consumergroup_lag{group="xxx",topic="xxx"}

STEP 4 · 简化

一句话总结:堆积排查 4 步——看 Lag 分布、分场景定位、找根因、分级方案。 记忆口诀:
  • 第一步:Lag 分布(均匀 vs 集中)
  • 三场景:均匀 / 集中 / 突增
  • 三根因:Consumer / Producer / Broker
  • 三级方案:P0 扩容+限流 / P1 降级 / P2 跳过
  • 三陷阱:消费者数 ≤ 分区数 / Rebalance 正常堆积 / 扩分区改路由

常见误区

堆积就直接扩容消费者
先看 Lag 分布,集中堆积扩消费者无效
消费者数 > 分区数能提升并行度
消费者数 ≤ 分区数,多出的空转
扩容分区是万能解
分区扩容会改变消息路由,影响顺序性
只说"确认情况"、"检查状态"
排查要具体命令和工具(kafka-consumer-groups.sh、Prometheus)
堆积恢复后直接消费
堆积恢复时可能重复消费,消费端必须幂等

延伸追问

堆积 100 万消息,业务方要求 10 分钟内追平,怎么做?
扩容消费者(注意分区数上限)、优化消费逻辑、跳过非核心消息、降级下游依赖、限流生产者
堆积集中在某一个分区,怎么排查?
检查消息 Key 是否特殊、消息体是否异常大、消费逻辑是否死锁、消费者是否被正确分配到该分区
Rebalance 频繁导致堆积,怎么解决?
用 CooperativeStickyAssignor;调大 session.timeout.ms 和 max.poll.interval.ms;减小 max.poll.records
如何判断是 Consumer 慢还是 Producer 突增?
看 Lag 分布:均匀增长通常是 Consumer 慢;从 0 突增通常是 Producer 突发
堆积恢复后可能重复消费,怎么办?
消费端必须做幂等:DB 唯一索引、Redis 记录、去重表、业务侧去重

速查表

4 步: 确认现象 → 分场景定位 → 根因分类 → 分级方案
命令: kafka-consumer-groups.sh --describe --group <group>
3 场景: 均匀 / 集中 / 突增
3 根因: Consumer / Producer / Broker 集群
3 方案: P0 扩容+限流 / P1 降级 / P2 跳过
3 陷阱: 消费者数 ≤ 分区数 / Rebalance 正常堆积 / 扩分区改路由

Anki 候选卡片

Q: Kafka 消费堆积排查的第一步?
A: kafka-consumer-groups.sh --describe --group 看 Lag 分布
Q: 堆积 3 种场景分别怎么排查?
A: 均匀 → Consumer 整体慢;集中 → 特定分区问题;突增 → Producer 突发或 Consumer 崩溃
Q: 消费者数能超过分区数吗?
A: 不能,多出的消费者空转,扩容无效
Q: 扩容分区有什么副作用?
A: 会改变消息路由,影响顺序性
Q: 堆积恢复后可能重复消费,怎么办?
A: 消费端必须做幂等——DB 唯一索引、Redis 记录、去重表

关联题目

  • ⚠️ 《消费组堆积从 1000 涨到 100 万,如何排查?》— Round 2 Q10, ⭐⭐(只有模糊描述,缺具体命令)
  • 关联知识

    4 步排查:确认现象 → 分场景定位 → 找根因 → 分级方案;关键看 Lag 分布
    医院排队:病人(Producer)、医生(Consumer)、排队人数(Lag)
    ✦ 记 忆 口 诀 ✦
    看 Lag 分布 → 分场景 → 找根因 → 分级方案
    关键可视化
    4 步排查流程
    flowchart TD
      A[第一步 确认现象] --> B[看 Lag 分布 kafka-consumer-groups.sh]
      B --> C[第二步 分场景定位]
      C --> D[均匀增长 Consumer 整体慢]
      C --> E[集中某分区 特定分区问题]
      C --> F[从 0 突增 Producer 突发或 Consumer 崩溃]
      D --> G[第三步 根因分类]
      E --> G
      F --> G
      G --> H[Consumer 侧 下游依赖 JVM 业务 配置 网络]
      G --> I[Producer 侧 突发 重试风暴]
      G --> J[Broker 集群侧 磁盘 副本 Controller]
      H --> K[第四步 分级方案]
      I --> K
      J --> K
      K --> L[P0 扩容 限流 下游优化]
      K --> M[P1 降级消费]
      K --> N[P2 跳过消息]
    3 个关键陷阱
    flowchart TD
      A[陷阱 1 消费者数 超过 分区数] --> B[多出的消费者空转 扩容无效]
      C[陷阱 2 Rebalance 期间] --> D[整个组暂停消费 堆积是正常的 几秒到几十秒]
      E[陷阱 3 扩容分区] --> F[会改变消息路由 影响顺序性]
    知识关系

    ⬆️ 前置(Prerequisite)

    kafka-rebalance

    🔄 延伸(Extends)

    暂无

    ⚡ 对比(Contrast)

    暂无
    🎯 概念 📏 规则 ⚠️ 误区 🔍 追问 ✨ 口诀 共 0 张卡,点击翻面