Skip to content

第 20 章 常见踩坑与排障案例集

每一个踩坑都是用别人加班换来的经验。 这一章把 12 个真实事故按「症状 → 现场快照 → 根因分析 → 解决方案 → 预防措施 → 相关参数」的统一格式复盘。

学完你会:

  • 看症状能快速猜中根因(80% 的 Kafka 问题就是这 12 类)
  • 知道每个踩坑对应要改的参数和命令
  • 能用「常见错误日志辞典」秒识别错误信息

📌 生活类比:踩坑案例集 = 「事故复盘报告汇编」。学最快的方式不是看官方手册,而是把别人的事故当自己的事故复盘一遍


案例 1:分区数选少了,扩容时业务全乱

症状

  • 上线时 partitions=4,业务跑了半年,QPS 涨到瓶颈
  • partitions 加到 16,第二天业务投诉「同一用户的订单顺序错乱」「Redis 缓存频繁穿透」

现场快照

bash
# 加分区前
kafka-topics.sh --describe --topic orders
# Topic: orders  PartitionCount: 4

# 加分区
kafka-topics.sh --alter --topic orders --partitions 16

# 业务投诉
2026-04-15 10:23:01 WARN  user_id=1234 cache miss again
2026-04-15 10:23:02 WARN  user_id=1234 order out-of-order

根因分析

  • Kafka 默认 Partitioner 用 hash(key) % partitions
  • 加分区后分母变了,同一个 key 之前落在分区 3,现在可能落在分区 11
  • 该 key 的旧消息和新消息在不同分区,单分区顺序保证失效
  • 下游消费者按 key 缓存的状态(如 Redis 里的「用户最新订单」),新消息走到不同消费者实例,缓存全部失效

解决方案

  • 回退:分区数无法减少,必须新建一个 Topic(orders.v2 用 16 分区),双写过渡
  • 过渡期:上游同时写新旧 Topic,下游切换消费新 Topic,旧 Topic 消费完后下线

预防措施

  1. 分区数一开始就预估到未来 2-3 年的峰值,宁多勿少
  2. 单 Broker 推荐 100-2000 分区上限;按 Broker 数 × 200 估算
  3. 用「自定义 Partitioner」而非 hash mod,例如 hash(key) % 1024 → 路由表 → 真实分区,加分区时只改路由表
  4. 关键 Topic 测试时就模拟「加分区是否破坏顺序」

相关参数

  • num.partitions(Topic 创建默认值)
  • partitioner.class(自定义 Partitioner)

案例 2:Key 设计不当,单分区热点

症状

  • 集群总吞吐 50K msg/s,看上去没问题
  • 但某个 Broker 的磁盘 IO / 网络打满,其他 Broker 闲置
  • 该 Broker 上的某个 Topic-Partition 占用 80% 流量

现场快照

Broker 1: BytesIn 200 MB/s   <-- 打满
Broker 2: BytesIn 30 MB/s
Broker 3: BytesIn 25 MB/s

Topic events: partition 7 = 35K msg/s, 其他分区 ~ 1K msg/s

业务方代码:

python
# 用「日期」作为 Key
producer.send("events", key=str(date.today()), value=event)

根因分析

  • Key 是「今天的日期」,所有当天消息哈希到同一个分区
  • 该分区的 Leader Broker 被独占,其他 Broker 干瞪眼
  • 消费端这个分区也卡住,整个 group 出现 Lag

解决方案

  • 立即改 Key 设计:用 user_id / order_id / device_id 等高基数字段
  • 临时降级:把 Key 设为 None(让 Producer 用 Sticky Partitioner 平均分布,但会牺牲顺序)
  • 修代码:producer.send("events", key=user_id, value=event)

预防措施

  1. Key 必须高基数(cardinality),至少 ≥ 分区数 × 100
  2. 禁止用 时间、固定 enum、城市名、机型名 这种低基数字段做 Key
  3. 上线前压测时专门检查「分区流量分布」(用 kafka-consumer-groups.sh --describe 看每分区 Lag)
  4. 监控告警:分区流量标准差 > 平均值 50% 报警

相关参数

  • 无需改 Broker 配置,纯业务设计问题
  • 可选 partitioner.class=RoundRobinPartitioner(牺牲顺序保平均)

案例 3:acks=1 下 Leader 切换丢消息

症状

  • 业务统计「订单数」与上游 Producer 的发送量对不上,少了 0.1%
  • 监控显示没有任何 Producer 报错
  • 偶发,随 Broker 重启 / 扩缩容时丢得多

现场快照

python
producer = Producer({
  "bootstrap.servers": "...",
  "acks": "1",         # ← 只等 Leader ack
  "retries": 3,
})

Broker 滚动重启时间表:

10:00 重启 broker-1(Leader 切到 broker-2)
10:00 ~ 10:00:30  Producer 收到 broker-1 的 ack,但 broker-2 还没拉到这些消息
10:00:31 broker-1 起来后角色变为 Follower,其本地超出 HW 的部分被截断 → 数据消失

