主题
附录 C:Kafka 踩坑案例总集
为什么要学踩坑案例?
学会一个分布式系统的「正确姿势」翻官方文档就够了;真正让你和新手拉开差距的,是踩过的坑。本集收录 Kafka 在生产中 25 + 个真实高频坑,覆盖设计 / Producer / Consumer / Rebalance / 副本 / Controller-KRaft / EOS / Compaction / 运维 / 跨机房 / Connect-Streams 十一大主题。
每条案例采用统一 7 段模板:
症状 → 现场快照 → 根因分析 → 解决方案 → 预防措施 → 相关参数 → 关联章节链接
建议先扫一眼下方目录与「踩坑高频度雷达图」,遇问题再精读对应条目;最后把「上线 Checklist」打印出来挂工位。
踩坑高频度 / 严重度雷达图(ASCII)
高频度(出现概率)
▲
严重度 ↑ │
┌───────────────────┼─────────────────────────┐
│ │ │
│ ★★★★★ unclean=true(数据丢) ★★★★★ Rebalance 风暴 │ ← 高频高危
│ ★★★★ acks=1 丢消息 ★★★★ 自动提交丢消息 │
│ ★★★★ ISR 抖 + min.isr 写阻塞 ★★★ 分区数选少了 │
│ ★★★ Compaction 不生效 ★★★ key 热点 │
│ ★★ EOS 用错拖垮吞吐 ★★ 大消息撑爆 │
│ ★ 双 transactional.id 同时跑 ★ Connect Offset 漂 │
│ │ │
└───────────────────┼─────────────────────────┘
│ 低 ─────────────→ 高
└─────► 高频度
📌 优先攻克「右上角」案例:高频 × 高危。目录
- 设计类(Topic / 分区 / 副本 / Key)
- Producer 端
- Consumer 端
- Rebalance 风暴
- 副本与 ISR
- Controller / KRaft
- EOS / 事务
- Compaction
- 运维
- 跨机房
- Connect / Streams / ksqlDB
一、设计类
Topic / 分区 / 副本 / Key 是 Kafka 集群的「四个基本面」,任意一项设计错了,后期改动代价都极高。
1.1 分区数选少了,扩容时 Key Hash 全打乱
症状
- 业务上线初期分区数选了 3,半年后流量爆炸到单分区 60MB/s,下游消费跟不上。
- 想加分区,运维
kafka-topics --alter --partitions 12后,之前依赖「同 Key → 同分区」的所有业务全乱了:同一个用户的事件被打散到不同分区,顺序消费失效,下游聚合出错。
现场快照
bash
$ kafka-topics.sh --bootstrap-server :9092 --describe --topic orders
Topic: orders PartitionCount: 12 ReplicationFactor: 3
# 但 user_id=1001 之前在 partition=2,现在变成了 partition=7
# 旧消息还在 partition=2,新消息进 partition=7 → 消费时序错乱根因分析
- Kafka 的默认分区器是 Murmur2(key) % numPartitions,分区数一变 → Hash 分布全变。
- Kafka 不像 Redis Cluster 那样支持「一致性 Hash + 重新分布」,加分区 = 破坏 Key 路由。
解决方案
- 方案 A:业务能容忍短期乱序,把「Key → 旧分区」的映射作为兼容层,新消息双写新旧两个 Topic,老消息消费完后切换。
- 方案 B:建一个新 Topic
orders_v2(分区数预留 5 ~ 10 倍流量),业务双写老 + 新;新 Topic 跑稳后切流,老 Topic 下线。 - 方案 C:自定义 Partitioner,写「分区数感知」逻辑(不推荐,复杂且脆弱)。
预防措施
- 上线前按 1 ~ 2 年的峰值流量预留分区,不要按当前流量。
- 经验公式:
目标分区数 = max(目标吞吐 / 单分区吞吐, 目标消费并行度),加 30% 余量。 - 单分区写吞吐保守按 10MB/s,消费按业务而定。
相关参数:num.partitions(Broker 默认)、自定义 partitioner.class
关联章节:Ch 06 Topic 设计与分区策略
1.2 副本数选少了,单 Broker 故障即数据丢
症状
- 测试环境用
replication-factor=1跑得好好的,上生产没改 → 一台 Broker 故障,那个 Broker 上所有分区直接没 Leader,业务读写全报错。
现场快照
bash
$ kafka-topics.sh --bootstrap-server :9092 --describe --unavailable-partitions
Topic: orders Partition: 5 Leader: -1 Replicas: 2 Isr:根因分析
replication-factor=1意味着只有一份数据,对应 Broker 挂掉就完蛋。- 内部 Topic(
__consumer_offsets/__transaction_state)默认副本因子读的是offsets.topic.replication.factor/transaction.state.log.replication.factor,如果集群启动时 Broker 数 < 该值,会被自动降级到当时的 Broker 数,从此永远不会自动提升。
解决方案
- 业务 Topic 重建:先双写新副本数 Topic 再切流。
- 内部 Topic:
kafka-reassign-partitions.sh显式扩副本到 3。
预防措施
- 业务 Topic 默认
--replication-factor 3。 - Broker 配置硬性约束:
default.replication.factor=3+offsets.topic.replication.factor=3+transaction.state.log.replication.factor=3。 - 集群启动时必须 Broker 数 ≥ 3,否则内部 Topic 副本数被锁死。
相关参数:default.replication.factor / min.insync.replicas / offsets.topic.replication.factor
关联章节:Ch 06 Topic 设计、Ch 09 副本与 ISR
1.3 Key 设计不当 → 热点 Partition
症状
- 6 个分区,监控发现 partition-0 的 BytesIn 是其他分区的 20 倍,磁盘 IO 跑满;其他 5 个分区闲得发慌。
- 消费侧,对应 partition-0 的 consumer 实例 Lag 飙升,其他实例没事干。
现场快照
partition-0: ▓▓▓▓▓▓▓▓▓▓▓▓▓▓▓▓▓▓▓▓ 200 MB/s
partition-1: ▓ 10 MB/s
partition-2: ▓ 10 MB/s
...根因分析
- Key 选择高度倾斜,比如:用
app_id当 Key,但流量 90% 集中在app_id=top1。 - 或者把
null当 Key 配 SticktPartitioner,整段时间消息粘在同一个分区。
解决方案
- 离散化 Key:在原 Key 上加随机后缀,如
user_id + ':' + (event_id % 16),把热点拆成 16 份再聚合。 - 复合 Key:
app_id + region + minute。 - 改用 RoundRobin Partitioner(牺牲顺序换均匀)。
预防措施
- 上线前对 Key 分布做离线统计(
SELECT key, count(*) FROM events GROUP BY key ORDER BY count DESC LIMIT 100)。 - 监控接入:
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec,topic=*各分区维度看分布。
相关参数:自定义 partitioner.class / partitioner.ignore.keys
关联章节:Ch 06 Topic 设计
1.4 Topic 命名混乱 / 一锅炖
症状
- 半年后
kafka-topics --list输出几百个名字:test_xxx、tmp_xxx、backup_orders_2023_07、orders、orders_v2、orders_new、orders_real,没人敢删。 - 排障时根本不知道哪个 Topic 还在用、哪个没人消费。
根因:没有命名规范。
解决方案
- 强制命名规范:
<env>.<biz>.<event>.<schema_version>,例如prod.order.created.v1。 - 过期 Topic 清理:建一个 cron 跑
kafka-consumer-groups --describe,对长期 Lag 不变的 Topic 标记为_TO_DELETE_<日期>通知业务方,30 天后清。
预防措施
- 建 Topic 必须走平台 / 工单(Kafka Manager / Conduktor / 自研平台)。
- 建 Topic 时强制要求填责任人 / TTL / 消费方。
相关参数:auto.create.topics.enable=false(必须关掉)
关联章节:Ch 06 Topic 设计、Ch 19 可观测性与运维
1.5 Topic 数量爆炸(百万 Topic 综合症)
症状
- ZK 时代:Broker 启动慢到几十分钟,Controller 切换 5 分钟拉不回来,元数据广播打满网络。
- KRaft 时代:好很多,但单个 Broker 挂载几万分区时打开文件数撑爆,segment 滚动慢。
根因
- 业务设计「一个客户一个 Topic」「一个设备一个 Topic」,IoT / SaaS 多租户场景非常常见。
解决方案
- 聚合 Topic:用 1 个 Topic + Key 维度区分租户,分区数足够即可。
- 多租户用前缀:
tenant.<id>.events,配合 ACL prefix 授权。 - 如果一定要多 Topic:升级 KRaft、控制 Broker 单实例分区数 ≤ 4000(KRaft 后可到上万,但仍需测试)。
预防措施
- 上线前把「Topic 数量预估」放进容量评估清单。
- 单个 Broker 分区数(含副本)告警阈值 4000(ZK 时代),10000(KRaft 时代)。
相关参数:num.io.threads / num.network.threads / OS 层 nofile、nproc
关联章节:Ch 02 架构总览、Ch 10 Controller 与 KRaft
二、Producer 端
2.1 acks=1 在 Leader 切换瞬间丢消息
症状
- 生产业务
acks=1(Leader 写入即返回 OK),运行半年都正常。 - 某次 Broker 滚动重启 / 网络抖动,监控显示 Producer 端 send 全成功,但下游消费少了几千条。
现场快照
T0 Producer send → Leader (Broker-1) 写本地 OK → 返回 ack
T0+5 Broker-1 还没把数据同步给 Follower 就宕机
T1 Controller 选 Broker-2 为新 Leader
→ Broker-2 上没那条数据 → 数据丢失根因分析
acks=1只等 Leader 本地写入,不等 Follower 同步;Leader 一切就丢。- 业务上线时按默认值(旧 Kafka 默认
acks=1,3.0+ 默认改all),没人审。
解决方案
- 改
acks=all+enable.idempotence=true+min.insync.replicas=2(3 副本时)。 - 已经丢的消息:从下游业务侧人工补数。
预防措施
- 平台对 Producer 客户端配置做静态扫描:业务关键 Topic 必须
acks=all。 - 上线 Checklist 第 1 条:
确认 acks 配置。
相关参数:acks / enable.idempotence / min.insync.replicas
关联章节:Ch 04 Producer 深入、Ch 09 副本与 ISR
2.2 没开幂等,重试导致重复
症状
- 业务消息有时会被消费方收到 2 次甚至 3 次。
- 排查后发现 Producer 收到了
TimeoutException但实际 Broker 已经写入,Producer 重发后变成两条。
根因:未开 enable.idempotence,Producer 重试时同一条消息被 Broker 接收两次。
解决方案
enable.idempotence=true(自动设acks=all、max.in.flight ≤ 5、retries=MAX)。- 消费侧业务设计为幂等(业务主键 + 去重表 / Redis SETNX)。
预防措施
- Producer 默认配置模板里硬编码
enable.idempotence=true。 - 关键业务消费端必须有幂等兜底。
相关参数:enable.idempotence
关联章节:Ch 04 Producer 深入、Ch 12 消息语义、Ch 13 EOS
2.3 batch.size 太大 → 延迟暴增;太小 → 吞吐很差
症状(太小)
- send rate 上不去,CPU / 网络都很闲;监控
batch-size-avg才 1KB(默认 16KB 的 1/16)。
症状(太大)
batch.size=4MB+linger.ms=1s→ 每条消息端到端延迟 1s 起步;buffer.memory经常 exhausted。
根因
- batch 太小:每个 RPC 只发 1 条消息,开销巨大。
- batch 太大:消息要等攒够 / 等 linger 时间,延迟高;buffer 容易满。
解决方案 / 调优
- 经验配置(吞吐优先):
batch.size=128KB ~ 256KB+linger.ms=20 ~ 50+compression.type=zstd。 - 经验配置(延迟优先):
batch.size=16KB(默认)+linger.ms=0+compression.type=lz4。 - 监控
batch-size-avg(接近 batch.size 上限说明攒得满)/record-queue-time-avg(攒太久)。
预防:每次调参一对(batch.size + linger.ms)一起改,单独动一个意义不大。
相关参数:batch.size / linger.ms / buffer.memory / compression.type
关联章节:Ch 04 Producer 深入、Ch 08 高吞吐底层
2.4 默认 Sticky Partitioner 把小流量都打到一个分区
症状
- 6 分区 Topic,只有 1 个分区一直在写,其他 5 个分区 0 流量。
- 切到新 Topic 也是这个现象。
根因
- 3.0+ 默认
DefaultPartitioner是 Sticky 行为:没有 Key时,会粘在同一分区直到攒满 batch / 触发 linger,然后才切换。 - 流量低 → 一直没攒满 → 一直粘在那个分区。
解决方案
- 让消息带 Key(最佳)。
- 或者强制
RoundRobinPartitioner(牺牲连续批次粘性)。 - 流量大起来后会自动均匀(攒满 batch 频繁切分区)。
预防:上线前低流量灰度时对分区分布有预期,不要看到「不均」就以为出问题。
相关参数:partitioner.class / partitioner.ignore.keys
关联章节:Ch 04 Producer 深入
2.5 选 gzip 压缩 → CPU 100%
症状
- Broker 上线压缩,选了
gzip,每秒几万消息,Broker CPU 飙到 100%,Producer 端 CPU 也居高不下。
根因
gzip是 5 种压缩算法里最慢但压缩率最高的;高 QPS 场景 CPU 顶不住。
对照表:
| 压缩算法 | CPU 开销 | 压缩率 | 推荐场景 |
|---|---|---|---|
none | 0 | 1.0 | 内网 + 巨大带宽 |
lz4 | 低 | 1.5 ~ 2× | 延迟优先(默认推荐) |
snappy | 低 | 1.5 ~ 2× | 老系统兼容 |
zstd | 中 | 2 ~ 3× | 吞吐 + 存储双优先(推荐) |
gzip | 高 | 2.5 ~ 3.5× | 离线归档、低 QPS |
解决方案:改 compression.type=zstd(KIP-110,0.10+ 引入;2.1+ 服务端支持)。
相关参数:compression.type
关联章节:Ch 04 Producer 深入
三、Consumer 端
3.1 自动提交 + 业务异常 → 消息丢失
症状
- 消费者重启或 crash 后发现业务实际处理了 800 条,但 offset 已经提交到 1000,中间 200 条永久丢失。
根因
enable.auto.commit=true时,提交 offset 由后台定时器自动触发,与业务处理是否成功无关。- poll 拿到 1000 条 → 业务处理到 800 条挂了 → 后台定时器已经提交 1000 → 重启从 1001 读起。
解决方案
- 关掉自动提交:
enable.auto.commit=false,业务处理完后手动commitSync()/commitAsync()。 - 标准模板:
python
for msg in consumer:
try:
process(msg)
consumer.commit(message=msg, asynchronous=False)
except Exception:
logger.exception("process failed, will retry next poll")
# 不 commit,下次 poll 还是从这条开始预防:Consumer 默认配置模板硬编码 enable.auto.commit=false。
相关参数:enable.auto.commit / auto.commit.interval.ms
关联章节:Ch 05 Consumer 深入、Ch 12 消息语义
3.2 手动提交位置错位 → 重复或丢失
症状
- 改成手动提交后,要么消息重复消费,要么漏消费,看代码也找不到错。
根因(典型 4 种)
process(msg)之前就commit(msg)→ 处理失败,offset 已提交,丢消息。commit(msg.offset)而不是commit(msg.offset + 1)→ 重复消费一条。- 多线程消费,子线程处理完不通知主线程,主线程错误提交了「最新 poll 的 offset」 → 提交超前。
- 在
onPartitionsRevoked回调外面提交 → Rebalance 时 offset 没及时落,下次起来重复。
解决方案
- 永远
process → commit,不要颠倒。 - Confluent / Java 的
commit(message=msg)内部已经是offset+1,自己拼装时记得加 1。 - 多线程消费用
kafka-python-ng / aiokafka的 batch commit 模式,不要自己 hack。 - 在
ConsumerRebalanceListener.onPartitionsRevoked里同步提交。
相关参数:partition.assignment.strategy / 回调 ConsumerRebalanceListener
关联章节:Ch 05 Consumer 深入、Ch 11 Rebalance
3.3 反序列化异常没 catch → 整个消费组卡住
症状
- 偶发 1 条 schema 错乱的消息进 Topic 后,整个消费组所有实例卡住,重启也不行。
根因
consumer.poll()内部反序列化抛异常 → 异常往外冒 → 业务代码若无 catch,进程退出 → 自动重启 → 又从同一条 offset 开始 poll → 又抛 → 死循环。
解决方案
- 用
value.deserializer=ByteArrayDeserializer拿 raw bytes,业务侧自己反序列化并 catch。 - 反序列化失败的消息:写到 DLQ Topic,跳过这条 offset 继续。
- 或 Confluent 的 Error Handling Deserializer。
预防
- 引入 Schema Registry + 兼容性策略 BACKWARD(生产者删字段、消费者老 schema 仍兼容)。
- DLQ 模式默认开启。
相关参数:value.deserializer / errors.tolerance / errors.dead.letter.queue.topic.name
关联章节:Ch 05 Consumer 深入、Ch 18 Schema Registry
3.4 auto.offset.reset=latest 导致初次上线丢历史数据
症状
- 新业务上线,期望从 Topic 头读历史,结果只读到了上线之后的数据;前几年累积的几亿条全没读到。
根因
auto.offset.reset=latest(默认):找不到已提交 offset 时,从最新位置开始读。- 新 Consumer Group 在
__consumer_offsets没记录 → 触发 reset → 从 latest。
解决方案
- 用历史回放的业务:上线前
auto.offset.reset=earliest;或kafka-consumer-groups --reset-offsets --to-earliest --execute显式重置。 - 或者首次启动时
seek(beginning)强制定位。
预防:每次新建 Group 上线前在审批清单里写明「期望初始位点:earliest / latest / 时间戳」。
相关参数:auto.offset.reset / kafka-consumer-groups --reset-offsets
关联章节:Ch 05 Consumer 深入、Ch 12 消息语义
四、Rebalance 风暴
4.1 处理慢 → max.poll.interval.ms 超时 → 死循环 Rebalance
症状
- 消费业务 CPU 偏高、Lag 一直涨;查日志发现每隔几分钟就一次
RebalanceInProgress,消费组永远稳定不下来。
根因
- 业务处理一批消息要 8 分钟,但
max.poll.interval.ms=5min(默认);5 分钟没新 poll → Coordinator 认为消费者死亡 → 触发 Rebalance。 - Rebalance 时正在处理的消息 commit 不上,新一轮起来又 poll 同样数据 → 又 8 分钟超时 → 再 Rebalance → 死循环。
解决方案
- 调大
max.poll.interval.ms(如 15min)。 - 调小
max.poll.records(默认 500,改 50 ~ 100)让单批处理更快。 - 业务并行:单线程处理慢则多线程;I/O 密集型用异步。
预防:监控告警 last-rebalance-seconds-ago 频繁掉低就报警。
相关参数:max.poll.interval.ms / max.poll.records / session.timeout.ms
关联章节:Ch 05 Consumer 深入、Ch 11 Rebalance
4.2 Cooperative Rebalance 没启用 → 全组 stop-the-world
症状
- 100 个消费实例的大消费组,每次扩容 / 缩容,全部实例都暂停 30 秒以上,Lag 在 Rebalance 期间疯涨。
根因
- 默认
partition.assignment.strategy=[RangeAssignor, CooperativeStickyAssignor](3.0+),但老客户端 / 老 Java SDK 默认仍是RangeAssignor→ Eager Rebalance → 所有人放下手里的分区 → 重新分配 → 重新启动消费。
解决方案
- 全部消费者改成
partition.assignment.strategy=CooperativeStickyAssignor。 - 这是滚动升级:不能直接整组切,要按 KIP-429 升级路径 走(先发 Cooperative + Range 双策略,等所有节点都升级,再去掉 Range)。
预防:新建消费组默认 CooperativeStickyAssignor。
相关参数:partition.assignment.strategy
关联章节:Ch 11 Rebalance
4.3 Static Membership 缺失 → 滚动重启全员 Rebalance
症状
- 消费组 30 个实例,业务每次 deploy(10 分钟内逐个重启)会触发 30 次 Rebalance。
根因
- 没设
group.instance.id,每个实例重启 = 一个新成员加入 = 旧成员超时离开 = 触发 Rebalance。
解决方案
- 给每个消费实例设
group.instance.id=consumer-<pod-name>(K8s 用 StatefulSet 拿稳定名)。 - 配合
session.timeout.ms调到 60s 左右,让短暂重启不踢出。
预防:所有有状态消费者默认开 Static Membership。
相关参数:group.instance.id / session.timeout.ms
关联章节:Ch 11 Rebalance
五、副本与 ISR
5.1 unclean.leader.election=true → 已 ack 数据丢失
症状
- 集群发生过一次 Broker 全挂事故;只有 Broker-2(一个非 ISR 副本)能起来,被选为 Leader;之后业务发现已经 ack 给 Producer 的消息凭空消失。
根因
unclean.leader.election.enable=true允许非 ISR 副本成为 Leader;该副本数据本来就落后于真正的 ISR。- 一旦选它当 Leader,新 Leader 的 LEO < 老 Leader 的 LEO → 数据被截断。
解决方案
- 立即关闭
unclean.leader.election.enable=false(2.0+ 默认 false,但被很多老集群继承下来)。 - 故障期间宁可分区不可用(OFFLINE)也不能让数据消失。
预防:硬性 Broker 配置 unclean.leader.election.enable=false,并在 Topic 级也保持 false。
相关参数:unclean.leader.election.enable
关联章节:Ch 09 副本与 ISR
5.2 min.insync.replicas 配置缺失 → acks=all 形同虚设
症状
- Producer 配
acks=all觉得万无一失,但还是丢消息。
根因
- 副本数 3,但
min.insync.replicas=1(默认)。 - ISR 内只剩 Leader 一个时,Leader 写入即返回 OK,和 acks=1 无区别。
- 此时 Leader 一挂就丢。
解决方案
- Topic 配置
min.insync.replicas=2(3 副本时)。 - Broker 默认
min.insync.replicas=2。
取舍:min.insync.replicas=2 时,挂 1 个 Broker(ISR 变 2)仍可写;挂 2 个时写入直接被拒(NotEnoughReplicasException),保不丢牺牲可写性。
相关参数:min.insync.replicas / acks
关联章节:Ch 09 副本与 ISR
5.3 replica.lag.time.max.ms 太短 → ISR 抖动不停
症状
- 集群没出故障,但监控里
IsrShrinksPerSec/IsrExpandsPerSec一直高频跳动。
根因
replica.lag.time.max.ms设得太小(如 5s),Follower 偶发慢一点就被踢出 ISR;恢复后又加回 → 反复抖。
解决方案
- 默认 30000ms 一般够用;偶发抖动调到 60000ms。
- 检查 Broker 网络 / 磁盘是不是真的不稳定,先治根本问题。
预防:监控 ISR 抖动告警阈值合理(每分钟少于 5 次)。
相关参数:replica.lag.time.max.ms / num.replica.fetchers
关联章节:Ch 09 副本与 ISR
六、Controller / KRaft
6.1 ZK → KRaft 迁移过程中误改 process.roles
症状
- 迁移过程中操作员手动把某 Broker 的
process.roles=broker改成broker,controller,重启后 Controller Quorum 选举混乱,元数据日志写入失败。
根因
- Controller Quorum 必须是预先设定的 voters;中途加 controller 角色不被支持(只能用「移除 + 重新加入」流程)。
解决方案
- 严格按 KIP-866 ZK → KRaft Migration 步骤走:
- 启用 KRaft Controller Quorum(独立 3 节点)
- Broker 配
zookeeper.metadata.migration.enable=true,仍连 ZK + 注册到 KRaft - Migration 完成后,Broker 配 KRaft only,ZK 下线
- 错误改了 → 回滚配置 → 清掉错误节点的 metadata 数据卷 → 重新加入。
预防:迁移前全员演练至少 1 次;KRaft 数据目录单独高速盘 + RPO 短的备份。
相关参数:process.roles / controller.quorum.voters / zookeeper.metadata.migration.enable
6.2 KRaft Quorum 配置错(少了一个 voter)
症状
- 3 节点 Controller Quorum 配置成
1@h1:9093,2@h2:9093(少了 voter 3),KRaft 启动后仍能选出 Leader(2 个就够 Quorum),但 voter 3 的元数据永远落不下来,扩缩容时报错。
根因
controller.quorum.voters配置不一致,所有节点必须配相同的 voters 列表,否则元数据视图不一致。
解决方案
- 同时改三台 Broker 的
controller.quorum.voters为完整 3 节点;滚动重启 Controller。
预防:用配置管理工具(Ansible / Puppet / K8s ConfigMap)保证一致;上线前 kafka-metadata-shell 确认 voters 一致。
相关参数:controller.quorum.voters / node.id
七、EOS / 事务
7.1 transactional.id 重复 → 永久 ProducerFenced
症状
- 两台机器跑同一个微服务,都用了相同的
transactional.id="order-service",结果第二台启动后第一台立刻报ProducerFencedException退出,业务一台台轮流挂。
根因
transactional.id是事务的全局唯一身份;后启动的 Producer 会bump epoch 把同 ID 的旧 Producer 踢掉(PID Fencing)。- 这是设计上保护跨实例双写的机制,不是 bug。
解决方案
- 每个实例用唯一稳定的 transactional.id:
tx-<service>-<pod-name>/tx-<service>-<host>-<port>。 - K8s StatefulSet 拿稳定 podName。
预防:硬编码模板里禁止用静态字符串当 transactional.id。
相关参数:transactional.id / transaction.timeout.ms
关联章节:Ch 13 EOS
7.2 事务超时 transaction.timeout.ms 过小 → 大批量 abort
症状
- 业务跑 batch consume → process → produce 模式,单事务里要处理 1000 条;偶发慢一点就
InvalidProducerEpochException,整批被 abort。
根因
transaction.timeout.ms默认 60s,业务实际可能要更长。- 同时不能超过 Broker 的
transaction.max.timeout.ms(默认 15min)。
解决方案
- 调大 Producer
transaction.timeout.ms到 5 ~ 10min,同时确保 ≤ Broker 的transaction.max.timeout.ms。 - 减小单事务消息数(
max.poll.records降到 100)。
预防:监控 record-error-rate + 事务 abort 率告警。
相关参数:transaction.timeout.ms / transaction.max.timeout.ms(Broker)
关联章节:Ch 13 EOS
7.3 滥用事务做小消息 → 吞吐爆跌
症状
- 只是想做幂等,开了
enable.idempotence=true不够,又开了事务(每条消息一个事务),结果吞吐从 50 万 QPS 掉到 5000 QPS。
根因
- 每个 commit 涉及两阶段提交:写 PREPARE 记录到
__transaction_state→ 给所有参与分区写控制消息 → 写 COMMITTED → 客户端等 ack。 - 事务非常重,单条事务化 = 把每条消息的 latency 从 10ms 拉到 100ms+。
解决方案
- 仅开
enable.idempotence=true即可保 Producer 端不重复(At Least Once + 幂等 = 类 EOS)。 - 真正需要 EOS 的场景:consume-process-produce 链路(Streams 用法),批量提交(一次事务包含一个 poll 的几百条)。
- 业务端可幂等的,不用事务。
预防:审计 transactional.id 使用,没必要的取消。
相关参数:transactional.id / enable.idempotence
关联章节:Ch 13 EOS
八、Compaction
8.1 Compaction 不生效 / 旧版本不删
症状
- 配了
cleanup.policy=compact的 Topic,磁盘一直涨;用kafka-dump-log看,同一 Key 的几百个旧版本都还在。
根因(4 选 1)
- Active segment 永远不被 compact(只 compact「关闭」的 segment);如果
segment.ms/segment.bytes滚不动,永远不 compact。 min.cleanable.dirty.ratio=0.5(默认):脏数据比例没到 50% 不触发。min.compaction.lag.ms=0(默认):但有些版本会有「保护期」。log.cleaner.enable=false被 Broker 强制关掉了。
解决方案
- 缩小
segment.ms/segment.bytes让 segment 更频繁滚动,老 segment 才能被 compact。 - 把
min.cleanable.dirty.ratio降到 0.1。 - 确认
log.cleaner.enable=true(默认 true)。 - 看
log.cleaner.threads是不是 0(曾经踩过log.cleaner.io.max.bytes.per.second限流过狠的坑)。
预防:监控 kafka.log.cleaner:type=LogCleanerManager,name=max-dirty-percent。
相关参数:cleanup.policy / min.cleanable.dirty.ratio / segment.ms / log.cleaner.enable / log.cleaner.threads
关联章节:Ch 14 Log Compaction
8.2 Cleaner 线程挂了,dirty ratio 飙升
症状
- 监控
max-dirty-percent一路涨到 99%,磁盘空间被旧版本占满。 log-cleaner.log里有 OOM 异常或IllegalStateException。
根因
- Cleaner 线程对单个 Key 出现了几亿次的 Topic 处理时,去重 buffer(
log.cleaner.dedupe.buffer.size)不够 → OOM → 线程挂掉。 - Cleaner 线程一旦挂掉不会自动重启,必须重启 Broker。
解决方案
- 加大
log.cleaner.dedupe.buffer.size(256MB → 1GB)。 - 加大
log.cleaner.threads到 4 ~ 8。 - 重启 Broker 让 Cleaner 起来。
预防:JMX 监控 cleaner 工作状态;磁盘告警阈值留 30% 余量给应急扩容。
相关参数:log.cleaner.dedupe.buffer.size / log.cleaner.threads / log.cleaner.io.max.bytes.per.second
关联章节:Ch 14 Log Compaction
8.3 Tombstone 不被清理 → 占空间
症状
- 用
null value标记删除(tombstone)后,过了 24 小时这些 tombstone 还在。
根因
- Tombstone 有保留期
delete.retention.ms(默认 24h):保证下游 Compaction Consumer 能看到删除事件。 - 这个时间内 tombstone 不被清。
解决方案
- 业务能容忍快速消失的,把
delete.retention.ms调小(如 1h)。 - 注意调小后,Compaction-Consumer 的处理 SLA 必须低于这个值。
相关参数:delete.retention.ms / min.compaction.lag.ms
关联章节:Ch 14 Log Compaction
九、运维
9.1 磁盘满 → Broker 全挂,数据写不进
症状
df -h看 100% 满了;server.log报KafkaStorageException: I/O error,Broker 进程挂掉。- 起不来:因为还要写
recovery-point-offset-checkpoint。
根因
- 没及时清理 / 扩容;
log.retention配错;磁盘告警阈值太宽松。
解决方案(紧急)
- 绝对不要
rm数据文件,会破坏 segment 索引。 - 用
kafka-configs临时把某 Topic 的retention.ms设短(如 1h)→ 等几分钟自动清理。 - 或者
kafka-delete-records.sh删指定 offset 之前的数据。 - 或挂载临时新盘到
log.dirs(KRaft 后支持热加 disk)。
预防
- 磁盘告警阈值 80%;持续 30min 不下降则扩容。
- 容量评估按
日均增量 × 保留天数 × 副本数 × 1.3 余量。
相关参数:log.retention.ms / log.retention.bytes / kafka-delete-records.sh
关联章节:Ch 19 可观测性与运维
9.2 文件句柄 / Socket 句柄不够 → Too many open files
症状
- 大集群 / 多分区 / 多客户端连接时,Broker 报
Too many open files,请求挂起。
根因
- 每个 segment 文件 + 每个 socket 连接都占一个 fd;分区数和客户端多了之后远超系统默认 1024 / 4096。
解决方案
ulimit -n 1048576(systemdLimitNOFILE、容器nofile65536+)。- 验证:
cat /proc/$(pgrep -f kafka)/limits | grep open。
预防:上线前压测把句柄打满做一次容量验证。
相关参数:OS nofile / nproc / vm.max_map_count
关联章节:Ch 19 可观测性与运维
9.3 JVM 堆给得太大 → Full GC 卡几十秒
症状
- Broker JVM 堆设了 32G(觉得越大越好),结果偶发 Full GC 30 秒,所有客户端报 timeout,Controller 选举触发,集群抖动。
根因
- Kafka 不在 JVM 堆里缓存消息,靠 PageCache。堆只放:网络 buffer、元数据、controller 状态、日志清理 buffer。
- 6 ~ 10G 堆已经够 99% 场景;32G 大堆 GC 算法(即使是 G1 / ZGC)在停顿时间上仍有限。
解决方案
- 堆设
-Xms6G -Xmx6G,操作系统留出 ≥ 50% 内存给 PageCache。 - GC 算法用 G1(默认);4.0+ 推荐 ZGC(generational)。
预防:上线前压测看 GC 时间分布;告警阈值 Full GC 单次 > 1s。
相关参数:KAFKA_HEAP_OPTS / G1 / ZGC / OS PageCache
关联章节:Ch 08 高吞吐底层、Ch 19 可观测性与运维
十、跨机房
10.1 ISR 跨机房抖动
症状
- 部署成「3 副本,分布在 3 个 AZ」,业务写入偶发延迟飙升(500ms+),ISR 频繁收缩。
根因
- 跨 AZ 网络延迟 5 ~ 30ms;副本同步走最慢链路。
acks=all+min.insync.replicas=2时,每条消息都要等到至少 2 个副本(含本地 + 至少 1 个远端)确认。
解决方案
- 业务可接受 At Least Once 的:
acks=1+ 本地副本 + 异步跨机房复制(MirrorMaker 2)。 - 必须强一致的:3 副本部署在 1 个机房(高可用换数据安全),用 MM2 做跨机房灾备。
- 调高
replica.lag.time.max.ms=60000,避免短抖动剔 ISR。
预防:跨机房部署前做网络延迟压测,评估业务可接受的 RTO / RPO。
相关参数:acks / min.insync.replicas / replica.lag.time.max.ms
关联章节:Ch 09 副本与 ISR、Ch 19 可观测性与运维
10.2 网络分区脑裂
症状
- 跨机房集群发生网络分区,分区 A 仍能写入(拥有 Leader),分区 B 也以为自己是 Controller,开始选举。
- 网络恢复后两边数据冲突。
根因
- Controller 切换 + Leader 切换在网络分区下不可靠(CAP 中 Kafka 选 CP)。
- ZK 时代 split-brain 风险更高;KRaft 用 Raft Quorum 缓解。
解决方案
- KRaft 模式:Quorum 大于一半才能选 Leader,自然避免脑裂(少数派直接不可写)。
- 业务感知 Topic 标记:网络分区期间业务降级(停写 / 重试到主机房)。
- 灾后用 Cruise Control / 手动 reassignment 修复。
预防:跨机房部署 5 节点 Controller Quorum(不是 3 节点跨 2 机房,那是 split brain 温床);推荐 3 + 1 + 1 拓扑。
关联章节:Ch 09 副本、Ch 10 KRaft
十一、Connect / Streams / ksqlDB
11.1 Connect Offset Topic 漂移导致 Source 重读
症状
- Debezium MySQL Source Connector 配置不变,重启后重新全量拉取几亿行历史 binlog,下游被洪水冲垮。
根因
connect-offsetsTopic 副本数选了 1,某次 Broker 故障导致 partition 不可用 → Connector 拿不到上次 offset → 重新从初始位置拉。
解决方案
connect-offsets副本数至少 3,分区 25(默认推荐)。- Connector 配置里
offset.storage.replication.factor=3。 - 灾后用
kafka-console-consumer.sh --topic connect-offsets --formatter org.apache.kafka.connect.runtime.OffsetUtils$Formatter恢复手工设置。
预防:起 Connect 集群前先校验内部 Topic 副本数;配置 monitoring。
相关参数:offset.storage.topic / offset.storage.replication.factor / offset.storage.partitions
关联章节:Ch 16 Kafka Connect
11.2 Streams State Store RocksDB 文件爆掉
症状
- Kafka Streams 应用跑几天后磁盘满,
/tmp/kafka-streams/<app-id>/...占用几百 GB。
根因
- 默认 State Store 用 RocksDB,本地路径
state.dir=/tmp/kafka-streams;机器重启/tmp清空 → 应用重启时从 changelog Topic 全量重建 → 几亿条数据写本地 → 撑爆。 - 或者 windowed store 没设 retention,窗口越攒越多。
解决方案
state.dir改到独立大盘,不要用 /tmp。- Windowed store 配 retention:
Materialized.as(...).withRetention(Duration.ofDays(7))。 - 监控 RocksDB JMX:
stream-state-metrics。
预防:上线前模拟「应用重启 + 全量回放 changelog」场景,看本地存储增长趋势。
相关参数:state.dir / Streams Materialized.withRetention()
关联章节:Ch 17 Streams
11.3 ksqlDB 跨版本 schema 不兼容
症状
- 升级 ksqlDB 0.28 → 0.29,重启 server 后所有 STREAM / TABLE 全部
DESERIALIZATION_ERROR。
根因
- ksqlDB 内部 metadata schema 在小版本间偶发不兼容;
_confluent-ksql-default__command_topic是 compact Topic,不能简单删除重建。
解决方案
- 严格按官方 release notes 的 upgrade guide 升级。
- 灾难恢复:备份
_confluent-ksql-default__command_topic,新建 ksqlDB 集群,按命令逐条重放 DDL。
预防:生产 ksqlDB 跳版本升级前在 staging 演练;订阅 release notes。
相关参数:ksql.service.id / ksql.streams.bootstrap.servers
上线 Checklist
把这份清单打印出来,每次新业务上 Kafka 前逐条核对,至少能避免 80% 的踩坑。
集群 / Broker 层(10 条)
- [ ] 1. 副本数 ≥ 3(业务 Topic + 内部 Topic)
- [ ] 2.
min.insync.replicas=2(3 副本时) - [ ] 3.
unclean.leader.election.enable=false - [ ] 4.
auto.create.topics.enable=false - [ ] 5.
delete.topic.enable=true(3.0+ 默认) - [ ] 6.
log.retention.*与磁盘容量匹配 - [ ] 7. JVM 堆 6 ~ 10G,OS 留 ≥ 50% 内存给 PageCache
- [ ] 8.
nofile≥ 1048576,vm.max_map_count≥ 262144 - [ ] 9. KRaft Controller Quorum 3 / 5 节点;
controller.quorum.voters各节点一致 - [ ] 10. 时钟同步 NTP / Chrony,所有 Broker 时间漂移 < 100ms
Topic / 分区设计(5 条)
- [ ] 11. 分区数按 1 ~ 2 年峰值流量 + 30% 余量预留
- [ ] 12. Topic 命名遵循
<env>.<biz>.<event>.<v>规范 - [ ] 13. Key 选择经过分布性评估(无 90% 热点)
- [ ] 14. 内部 Topic(
__consumer_offsets/__transaction_state/_schemas/connect-*)副本因子 ≥ 3 - [ ] 15. 业务 Topic 必须填责任人 + TTL + 上下游消费方
Producer 配置(5 条)
- [ ] 16.
acks=all(关键业务);enable.idempotence=true - [ ] 17.
compression.type=zstd(吞吐优先)或lz4(延迟优先) - [ ] 18.
linger.ms+batch.size一对调参,不要单独动 - [ ] 19.
transactional.id唯一稳定(如<service>-<pod>) - [ ] 20. Producer 客户端配
client.id,便于排障 / 配额
Consumer 配置(5 条)
- [ ] 21.
enable.auto.commit=false,业务处理完手动提交 - [ ] 22.
auto.offset.reset上线前显式确认 earliest / latest - [ ] 23.
max.poll.interval.ms≥ 业务最长处理时间 × 1.5;max.poll.records调到能在该时间内处理完 - [ ] 24.
partition.assignment.strategy=CooperativeStickyAssignor - [ ] 25. 有状态消费者开
group.instance.idStatic Membership
监控告警(5 条)
- [ ] 26.
UnderReplicatedPartitions > 0持续 5min 告警 - [ ] 27.
OfflinePartitionsCount > 0立即告警 - [ ] 28.
UncleanLeaderElectionsPerSec > 0立即告警 - [ ] 29. Consumer
records-lag-max业务侧阈值告警 - [ ] 30. Cluster Controller 数 = 1(>1 或 0 都是异常)
- [ ] 31. JVM Full GC 单次 > 1s 告警
- [ ] 32. 磁盘使用率 > 80% 告警;> 90% 立即扩容
应急预案(必备)
- [ ] 33. 拉一份「常见错误 → 处置 SOP」文档(参考本附录第 7 节)
- [ ] 34. 演练过:删 Topic / 改副本 / 重置 Offset / 分区迁移 / Leader 切换 / Broker 重启 至少一次
- [ ] 35. 备份策略:MirrorMaker 2 / Cluster Linking / 跨机房灾备演练过
📌 本附录的姐妹篇:
appendix_cheatsheet.md— 命令 / 参数 / JMX 速查appendix_kafka_vs_others.md— 横向对比interview.md— 面试题总索引(很多踩坑直接对应高频题)- Ch 20 常见踩坑章节 — 12 个最高频的「教学版」踩坑案例(更详细但数量少)