这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本章是“急救包”:按症状组织,给出排查路径、关键命令与常见根因。遇到问题时可以直接对照执行。
20.1 排查总原则
1. 先看范围:单 topic?单分区?单消费者组?全集群?
2. 再看时间:什么时候开始?有没有变更(发布/扩容/配置)?
3. 分层定位:客户端 -> 网络 -> Broker -> 下游
4. 保留现场:日志、指标截图、配置快照
5. 最小干预:先恢复(扩容/限流/切流),再根因
万能三件套:
kafka-topics.sh --describe --topic <t>
kafka-consumer-groups.sh --describe --group <g>
kafka-broker-api-versions.sh --bootstrap-server <bs> # 确认集群可达
20.2 消息丢失
排查路径
1. Broker 端确认写入了吗?
kafka-dump-log.sh / kafka-get-offsets.sh 查看分区末位移是否增长
2. 生产端真的发成功了吗?
检查回调/日志;检查是否 acks=0;是否未 flush 退出;
是否异步 send 后进程崩溃
3. 消费端跳过了吗?
查 committed offset 与 log end offset;
查 auto.offset.reset、位移重置记录、rebalance 时间点
4. 数据被删了吗?
retention 是否过短;是否误操作 reset;磁盘清理是否异常
高频根因
| 层 | 根因 | 证据 |
|---|---|---|
| 生产 | send() 后直接退出 |
无回调,日志缺失 |
| 生产 | acks=0 或 1 且 Leader 闪退 | producer 配置 |
| Broker | RF=1 + 磁盘损坏 | topic describe |
| Broker | unclean election | controller.log |
| 消费 | 先 commit 后 process | 代码审查 |
| 消费 | reset 到 latest/未来 offset | OffsetReset 日志 |
20.3 消息重复
排查路径
1. 判断重复发生在哪一层:
- broker 日志里物理上两条? -> 生产端重试/业务重发
- 日志一条但业务处理两次? -> 消费端重复消费(rebalance/提交失败)
2. 生产端:
关闭幂等?transactional.id 不稳定?业务层超时后自己重发?
3. 消费端:
处理完但 commit 失败?再平衡导致重复处理?
异步提交丢失(commitAsync 未等待确认)?
4. 兜底:
下游是否有幂等键(requestId / offset / 唯一约束)?
处置
- 物理重复无法事后删除,只能下游去重或重放修正;
- 预防:稳定 transactional.id、手动同步提交、幂等键设计;
- 对账任务:按业务主键统计多处理记录,发现重复模式。
20.4 消息堆积(Lag 告警)
定位瓶颈
1. 写入速率是否突增?
MessagesInPerSec、BytesInPerSec 按 topic 查看
2. 消费速率是否下降?
消费者处理耗时日志、下游 DB/接口指标
3. 消费者数量是否减少?
rebalance 记录、实例存活、空闲消费者
4. 是否分区倾斜?
按 partition 看 lag;某分区 key 热点或数据倾斜
恢复手段
| 手段 | 适用 |
|---|---|
| 扩消费者实例 | 实例数 < 分区数 |
| 批量化下游写库 | 单条处理慢 |
| 降级非核心逻辑 | 处理链路太重 |
| 临时扩分区 + 消费者 | 长期容量不足(注意 key 顺序影响) |
| 跳过历史(慎用) | 业务允许丢弃过期数据 |
跳过堆积只适合明确可丢弃的数据(如过期告警),执行前必须业务确认并留审计。
20.5 消息乱序
常见原因:
- 未用相同 key,或 key 设计不对(同一实体被路由到不同分区);
- 生产端重试导致乱序(未开幂等,in-flight > 1);
- 消费端多线程/并发处理同一分区的消息;
- 增加分区导致 key 落点变化;
- 上游本身乱序产生(多线程写库后发消息)。
排查:
# 消费时打印分区与 offset,确认同一实体的消息是否同分区连续
kafka-console-consumer.sh ... --property print.key=true \
--property print.partition=true --property print.offset=true
修复:
- key = 需要顺序的实体 ID;
- 开启幂等 + in-flight <= 5;
- 消费端按 key 哈希到固定工作线程;
- 版本号/时间戳拒绝旧事件。
20.6 ISR 频繁收缩
现象:UnderReplicatedPartitions 抖动,IsrShrinksPerSec 高。
排查顺序:
- Follower 所在 Broker 资源:CPU、磁盘 util、GC;
num.replica.fetchers是否不足;- 网络是否抖动(Broker 间延迟);
- 是否有大分区迁移/均衡占满 IO;
replica.lag.time.max.ms是否设得过小。
处置:加大 fetcher 线程、限流迁移任务、错峰均衡、修复慢盘。
20.7 Rebalance 风暴
现象:消费反复停顿,日志出现 heartbeat failed / max poll interval / rebalancing。
排查:
1. 找出被踢的成员与时间点
2. 对应实例 GC 日志 / CPU / 处理耗时
3. 检查订阅是否一致(同一组订阅不同 topic 集合)
4. 检查是否有实例频繁发布(无静态成员)
处置参考第 15 章治理清单:调小单次处理量、放宽超时、静态成员、协作式策略、统一订阅。
20.8 磁盘满与日志损坏
磁盘满
应急:
1. 扩容磁盘(云盘在线扩)
2. 缩短低价值 topic 的 retention
3. 删除已确认可删的僵尸 topic
长期:
1. 容量水位告警(80%)
2. topic 生命周期治理
3. 压缩策略启用评估
注意:磁盘 100% 时 Broker 可能无法写日志/元数据,处置优先级最高。
段文件损坏
症状:启动报 CorruptIndexException、Unkown magic value 等。
处理:
- 停 Broker,备份整个分区目录;
- 删除对应
.index/.timeindex(会自动重建); - 若
.log损坏,评估删除损坏 segment 的数据损失; - 依赖副本恢复:把损坏副本下线,让其他副本重新同步。
20.9 速查表
| 症状 | 第一条命令 |
|---|---|
| 集群不可达 | kafka-broker-api-versions.sh |
| 分区不可用 | kafka-topics.sh --describe --unavailable-partitions |
| 副本落后 | kafka-topics.sh --describe --under-replicated-partitions |
| 消费堆积 | kafka-consumer-groups.sh --describe --group |
| 消息没写入 | kafka-get-offsets.sh --topic + dump-log |
| KRaft 异常 | kafka-metadata-quorum.sh describe --status |
| 配置疑议 | kafka-configs.sh --describe --all |
本章小结
- 排查先定范围与时间线,再分层定位,保留现场;
- 丢失要区分“没写入”与“没消费”,重复要区分“物理重复”与“逻辑重复”;
- 堆积先看写入增速、消费降速、消费者数量与分区倾斜;
- 乱序通常出在 key 设计、生产端重试与消费端并发;
- 磁盘满与段损坏是最高优先级的基础设施故障。
思考题
- 业务说“丢消息”,但你查到日志里有这条记录,问题出在哪个环节?下一步查什么?
- 为什么删除损坏的
.index文件相对安全,删除.log却有真实数据损失? - 为你们团队写一份“Kafka 故障应急预案”,包含哪些角色和动作?