根因分析

  • acks=1:Leader 写入本地 log 立刻 ack,不等 Follower 同步
  • 如果 Leader 在 ack 之后、Follower 拉到之前宕机
    • 新 Leader 选 Follower,Follower 没有这些消息
    • 旧 Leader 起来时发现自己 LEO > 新 HW,会截断超出部分
    • 这些消息永久丢失

解决方案

  • 关键场景acks=all + min.insync.replicas=2(3 副本下保证至少 2 个写成功才返回)
  • 配合 enable.idempotence=true(幂等 Producer)避免重复
  • Producer 端开启 retries=Integer.MAX_VALUE + delivery.timeout.ms=120000

预防措施

  1. 重要业务一律 acks=all,吞吐损失通常 < 20%
  2. 滚动重启 Broker 时控制速度(一台一台),并监控 Producer 错误率
  3. 用「业务对账」兜底:定期 diff 上游发送 ID 与下游消费 ID

相关参数

  • Producer: acks=all, retries=Integer.MAX_VALUE, enable.idempotence=true, max.in.flight.requests.per.connection<=5
  • Broker: min.insync.replicas=2(Topic 级别)

案例 4:auto.offset.reset=latest 初次上线丢历史

症状

  • 新接入一个消费组,消费了 10 分钟才发现:只消费到了上线后的消息,前面 7 天的全没消费
  • 业务统计少了一大块

现场快照

python
consumer = Consumer({
  "group.id": "new-team-stats",
  "auto.offset.reset": "latest",   # ← 默认值
})
consumer.subscribe(["orders.events"])

Topic 已存在 7 天,有 1000 万条历史消息。新 group 第一次连,没有 __consumer_offsets 记录,按 latest 跳到末尾。

根因分析

  • auto.offset.reset 仅在「没有有效 offset」时生效(首次接入或 offset 过期被清理)
  • latest = 跳到末尾 → 历史全部丢失
  • earliest = 从头开始
  • none = 直接报错(最严格)

解决方案

  • 立即操作:用 kafka-consumer-groups.sh --reset-offsets --to-earliest 重置到最早
  • 重新启动消费者,把历史补回来

预防措施

  1. 新接入业务先评估「是否需要历史」
    • 实时风控、统计:通常需要历史 → earliest
    • 实时通知、排行榜:通常不需要 → latest
  2. 关键业务用 none:强制人工设置初始 offset,避免静默错误
  3. 消费组首次上线后立即检查 Lag,确认是否在「合理范围」

相关参数

  • auto.offset.reset = earliest | latest | none
  • 重置命令:kafka-consumer-groups.sh ... --reset-offsets --to-earliest --execute

案例 5:消费慢 + max.poll.interval.ms 超时引发 Rebalance 风暴

症状

  • 凌晨业务突增,Consumer 处理变慢
  • 突然所有 Consumer 不再消费,Lag 暴涨
  • 日志疯狂打 Member xxx has been removed from group ...Rebalance triggered
  • Rebalance 一次接一次,永远不结束

现场快照

[10:01:23] Consumer-1 process took 360s
[10:01:30] coordinator: Member consumer-1 has left the group (max.poll.interval.ms exceeded)
[10:01:31] Rebalance start...
[10:01:35] consumer-2 joined, partition reassigned
[10:01:36] consumer-2 process took 350s ← 同样慢
[10:01:40] coordinator: Member consumer-2 has left ...
... 无限循环

根因分析

  • max.poll.interval.ms 默认 5 分钟(300_000ms)
  • Consumer 在两次 poll() 之间的处理超过这个时间,会被踢出 group
  • 踢出 → Rebalance → 重新分配分区给其他 Consumer → 它们也处理慢被踢
  • 整个 group 陷入「死循环 Rebalance」

解决方案

  • 临时止血:把 max.poll.interval.ms 调高到 30 分钟(1_800_000ms)
  • 缩小批次max.poll.records 从 500 调到 50,让每次 poll 处理时间短一些
  • 业务侧异步化:消费线程只入队,单独的 worker 处理,poll 永不阻塞

预防措施

  1. 理解 max.poll.interval.ms 是「业务处理时间上限」,不是网络超时
  2. 批量大小 × 单条处理时间 < max.poll.interval.ms 的 70%
  3. 慢操作(DB / RPC)必须有超时,不能阻塞 poll 线程
  4. CooperativeSticky 分配策略,减少 Rebalance 时的 STW

相关参数

  • max.poll.interval.ms(Consumer,默认 300_000)
  • max.poll.records(Consumer,默认 500)
  • session.timeout.ms / heartbeat.interval.ms
  • partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor

案例 6:unclean.leader.election.enable=true 导致已 ack 消息丢失

症状

  • 风控团队发现某些「已经 ack 给上游」的消息消费不到
  • 监控告警 UncleanLeaderElectionsPerSec > 0
  • 几分钟前刚发生过一次 Broker 大面积宕机

现场快照

