主题
第 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 消费完后下线
预防措施
- 分区数一开始就预估到未来 2-3 年的峰值,宁多勿少
- 单 Broker 推荐 100-2000 分区上限;按 Broker 数 × 200 估算
- 用「自定义 Partitioner」而非 hash mod,例如
hash(key) % 1024 → 路由表 → 真实分区,加分区时只改路由表 - 关键 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)
预防措施
- Key 必须高基数(cardinality),至少 ≥ 分区数 × 100
- 禁止用 时间、固定 enum、城市名、机型名 这种低基数字段做 Key
- 上线前压测时专门检查「分区流量分布」(用
kafka-consumer-groups.sh --describe看每分区 Lag) - 监控告警:分区流量标准差 > 平均值 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
预防措施
- 重要业务一律
acks=all,吞吐损失通常 < 20% - 滚动重启 Broker 时控制速度(一台一台),并监控 Producer 错误率
- 用「业务对账」兜底:定期 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重置到最早 - 重新启动消费者,把历史补回来
预防措施
- 新接入业务先评估「是否需要历史」:
- 实时风控、统计:通常需要历史 →
earliest - 实时通知、排行榜:通常不需要 →
latest
- 实时风控、统计:通常需要历史 →
- 关键业务用
none:强制人工设置初始 offset,避免静默错误 - 消费组首次上线后立即检查 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 永不阻塞
预防措施
- 理解
max.poll.interval.ms是「业务处理时间上限」,不是网络超时 - 批量大小 × 单条处理时间 <
max.poll.interval.ms的 70% - 慢操作(DB / RPC)必须有超时,不能阻塞 poll 线程
- 用 CooperativeSticky 分配策略,减少 Rebalance 时的 STW
相关参数
max.poll.interval.ms(Consumer,默认 300_000)max.poll.records(Consumer,默认 500)session.timeout.ms/heartbeat.interval.mspartition.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) - 已丢的消息无法恢复,只能业务侧补救(人工对账 / 重新生成)
预防措施
- 新建 Topic 默认
false - 监控
UncleanLeaderElectionsPerSec,一旦 > 0 即 P0 告警 - 极端情况下「宁可分区不可用,也不能丢数据」是默认原则
- 配合
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)
预防措施
- Kafka 不是文件传输系统,单条消息控制在 < 1MB
- 大文件用「Kafka + 对象存储」模式:消息只放 URL,文件存 S3 / OSS
- 真要传大消息,先开启压缩(
compression.type=zstd) - 上线前检查所有客户端的
max.request.size/fetch.max.bytes
相关参数
| 端 | 参数 |
|---|---|
| Broker | message.max.bytes, replica.fetch.max.bytes |
| Topic | max.message.bytes |
| Producer | max.request.size, compression.type |
| Consumer | fetch.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 类)
- Cleaner 线程挂了(OOM):堆给 Cleaner 的 buffer 不够
- Active Segment 不参与 Compaction:当前正在写的段不会被清理
min.cleanable.dirty.ratio太高:默认 0.5,要等脏数据占一半才动手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预防措施
- Compaction Topic 必须监控
log-cleaner.log与LogCleanerManagerJMX segment.ms要小(小时级),让段尽快从 active 变成 inactive- 写 Tombstone (null value) 实现「删除」时,注意
delete.retention.ms(默认 24h,超过才真删)
相关参数
cleanup.policy=compactmin.cleanable.dirty.ratio(默认 0.5)segment.ms/segment.byteslog.cleaner.threadslog.cleaner.dedupe.buffer.sizedelete.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/ 数据库唯一约束)兜底
预防措施
- Connect Worker 配置用同一个 properties 文件分发,禁止手抄
- 关键 Topic 名(offset/config/status)放进配置中心或 Helm Chart
- 新 Worker 上线前 dry-run:检查它要连接的 3 个 Topic 是否已存在并配置正确
相关参数(Connect Worker 必须一致)
offset.storage.topic、offset.storage.replication.factor、offset.storage.partitionsconfig.storage.topic、config.storage.replication.factorstatus.storage.topic、status.storage.replication.factor、status.storage.partitionsgroup.id(多个 Worker 必须相同)
案例 10:跨机房网络抖动 → ISR 频繁收缩
症状
- 集群部署在 2 个机房(北京 + 上海),副本跨机房分布
ISRShrinksPerSec与ISRExpandsPerSec同时高- 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实现就近读)
预防措施
- 跨机房集群要专门测试 副本同步在网络抖动下的鲁棒性
- 关键 Topic 选「3 副本同机房 + 1 副本跨机房」的不对称布局
- 跨机房推荐用 MirrorMaker 2 异步复制,而非同集群多机房
- 关注
RequestQueueTimeMsp99,看抖动是否堆积请求
相关参数
replica.lag.time.max.ms(默认 30000)replica.fetch.max.wait.ms(默认 500)replica.fetch.min.bytesnum.replica.fetchers(每 Broker 拉取线程数)
案例 11:EOS 滥用拖垮吞吐
症状
- 业务追求「不丢不重」,全量打开事务
- Producer TPS 从 100K 掉到 5K,p99 延迟从 20ms 升到 800ms
- 监控
__transaction_stateTopic 写入暴涨
现场快照
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
预防措施
- 事务 = 性能 ↓ + 复杂度 ↑,先评估是否真的需要
- 用 EOS 的场景:扣款 / 库存 / 计费等「绝不能重复」的业务
- 用 At-Least-Once + 幂等的场景:日志、行为埋点、绝大多数业务
- 事务 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+,分布在不同机架
预防措施
offsets.topic.replication.factor=3必须配置(默认 1,灾难配置)- 上线时检查
__consumer_offsets是否真的 3 副本:bashkafka-topics.sh --describe --topic __consumer_offsets | grep -v "Replicas: 1,2,3" - 每个 Broker 要打
broker.rack,让 Kafka 跨机架放副本 __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 | 客户端发到的分区已不是 Leader | Producer/Consumer 自动刷新元数据重发,频繁出现就查 Leader 切换原因 |
NotEnoughReplicasException | ISR 数 < min.insync.replicas | 检查副本健康;不可降级时报 acks=1 |
RecordTooLargeException | 消息超出 max.request.size / message.max.bytes | 调参或分片,见案例 7 |
OffsetOutOfRangeException | Offset 在 Topic 范围外(已被清理或越界) | 按 auto.offset.reset 处理;评估是不是 retention 太短 |
CommitFailedException | Offset 提交失败(超时 / Group 已 Rebalance) | 减小 batch、加 max.poll.interval.ms,避免 poll 间隔过长 |
CoordinatorNotAvailableException | __consumer_offsets 分区 Leader 不可用 | 见案例 12 |
UnknownTopicOrPartitionException | Topic 不存在 / 已被删 | 检查 Topic 名拼写 / auto.create.topics.enable |
ProducerFencedException | 同 transactional.id 已被新 Producer 接管 | 旧实例必须立即退出,避免双写 |
InvalidPidMappingException | PID 失效(事务过期 / 元数据丢失) | 重启 Producer 重新申请 PID |
Disconnected during request to broker xxx | 与 Broker 断连 | 网络抖动 / Broker 宕机;查 server.log |
Discovered new coordinator | Coordinator 切换 | 通常是上一个挂了;看 __consumer_offsets 是否健康 |
OutOfMemoryError: Java heap space | JVM 堆 OOM | 检查 Cleaner / JMX 内存配置;分析 GC 日志 |
Bad config value | 配置非法 | 看具体字段,对照官方文档 |
Marking the coordinator dead | 心跳超时 | 网络 / Coordinator 宕;多次出现做 Rebalance |
本章面试高频题
Q1:acks=1 和 acks=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 间隔超阈值
- 避免:
- 调大
max.poll.interval.ms - 调小
max.poll.records,缩短单次处理时间 - 业务异步化:poll 线程只入队,独立 worker 处理
- 用
CooperativeStickyAssignor(增量再平衡,减小 STW) - 用 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,需要改哪些参数?
答案表:
| 端 | 参数 | 改成 |
|---|---|---|
| Broker | message.max.bytes | 10MB |
| Broker | replica.fetch.max.bytes | 10MB(必须 ≥ message.max.bytes) |
| Topic | max.message.bytes | 10MB(覆盖 Broker 全局) |
| Producer | max.request.size | 10MB |
| Producer | compression.type | zstd(建议) |
| Consumer | fetch.max.bytes | 100MB(建议 ≥ message.max.bytes × 10) |
| Consumer | max.partition.fetch.bytes | 10MB |
更优解:消息存 URL,文件存 OSS(Kafka 不擅长大消息)。
Q5:消费 Lag 一直涨,怎么排查?
(这道题与第 19 章 Q2 重复,但踩坑视角补充)
踩坑维度的 6 类原因:
- 分区数太少 → 加分区(注意 Key 顺序)
- Key 设计热点 → 单分区 Lag 高 → 重新设计 Key
- 业务卡住(DB 慢查询、外部 RPC 超时)→ 看应用日志
- Rebalance 风暴 → 改
max.poll.interval.ms - Consumer 实例数 > 分区数 → 多余的实例空跑,分区数才是真正的并行度上限
__consumer_offsets副本异常 → 提交不上去,反复重消费
小结
- Kafka 的坑 80% 出在配置,剩下 20% 出在 Key 设计与业务处理
- 三个最常见根因:
- 副本与 ack 配置不当(案例 3、6、12)
- 消费者超时配置不当(案例 5、4)
- 分区与 Key 设计不当(案例 1、2)
- 每次配置变更都要问自己:「故障场景下会发生什么?」
- 每个生产 Kafka 集群都应该做的 5 件事:
__consumer_offsets副本 ≥ 3- 关键 Topic
min.insync.replicas=2 + acks=all unclean.leader.election.enable=false- JVM 6-8 GB + G1
- 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:
passheavy_key_simulator.py ↗ · large_message_handle.py ↗ · rebalance_storm_repro.py ↗