Skip to content

附录 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 漂 │
   │                   │                         │
   └───────────────────┼─────────────────────────┘
                       │  低 ─────────────→ 高
                       └─────► 高频度
                       
   📌 优先攻克「右上角」案例:高频 × 高危。

目录

  1. 设计类(Topic / 分区 / 副本 / Key)
  2. Producer 端
  3. Consumer 端
  4. Rebalance 风暴
  5. 副本与 ISR
  6. Controller / KRaft
  7. EOS / 事务
  8. Compaction
  9. 运维
  10. 跨机房
  11. Connect / Streams / ksqlDB

末尾:上线 Checklist(30+ 项)


一、设计类

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 份再聚合。
  • 复合 Keyapp_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_xxxtmp_xxxbackup_orders_2023_07ordersorders_v2orders_neworders_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 层 nofilenproc

关联章节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=allmax.in.flight ≤ 5retries=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+ 默认 DefaultPartitionerSticky 行为:没有 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 开销压缩率推荐场景
none01.0内网 + 巨大带宽
lz41.5 ~ 2×延迟优先(默认推荐)
snappy1.5 ~ 2×老系统兼容
zstd2 ~ 3×吞吐 + 存储双优先(推荐)
gzip2.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 种)

  1. process(msg) 之前就 commit(msg) → 处理失败,offset 已提交,丢消息。
  2. commit(msg.offset) 而不是 commit(msg.offset + 1) → 重复消费一条。
  3. 多线程消费,子线程处理完不通知主线程,主线程错误提交了「最新 poll 的 offset」 → 提交超前。
  4. 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 步骤走:
    1. 启用 KRaft Controller Quorum(独立 3 节点)
    2. Broker 配 zookeeper.metadata.migration.enable=true,仍连 ZK + 注册到 KRaft
    3. Migration 完成后,Broker 配 KRaft only,ZK 下线
  • 错误改了 → 回滚配置 → 清掉错误节点的 metadata 数据卷 → 重新加入。

预防:迁移前全员演练至少 1 次;KRaft 数据目录单独高速盘 + RPO 短的备份。

相关参数process.roles / controller.quorum.voters / zookeeper.metadata.migration.enable

关联章节Ch 10 Controller 与 KRaft


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

关联章节Ch 10 Controller 与 KRaft


七、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)

  1. Active segment 永远不被 compact(只 compact「关闭」的 segment);如果 segment.ms / segment.bytes 滚不动,永远不 compact。
  2. min.cleanable.dirty.ratio=0.5(默认):脏数据比例没到 50% 不触发。
  3. min.compaction.lag.ms=0(默认):但有些版本会有「保护期」。
  4. 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.logKafkaStorageException: 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(systemd LimitNOFILE、容器 nofile 65536+)。
  • 验证: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 副本与 ISRCh 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-offsets Topic 副本数选了 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

关联章节Ch 17 Streams 与 ksqlDB


上线 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.id Static 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 / 跨机房灾备演练过

📌 本附录的姐妹篇