broker.properties:
unclean.leader.election.enable=true   ← 误配置!

事件时间线:
14:00 broker-2(Leader)宕机
14:00 ISR 只剩 broker-2 一个(原本 broker-1 早就掉队,broker-3 也掉队)
14:00 此时 ISR = [], 但 unclean=true → Controller 选 broker-1 当 Leader(虽然它不在 ISR)
14:01 broker-1 接手,但它的 LEO 比 broker-2 少了 50万条
14:01 这 50万条消息在「Leader 切换瞬间」永久丢失
14:02 业务报警,对账少了 50万订单

根因分析

  • unclean.leader.election.enable=true 允许从「不在 ISR 的副本」中选 Leader
  • 这种副本没有最新数据,选成 Leader 后高水位回退,已 ack 消息消失
  • 0.11+ 默认是 false,老集群升级时常被遗留为 true

解决方案

  • 立刻关掉:kafka-configs.sh --alter --add-config unclean.leader.election.enable=false --entity-type topics --entity-name xxx
  • 全集群默认值改为 false(broker.properties)
  • 已丢的消息无法恢复,只能业务侧补救(人工对账 / 重新生成)

预防措施

  1. 新建 Topic 默认 false
  2. 监控 UncleanLeaderElectionsPerSec,一旦 > 0 即 P0 告警
  3. 极端情况下「宁可分区不可用,也不能丢数据」是默认原则
  4. 配合 min.insync.replicas=2 + acks=all,从源头避免 ISR 收缩

相关参数

  • unclean.leader.election.enable=false(Broker 全局 + Topic 级别)

案例 7:大消息撑爆 message.max.bytes

症状

  • 业务突然增加了「附件全文」字段,单条消息从 10KB 飙到 5MB
  • Producer 报错 RecordTooLargeException: The message is 5242880 bytes when serialized which is larger than the maximum...
  • 部分老消息成功写入,但Consumer 拉不出来The message is too large to receive
  • Broker 副本同步失败:Replica fetch failed: message exceeds replica.fetch.max.bytes

现场快照

配置(默认值):
message.max.bytes        = 1048588   (~ 1MB) Broker
max.request.size         = 1048576   (~ 1MB) Producer
fetch.max.bytes          = 52428800  (~ 50MB) Consumer
max.partition.fetch.bytes = 1048576  (~ 1MB) Consumer
replica.fetch.max.bytes  = 1048576   (~ 1MB) Broker 副本同步

任何一处 < 实际消息大小,链路就断。

根因分析

  • Kafka 消息大小有 4 道闸门:Producer → Broker → Replica 同步 → Consumer
  • 任何一道闸门小于实际消息都会失败
  • 而且 replica.fetch.max.bytes 不够时,单个大消息永远追不上,会卡住 Follower

解决方案

  • 改 4 个参数(保守翻倍到 10MB):
    # Broker
    message.max.bytes=10485760
    replica.fetch.max.bytes=10485760
    
    # Topic 级别可单独覆盖
    kafka-configs.sh --alter --entity-type topics --entity-name big-topic \
      --add-config max.message.bytes=10485760
    
    # Producer
    max.request.size=10485760
    
    # Consumer
    fetch.max.bytes=104857600   # 至少 message.max.bytes × 10
    max.partition.fetch.bytes=10485760
  • 更优雅的方案:分片(见 code/large_message_handle.py

预防措施

  1. Kafka 不是文件传输系统,单条消息控制在 < 1MB
  2. 大文件用「Kafka + 对象存储」模式:消息只放 URL,文件存 S3 / OSS
  3. 真要传大消息,先开启压缩compression.type=zstd
  4. 上线前检查所有客户端的 max.request.size / fetch.max.bytes

相关参数

参数
Brokermessage.max.bytes, replica.fetch.max.bytes
Topicmax.message.bytes
Producermax.request.size, compression.type
Consumerfetch.max.bytes, max.partition.fetch.bytes

案例 8:Compaction 不生效

症状

  • 用 Compaction Topic 存「用户最新状态」,期望每个 user_id 只保留最新值
  • 跑了一周,磁盘占用还在线性增长,旧值没被清理
  • kafka-dump-log 看到大量重复 Key

现场快照

bash
kafka-configs.sh --describe --entity-type topics --entity-name user.state
# cleanup.policy=compact
# min.cleanable.dirty.ratio=0.5  ← 50% 才触发
# segment.ms=604800000           ← 7 天

du -sh /data/kafka/user.state-0
# 250G  ← 一直涨

log-cleaner.log

ERROR Cleaner thread 0 panicked: java.lang.OutOfMemoryError: Java heap space

根因分析(4 类)

  1. Cleaner 线程挂了(OOM):堆给 Cleaner 的 buffer 不够
  2. Active Segment 不参与 Compaction:当前正在写的段不会被清理
  3. min.cleanable.dirty.ratio 太高:默认 0.5,要等脏数据占一半才动手
  4. segment.ms 太长:段不滚动,Active Segment 占绝大部分数据

解决方案

bash
# 1. 调小 dirty ratio,让 Compaction 更勤快
kafka-configs.sh --alter --entity-type topics --entity-name user.state \
  --add-config min.cleanable.dirty.ratio=0.1

# 2. 让段更频繁滚动
kafka-configs.sh --alter --entity-type topics --entity-name user.state \
  --add-config segment.ms=3600000   # 1 小时滚一次

# 3. 给 Cleaner 更多内存(broker 配置)
log.cleaner.dedupe.buffer.size=2000000000   # 2GB
log.cleaner.threads=2

预防措施

  1. Compaction Topic 必须监控 log-cleaner.logLogCleanerManager JMX
  2. segment.ms 要小(小时级),让段尽快从 active 变成 inactive
  3. 写 Tombstone (null value) 实现「删除」时,注意 delete.retention.ms(默认 24h,超过才真删)

相关参数

  • cleanup.policy=compact
  • min.cleanable.dirty.ratio(默认 0.5)
  • segment.ms / segment.bytes
  • log.cleaner.threads
  • log.cleaner.dedupe.buffer.size
  • delete.retention.ms

案例 9:Connect Offset Topic 配置漂移导致 Connector 重置

症状

  • 新搭一个 Connect Worker,跑了几天突然所有 Source Connector 从头重传
  • ClickHouse 收到几亿条重复数据
  • 看 Connect 日志:Restoring offsets for connector: starting from beginning

现场快照

Worker 1 配置(旧):
offset.storage.topic=connect-offsets
offset.storage.replication.factor=3
offset.storage.partitions=25

Worker 2 配置(新):
offset.storage.topic=connect-offsets-v2   ← 不一样!

根因分析

  • Connect Worker 用 Kafka Topic 存「Source Connector 的进度」
  • 多个 Worker 必须指向同一个 Topic才能共享进度
  • 新 Worker 用了不同的 Topic,里面是空的 → 所有 Connector 视为从未跑过 → 重头开始
  • 已有的 Sink Connector 也会因 consumer.group.id 不一致重新消费

解决方案

  • 立刻把 Worker 2 配置改回与 Worker 1 一致并重启
  • 已经重发的数据:依赖下游幂等(ClickHouse ReplacingMergeTree / 数据库唯一约束)兜底

预防措施

  1. Connect Worker 配置用同一个 properties 文件分发,禁止手抄
  2. 关键 Topic 名(offset/config/status)放进配置中心或 Helm Chart
  3. 新 Worker 上线前 dry-run:检查它要连接的 3 个 Topic 是否已存在并配置正确

相关参数(Connect Worker 必须一致)

  • offset.storage.topicoffset.storage.replication.factoroffset.storage.partitions
  • config.storage.topicconfig.storage.replication.factor
  • status.storage.topicstatus.storage.replication.factorstatus.storage.partitions
  • group.id(多个 Worker 必须相同)

案例 10:跨机房网络抖动 → ISR 频繁收缩

症状

  • 集群部署在 2 个机房(北京 + 上海),副本跨机房分布
  • ISRShrinksPerSecISRExpandsPerSec 同时高
  • Producer 偶发超时,Consumer 没事
  • 跨机房 RTT 偶尔 100ms+

现场快照

server.log:
[01:23:45] Shrinking ISR for partition orders-2 from (1,2,3) to (1,2)
[01:23:48] Expanding ISR for partition orders-2 from (1,2) to (1,2,3)
[01:24:01] Shrinking ISR for partition orders-2 from (1,2,3) to (1,2)
[01:24:04] Expanding ...

根因分析

  • replica.lag.time.max.ms 默认 30 秒
  • 跨机房 100ms RTT 下,瞬间网络抖动可能让 Follower 30 秒内追不上 → 被踢
  • 抖动结束后又能追上 → 加回
  • 频繁 Shrink/Expand 本身会触发元数据更新,对 Controller 有压力

解决方案

  • 调大 replica.lag.time.max.ms 到 60 秒(或 90 秒),容忍短暂抖动
  • 调大 replica.fetch.max.wait.ms,减少跨机房拉取频次
  • 副本布局优化:Leader 优先放本机房,Follower 跨机房(用 replica.selector.class 实现就近读)

预防措施

  1. 跨机房集群要专门测试 副本同步在网络抖动下的鲁棒性
  2. 关键 Topic 选「3 副本同机房 + 1 副本跨机房」的不对称布局
  3. 跨机房推荐用 MirrorMaker 2 异步复制,而非同集群多机房
  4. 关注 RequestQueueTimeMs p99,看抖动是否堆积请求

相关参数

  • replica.lag.time.max.ms(默认 30000)
  • replica.fetch.max.wait.ms(默认 500)
  • replica.fetch.min.bytes
  • num.replica.fetchers(每 Broker 拉取线程数)

案例 11:EOS 滥用拖垮吞吐

症状

  • 业务追求「不丢不重」,全量打开事务
  • Producer TPS 从 100K 掉到 5K,p99 延迟从 20ms 升到 800ms
  • 监控 __transaction_state Topic 写入暴涨

现场快照

python
# 每条消息一个事务!
for event in stream:
    producer.beginTransaction()
    producer.send("topic", event)
    producer.commitTransaction()    # ← 每次都同步 2PC

根因分析

  • 事务 Producer 走两阶段提交:beginTxn → send → commitTxn
  • commitTxn 要写 __transaction_state 主从复制 + 标记所有相关分区,至少 4 次网络往返
  • 每条消息一事务 = 每条消息 4 次额外 RTT,吞吐暴跌

解决方案

  • 批量事务:100~1000 条消息一个事务
    python
    producer.beginTransaction()
    for i, event in enumerate(batch):
        producer.send("topic", event)
        if (i+1) % 500 == 0:
            producer.commitTransaction()
            producer.beginTransaction()
    producer.commitTransaction()
  • 能不用事务就不用:99% 场景下 enable.idempotence=true + 业务幂等表 已经够
  • 如果真要 EOS,processing.guarantee=exactly_once_v2 性能远好于 v1

预防措施

  1. 事务 = 性能 ↓ + 复杂度 ↑,先评估是否真的需要
  2. 用 EOS 的场景:扣款 / 库存 / 计费等「绝不能重复」的业务
  3. 用 At-Least-Once + 幂等的场景:日志、行为埋点、绝大多数业务
  4. 事务 Topic 副本数 ≥ 3,因为它是 EOS 链路的命脉

相关参数

  • enable.idempotence(默认 false,强烈建议 true
  • transactional.id(事务 Producer 必填)
  • transaction.timeout.ms(默认 60s)
  • processing.guarantee(Streams 配置:at_least_once / exactly_once_v2)

案例 12:__consumer_offsets 副本不健康导致整组 Rebalance 失败

症状

  • 所有消费组无法重新加入,新 Consumer 启动卡在 JoinGroup
  • 老 Consumer 报 CoordinatorNotAvailableException
  • __consumer_offsets 某些分区 UnderReplicated

现场快照

bash
kafka-topics.sh --describe --topic __consumer_offsets | head -5
# Topic: __consumer_offsets  Partition: 12
# Leader: -1  Replicas: 1,2,3  Isr: -
# ↑ Leader = -1 → 该分区无 Leader!

某 Consumer Group 的 hash 正好落在分区 12 → 它的 Coordinator 不可用 → Group 无法运作。

根因分析

  • 消费组的 Coordinator 由 __consumer_offsets 分区的 Leader 担任
  • 该分区的所有副本所在 Broker 都宕了 / 磁盘满了
  • 没有 Leader → 没有 Coordinator → JoinGroup 永远失败

解决方案

  • 立刻拉起对应 Broker(Replicas 列里那 3 个)
  • 暂时不可恢复时:临时改 unclean.leader.election.enable=true(注意会丢部分 Offset,下次启动从 latest 开始)
  • 长期:把 __consumer_offsets 副本数调到 3+,分布在不同机架

预防措施

  1. offsets.topic.replication.factor=3 必须配置(默认 1,灾难配置)
  2. 上线时检查 __consumer_offsets 是否真的 3 副本:
    bash
    kafka-topics.sh --describe --topic __consumer_offsets | grep -v "Replicas: 1,2,3"
  3. 每个 Broker 要打 broker.rack,让 Kafka 跨机架放副本
  4. __transaction_state 同样问题,配置 transaction.state.log.replication.factor=3

相关参数

  • offsets.topic.replication.factor(默认 1,生产必须 ≥ 3
  • offsets.topic.num.partitions(默认 50)
  • transaction.state.log.replication.factor(默认 1,必须改 3)
  • min.insync.replicas(针对内部 Topic)

常见错误日志辞典

错误信息含义处理
NotLeaderForPartitionException客户端发到的分区已不是 LeaderProducer/Consumer 自动刷新元数据重发,频繁出现就查 Leader 切换原因
NotEnoughReplicasExceptionISR 数 < min.insync.replicas检查副本健康;不可降级时报 acks=1
RecordTooLargeException消息超出 max.request.size / message.max.bytes调参或分片,见案例 7
OffsetOutOfRangeExceptionOffset 在 Topic 范围外(已被清理或越界)auto.offset.reset 处理;评估是不是 retention 太短
CommitFailedExceptionOffset 提交失败(超时 / Group 已 Rebalance)减小 batch、加 max.poll.interval.ms,避免 poll 间隔过长
CoordinatorNotAvailableException__consumer_offsets 分区 Leader 不可用见案例 12
UnknownTopicOrPartitionExceptionTopic 不存在 / 已被删检查 Topic 名拼写 / auto.create.topics.enable
ProducerFencedExceptiontransactional.id 已被新 Producer 接管旧实例必须立即退出,避免双写
InvalidPidMappingExceptionPID 失效(事务过期 / 元数据丢失)重启 Producer 重新申请 PID
Disconnected during request to broker xxx与 Broker 断连网络抖动 / Broker 宕机;查 server.log
Discovered new coordinatorCoordinator 切换通常是上一个挂了;看 __consumer_offsets 是否健康
OutOfMemoryError: Java heap spaceJVM 堆 OOM检查 Cleaner / JMX 内存配置;分析 GC 日志
Bad config value配置非法看具体字段,对照官方文档
Marking the coordinator dead心跳超时网络 / Coordinator 宕;多次出现做 Rebalance

本章面试高频题

Q1:acks=1acks=all 在故障下分别会发生什么?

答案

  • acks=1:Leader 写完 ack。Leader 在 ack 之后、Follower 同步之前宕机 → 数据丢失(旧 Leader 重启会被截断)
  • acks=all:等所有 ISR 都写完才 ack。只要 ISR 中至少 min.insync.replicas 个副本存活,数据不丢
  • 副作用:acks=all 延迟更高(多一次跨 Broker 同步),但通常 < 30%

加分:能讲清楚 acks=all + min.insync.replicas=2 + enable.idempotence=true 三件套是「不丢不重」的标配。


Q2:什么是 Rebalance 风暴?怎么避免?

答案

  • 现象:消费组无限循环 Rebalance,每次都因 max.poll.interval.ms 超时被踢
  • 根因:业务处理慢,poll 间隔超阈值
  • 避免:
    1. 调大 max.poll.interval.ms
    2. 调小 max.poll.records,缩短单次处理时间
    3. 业务异步化:poll 线程只入队,独立 worker 处理
    4. CooperativeStickyAssignor(增量再平衡,减小 STW)
    5. 用 Static Membership(group.instance.id)避免短暂抖动触发

加分:能解释 Eager Rebalance vs Cooperative 的差异。


Q3:unclean.leader.election.enable=true 为什么危险?

答案

  • 允许从「不在 ISR」的副本中选 Leader
  • 这些副本数据落后,被选成 Leader 后 HW 回退,已 ack 消息消失
  • 对要求不丢数据的业务是灾难
  • 配置原则:min.insync.replicas=2 + acks=all + unclean=false 三件套

Q4:单条消息 5MB 想发到 Kafka,需要改哪些参数?

答案表

参数改成
Brokermessage.max.bytes10MB
Brokerreplica.fetch.max.bytes10MB(必须 ≥ message.max.bytes)
Topicmax.message.bytes10MB(覆盖 Broker 全局)
Producermax.request.size10MB
Producercompression.typezstd(建议)
Consumerfetch.max.bytes100MB(建议 ≥ message.max.bytes × 10)
Consumermax.partition.fetch.bytes10MB

更优解:消息存 URL,文件存 OSS(Kafka 不擅长大消息)。


Q5:消费 Lag 一直涨,怎么排查?

(这道题与第 19 章 Q2 重复,但踩坑视角补充)

踩坑维度的 6 类原因

  1. 分区数太少 → 加分区(注意 Key 顺序)
  2. Key 设计热点 → 单分区 Lag 高 → 重新设计 Key
  3. 业务卡住(DB 慢查询、外部 RPC 超时)→ 看应用日志
  4. Rebalance 风暴 → 改 max.poll.interval.ms
  5. Consumer 实例数 > 分区数 → 多余的实例空跑,分区数才是真正的并行度上限
  6. __consumer_offsets 副本异常 → 提交不上去,反复重消费

小结

  • Kafka 的坑 80% 出在配置,剩下 20% 出在 Key 设计与业务处理
  • 三个最常见根因
    1. 副本与 ack 配置不当(案例 3、6、12)
    2. 消费者超时配置不当(案例 5、4)
    3. 分区与 Key 设计不当(案例 1、2)
  • 每次配置变更都要问自己:「故障场景下会发生什么?」
  • 每个生产 Kafka 集群都应该做的 5 件事
    1. __consumer_offsets 副本 ≥ 3
    2. 关键 Topic min.insync.replicas=2 + acks=all
    3. unclean.leader.election.enable=false
    4. JVM 6-8 GB + G1
    5. JMX → Prometheus → Grafana 全套监控

下一章我们把 Schema、监控、踩坑经验全部串起来,落地一个真实可跑的「实时订单总线 + CDC 数仓链路」

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
#!/usr/bin/env python3
"""
热点 Key 模拟器:复现「单分区热点」踩坑(案例 2)。

3 种模式:
  good     : 用 user_id (高基数),分布均匀
  bad_date : 用 today() 做 key,所有消息进同一分区
  bad_city : 5 个城市做 key,5 个分区被打爆,其余空闲

跑完会打印每个分区的消息数,让你直观看到倾斜。

依赖:pip install confluent-kafka
前置:先建一个 12 分区的 Topic
  kafka-topics.sh --create --bootstrap-server localhost:9092 \
    --topic learn.20.heavy --partitions 12 --replication-factor 1

运行:
  MODE=bad_date python heavy_key_simulator.py
  MODE=good     python heavy_key_simulator.py
  MODE=bad_city python heavy_key_simulator.py
"""

import os
import random
import time
from collections import Counter
from datetime import date

from confluent_kafka import Producer, Consumer, TopicPartition

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
TOPIC = os.getenv("TOPIC", "learn.20.heavy")
MODE = os.getenv("MODE", "bad_date")  # good / bad_date / bad_city
COUNT = int(os.getenv("COUNT", "10000"))

CITIES = ["beijing", "shanghai", "shenzhen", "hangzhou", "chengdu"]


def gen_key(i):
    if MODE == "good":
        return f"user-{random.randint(1, 100000)}"
    if MODE == "bad_date":
        return str(date.today())
    if MODE == "bad_city":
        return random.choice(CITIES)
    raise ValueError(MODE)


def produce():
    p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 10})
    print(f"模式 = {MODE}, 发送 {COUNT} 条消息到 {TOPIC} ...")
    t0 = time.time()
    for i in range(COUNT):
        k = gen_key(i)
        p.produce(TOPIC, key=k.encode(), value=f"msg-{i}".encode())
        if (i + 1) % 1000 == 0:
            p.poll(0)
    p.flush(20)
    print(f"发送完成,耗时 {time.time()-t0:.1f}s")


def show_distribution():
    """从 Topic 元数据读分区数,遍历所有分区算 hi-lo 算消息数。"""
    c = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "heavy-key-stat",
        "enable.auto.commit": False,
    })
    md = c.list_topics(TOPIC, timeout=10)
    parts = list(md.topics[TOPIC].partitions.keys())
    print(f"\n分区数 = {len(parts)}")
    counter = Counter()
    for pid in parts:
        lo, hi = c.get_watermark_offsets(TopicPartition(TOPIC, pid), timeout=5)
        counter[pid] = hi - lo
    c.close()

    total = sum(counter.values())
    print(f"\n=== 分区消息数分布 (共 {total} 条) ===")
    print(f"{'分区':>4} | {'消息数':>10} | {'占比':>8} | 直方图")
    print("-" * 60)
    for pid in sorted(parts):
        n = counter[pid]
        pct = n / total * 100 if total else 0
        bar = "█" * int(pct / 2)
        print(f"{pid:>4} | {n:>10} | {pct:>7.2f}% | {bar}")

    # 倾斜度
    if total:
        mx = max(counter.values())
        avg = total / len(parts)
        print(f"\n最大/平均 = {mx/avg:.2f}x" + (" 严重倾斜!" if mx/avg > 3 else ""))


if __name__ == "__main__":
    produce()
    show_distribution()
python
#!/usr/bin/env python3
"""
大消息分片发送范式(应对案例 7)。

思路:
- 把 10MB 的大消息切成多个 ≤ 512KB 的分片
- 给每片打上同一个 message_id 与 (idx, total)
- 用 message_id 作为 Kafka Key,保证所有分片落同一分区(顺序)
- 消费端按 message_id 收齐后再 reassemble

依赖:pip install confluent-kafka
"""

import os
import json
import uuid
import hashlib
from collections import defaultdict
from confluent_kafka import Producer, Consumer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
TOPIC = os.getenv("TOPIC", "learn.20.large")
CHUNK_SIZE = 512 * 1024   # 512KB / 片


# =================== Producer 端 ===================
def send_large(producer: Producer, payload: bytes, meta: dict = None):
    """把大 payload 切片发送。"""
    msg_id = uuid.uuid4().hex
    digest = hashlib.md5(payload).hexdigest()
    total = (len(payload) + CHUNK_SIZE - 1) // CHUNK_SIZE

    for idx in range(total):
        chunk = payload[idx * CHUNK_SIZE:(idx + 1) * CHUNK_SIZE]
        envelope = {
            "msg_id": msg_id,
            "idx": idx,
            "total": total,
            "digest": digest,         # MD5 校验
            "size": len(payload),
            "meta": meta or {},
            "data_b64": chunk.hex(),  # 用 hex 简化(生产可用 base64 + 直接 bytes)
        }
        producer.produce(
            topic=TOPIC,
            key=msg_id.encode(),       # 同一 message_id 保证同分区有序
            value=json.dumps(envelope).encode(),
        )
        producer.poll(0)
    producer.flush(10)
    print(f"  ✓ 发送 msg_id={msg_id}{total} 片,原大小 {len(payload)} bytes")


# =================== Consumer 端 ===================
def consume_large(consumer: Consumer):
    """收齐所有分片后 reassemble。"""
    buffers = defaultdict(dict)   # msg_id -> {idx: chunk}
    metas = {}

    print("等待消息 ... Ctrl+C 退出")
    while True:
        msg = consumer.poll(1.0)
        if msg is None or msg.error():
            continue

        env = json.loads(msg.value())
        mid = env["msg_id"]
        buffers[mid][env["idx"]] = bytes.fromhex(env["data_b64"])
        if mid not in metas:
            metas[mid] = env

        # 收齐了?
        if len(buffers[mid]) == env["total"]:
            data = b"".join(buffers[mid][i] for i in sorted(buffers[mid]))
            ok = hashlib.md5(data).hexdigest() == env["digest"]
            print(f"  ✓ 收齐 msg_id={mid} size={len(data)} 校验={'OK' if ok else 'FAIL'}")
            del buffers[mid]
            del metas[mid]
            consumer.commit(msg)
            yield data, env["meta"]


# =================== 演示入口 ===================
def demo():
    p = Producer({
        "bootstrap.servers": BOOTSTRAP,
        "compression.type": "zstd",   # 大消息一定开压缩
        "enable.idempotence": True,
        "acks": "all",
    })

    # 造一个 5MB 的「大消息」
    big = b"X" * (5 * 1024 * 1024)
    print("Sending one 5MB message via chunking ...")
    send_large(p, big, meta={"filename": "huge.bin"})


if __name__ == "__main__":
    if os.getenv("ROLE", "producer") == "consumer":
        c = Consumer({
            "bootstrap.servers": BOOTSTRAP,
            "group.id": "large-msg-demo",
            "auto.offset.reset": "earliest",
            "enable.auto.commit": False,
            "fetch.max.bytes": 10 * 1024 * 1024,
        })
        c.subscribe([TOPIC])
        try:
            for data, meta in consume_large(c):
                pass
        except KeyboardInterrupt:
            c.close()
    else:
        demo()
python
#!/usr/bin/env python3
"""
Rebalance 风暴复现脚本(案例 5)。

操作:
1. 起 N 个 Consumer,故意把每条消息处理设置为 5 分钟(> max.poll.interval.ms 默认 5 分钟)
2. 你会观察到:每个 Consumer 都在第 5 分钟左右被踢,触发 Rebalance
3. Rebalance 后分区重新分配,新接手的 Consumer 也会因相同原因被踢
4. 整个 group 永远在 Rebalance 状态

可调环境变量:
    PROCESS_SECS=400        # 每条消息处理时间(秒),> 300 必现风暴
    MAX_POLL_MS=300000      # max.poll.interval.ms(毫秒)
    INSTANCE_ID=worker-1    # 多终端起多实例

修复方式(取消注释 fix block):
- PROCESS_SECS 调小 < MAX_POLL_MS / 1000
- 改用 CooperativeStickyAssignor 减小 STW
- max.poll.records 调小

依赖:pip install confluent-kafka
"""

import os
import time
from confluent_kafka import Consumer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
TOPIC = os.getenv("TOPIC", "learn.20.storm")
GROUP = os.getenv("GROUP", "storm-group")
INSTANCE_ID = os.getenv("INSTANCE_ID", f"worker-{os.getpid()}")
PROCESS_SECS = int(os.getenv("PROCESS_SECS", "400"))
MAX_POLL_MS = int(os.getenv("MAX_POLL_MS", "300000"))


def on_assign(c, parts):
    print(f"  [{INSTANCE_ID}] ASSIGN -> {[(p.topic, p.partition) for p in parts]}")

def on_revoke(c, parts):
    print(f"  [{INSTANCE_ID}] REVOKE <- {[(p.topic, p.partition) for p in parts]}")


def main():
    cfg = {
        "bootstrap.servers": BOOTSTRAP,
        "group.id": GROUP,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
        "max.poll.interval.ms": MAX_POLL_MS,
        "session.timeout.ms": 45000,
        # === 复现关键点:单次 poll 处理 PROCESS_SECS 秒 ===
        # === 修复方案 1: 改成 CooperativeStickyAssignor ===
        # "partition.assignment.strategy": "cooperative-sticky",
        # === 修复方案 2: 用 group.instance.id 静态成员,避免短抖动触发 ===
        # "group.instance.id": INSTANCE_ID,
    }
    c = Consumer(cfg)
    c.subscribe([TOPIC], on_assign=on_assign, on_revoke=on_revoke)

    print(f"[{INSTANCE_ID}] start, MAX_POLL={MAX_POLL_MS}ms, PROCESS={PROCESS_SECS}s")
    print(f"[{INSTANCE_ID}] 预期:处理时间 > max.poll.interval.ms 时会被踢出 group")

    while True:
        msg = c.poll(1.0)
        if msg is None or msg.error():
            continue

        print(f"[{INSTANCE_ID}] consume p={msg.partition()} off={msg.offset()} → 模拟处理 {PROCESS_SECS}s")
        # ↓↓↓ 故意慢处理 ↓↓↓
        time.sleep(PROCESS_SECS)

        try:
            c.commit(msg)
            print(f"[{INSTANCE_ID}] commit OK")
        except Exception as e:
            print(f"[{INSTANCE_ID}] commit FAIL: {e}  ← 已经被踢出 group")


if __name__ == "__main__":
    try:
        main()
    except KeyboardInterrupt:
        pass

heavy_key_simulator.py ↗ · large_message_handle.py ↗ · rebalance_storm_repro.py ↗