Kafka 生产事故:发消息超时、堆积、消费延迟排查
提出问题
Kafka 靠高吞吐和可靠性标签赢得无数生产部署,但跑起来之后,该踩的坑一个都不会少。发消息超时导致业务链路中断、Consumer Lag 飙到几十万、端到端延迟从毫秒级跳到秒级——这些不是文档里看得到的,而是真实的线上事故。
更棘手的是,Kafka 的故障现象往往很相似,但根因完全相反:同样是 TimeoutException,可能是磁盘满了、ISR 收缩了、也可能是网络带宽打满了。我之前遇到过一个真实案例:某支付对账系统凌晨 3 点 Producer 批量超时,排查了半小时发现是凌晨的日志备份任务把磁盘 IO 打满,导致 Broker 刷盘超时。同一个问题换一个时间点,根因可能变成 ISR 副本宕机。
这套排查体系不是靠背文档能建立的,得靠实际踩过的坑和系统的排查工具链来支撑。
分析问题
发消息超时(TimeoutException)的根因与排查
发消息超时是 Kafka 生产端最常见的故障,表象是 Producer 抛出 org.apache.kafka.common.errors.TimeoutException。核心原因是 Producer 在 request.timeout.ms(默认 30 秒)内没有收到 Broker 的确认。
根因一:ISR 不可达(最常见)
当 acks=all 且 min.insync.replicas=2 时,Broker 需要等待至少 2 个 ISR 副本写入成功。如果 ISR 只剩 1 个(另外 1 个宕机或网络不可达),Leader 永远无法满足条件,Producer 只能等到超时。
真实案例:某日志采集系统,Topic 副本数 3、min.insync.replicas=2。某天两台 Broker 机器同时做内核升级重启,ISR 只剩 1 个。Producer 全部超时,6 个业务 Topic 写入中断 8 分钟。事后复盘发现:acks=all 给了一个"安全"的错觉,但没人配置 min.insync.replicas 的告警。
# 排查 ISR 状态
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic order-events
# 输出示例
Topic: order-events Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2
# Isr 缺失了 3,检查 Broker 3 是否存活
# 如果 Isr 只剩 1 且 Leader 也在 1(单一副本),说明 ISR 已严重收缩
# 对比:正常状态
Topic: order-events Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3关键参数:
| 参数 | 默认值 | 推荐的极限值 | 说明 |
|---|---|---|---|
request.timeout.ms | 30000 | 60000 | 生产端超时,太短容易误判瞬时抖动 |
delivery.timeout.ms | 120000 | 120000-300000 | 包含重试的总超时,必须 > request.timeout.ms |
min.insync.replicas | 1 | 2 | Topic 级别配置,1 时 acks=all 形同虚设 |
retries | Integer.MAX_VALUE | 保持默认 | 重试次数,设为 0 会导致一次超时直接丢消息 |
踩坑:delivery.timeout.ms 必须大于 request.timeout.ms + retry.backoff.ms,否则还没重试完就超时了。我见过一个配置:request.timeout.ms=30000、delivery.timeout.ms=30000,第二次重试必然超时。
根因二:磁盘或带宽瓶颈
Broker 端磁盘 IO 打满(日志刷盘卡住)或网络带宽被打满,导致 Broker 无法及时处理 Produce 请求。
真实数据:某 Kafka 集群 3 台 Broker,每台挂 4 块 4TB HDD 做 RAID 0。白天高峰期磁盘 IO 利用率稳定在 60-70%,凌晨日志备份任务启动后飙到 95%+,IO 等待时间(await)从 5ms 涨到 200ms+。Producer 写入从 5ms 响应跳到 20s 超时。
# 排查工具链
# 1. 磁盘 IO 利用率
iostat -x 1 5
# 关键指标:%util(接近 100% 说明 IO 饱和)、await(> 30ms 说明磁盘慢)
# 2. 具体分区磁盘使用
kafka-log-dirs.sh --bootstrap-server localhost:9092 --describe --topic-list order-events
# 3. 网络带宽
# 在 Broker 上
sar -n DEV 1 5
# 看 rxkB/s 和 txkB/s,对比网卡带宽上限(比如 10Gbps 约 1250MB/s)根因三:Leader 选举或重平衡中
Partition 的 Leader 正在切换,或者 Controller 正在做重平衡,这段时间内该 Partition 不可写。
特征:Producer 日志中同时出现 LEADER_NOT_AVAILABLE 和 TimeoutException。kafka-topics.sh --describe 看到多个 Partition 的 Leader 为 -1。
排查时序:先看 Controller 是否频繁切换(kafka-controller.log 中 New controller elected 出现次数),再检查 ZooKeeper 连接是否稳定。
# 查看 Controller 状态
zookeeper-shell.sh localhost:2181 get /controller
# 输出:{"version":1,"brokerid":2,"timestamp":"..."}
# brokerid 稳定不变才是正常的,如果频繁变化说明 Controller 在颠簸消息堆积(Consumer Lag 飙升)的应对策略
消息堆积的本质是消费者的处理速度跟不上生产者的写入速度。排查的第一步是确认 Lag 有多大:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group order-processor --describe
# 输出示例
# GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
# order-processor order-events 0 1500000 2000000 500000
# order-processor order-events 1 1200000 1800000 600000
# 总 Lag = 110 万条,假设每条消息平均 1KB,约 1GB 积压场景一:分区数不足
如果分区数 ≤ 消费者数,即使增加消费者也无法提升并行度。Kafka 的消费并行度上限就是分区数。
真实案例:某订单处理系统,Topic 6 个分区,部署了 10 个消费者实例。高峰时 Lag 持续增长。加消费者到 20 个,Lag 不变。排查发现分区才 6 个,20 个消费者里有 14 个在空转。
# 查看分区数
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic order-events | wc -l
# 输出 7(1 行标题 + 6 行分区信息)
# 增加分区(注意:会导致 Key 路由变化,可能乱序)
kafka-topics.sh --alter --topic order-events --partitions 12典型流量场景对照:
| 场景 | 生产者 TPS | 单消费者处理 TPS | 分区数 | 消费组实例数 | 能否追上 |
|---|---|---|---|---|---|
| 正常业务 | 5000 | 2000 | 6 | 6 | 能(理论 12000,有富余) |
| 分区不足 | 10000 | 2000 | 6 | 12 | 不能(瓶颈在分区数) |
| 扩容后 | 10000 | 2000 | 12 | 12 | 能(理论 24000) |
场景二:消费逻辑阻塞(最常见)
消费线程池被一个慢调用占满,后续消息排队等待,Lag 飙升。典型场景:消费线程调了外部 API,该 API 超时 30 秒,线程池全部阻塞。
踩坑实况:某风控系统,消费线程里调了第三方黑名单接口,接口偶尔超时 30 秒。max.poll.records=500,max.poll.interval.ms=300000(5 分钟)。一个线程阻塞 30 秒,剩下 4 个线程很快也阻塞,5 分钟内无法完成 poll,Consumer 被踢出组,触发 Rebalance,Rebalance 期间停止消费,Lag 进一步增加。恶性循环。
时序图(文字描述):
t=0s: Consumer poll 拉取 500 条消息,分派给 5 个线程
t=1s: 4 个线程处理完,1 个线程调黑名单接口阻塞
t=31s: 阻塞线程超时返回,但剩余 4 个线程已经空闲了 30 秒
t=60s: 500 条处理完,但只有 200 条是有效处理,300 条被阻塞拖慢
t=180s: 连续 3 轮 poll 都出现阻塞,总处理时间超过 max.poll.interval.ms
t=180s+1: Consumer 被踢出组,触发 Rebalance
t=180s+5: Rebalance 完成,该 Consumer 重新分配分区,但 Lag 已经翻了 3 倍解法:消费与处理分离
Consumer 线程只负责拉取消息并放入本地内存队列,业务线程池从本地队列取消息处理。
@Component
public class DecoupledKafkaConsumer {
private final ExecutorService bizExecutor = new ThreadPoolExecutor(
8, 16, 60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1000), // 有界队列,防止内存溢出
new ThreadPoolExecutor.CallerRunsPolicy()
);
@KafkaListener(topics = "order-events", concurrency = "3")
public void onMessage(List<ConsumerRecord<String, String>> records) {
for (ConsumerRecord<String, String> record : records) {
bizExecutor.submit(() -> processRecord(record));
}
}
private void processRecord(ConsumerRecord<String, String> record) {
try {
// 这里是真正的业务逻辑,可以慢
blacklistService.check(record.value());
} catch (Exception e) {
log.error("处理失败,偏移量: {}", record.offset(), e);
// 注意:这里不能无限重试,否则线程池会被失效消息填满
// 应该投递到死信队列或者记录偏移量后跳过
}
}
}场景三:持久堆积需要扩容
如果分区数足够,但消费能力确实不够,临时扩容消费者群组是有效的。但要注意,消费者的 group.id 相同才能分担负载,新消费者加入会触发 Rebalance,期间消费暂停。
扩容操作:滚动增加消费者实例,观察 kafka-consumer-groups 的 LAG 下降趋势。如果 LAG 不降反升,说明瓶颈不在消费者数量,而在分区数或外部依赖。
消费延迟的端到端排查
消费延迟(End-to-End Latency)和消息堆积(Lag)是两回事。Lag 是积压的条数,延迟是消息从发送到被消费的时间差。可能出现 Lag 很小但延迟很大的情况——比如事务未提交导致 read_committed 的消费者需要等待。
真实案例:某金融场景,isolation.level=read_committed,Producer 端开启事务,每条消息发送后 commitTransaction() 平均耗时 500ms。消费者看到的消息是 500ms 前发出的,但 Lag 一直是 0。业务方投诉"消息延迟 500ms",监控显示 Lag=0,开发者以为没有问题。真相是 Lag 监控不管用,需要端到端时间戳埋点。
端到端埋点方案:
// Producer 端:在消息头写入时间戳
ProducerRecord<String, String> record = new ProducerRecord<>("topic", "key", "value");
byte[] tsBytes = String.valueOf(System.currentTimeMillis()).getBytes(StandardCharsets.UTF_8);
record.headers().add("produce-timestamp", tsBytes);
producer.send(record);
// Consumer 端:计算延迟,超过阈值告警
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
Header header = record.headers().lastHeader("produce-timestamp");
if (header != null) {
long produceTime = Long.parseLong(new String(header.value(), StandardCharsets.UTF_8));
long latency = System.currentTimeMillis() - produceTime;
// 监控指标:latency_p99、latency_p999
MetricsRecorder.record("kafka.end_to_end.latency", latency, "topic", record.topic());
if (latency > 5000) { // 5 秒阈值
log.warn("端到端延迟超阈值: topic={}, partition={}, offset={}, latency={}ms",
record.topic(), record.partition(), record.offset(), latency);
}
}
}触发延迟的可能原因及对比:
| 原因 | 表象 | 排查手段 | 修复方案 |
|---|---|---|---|
| Fetch 参数保守 | 延迟稳定在 500ms 左右 | 检查 fetch.max.wait.ms | 延迟敏感场景调小到 50-100ms |
| Rebalance 频繁 | 延迟间歇性跳变(秒级) | 查看 kafka-consumer-groups 的 COORDINATOR-REBALANCE | 调大 session.timeout.ms、减少成员变动 |
| 事务阻塞 | 延迟 = 事务提交时间 | 查看 Producer 事务日志 | 缩短事务范围、或者用 read_uncommitted |
| Consumer GC 停顿 | 延迟周期性跳变 | 查看 Consumer GC 日志 | 优化 GC 参数、减少 max.poll.records |
fetch.min.bytes 和 fetch.max.wait.ms 的取舍:
fetch.min.bytes=1(默认) → 频繁拉取,延迟低(1ms),吞吐低
fetch.min.bytes=65536(64KB) → 批量拉取,延迟高(取决于等待时间),吞吐高
fetch.max.wait.ms=500(默认) → 最多等 500ms 攒够数据延迟敏感场景(如实时支付通知):fetch.min.bytes=1, fetch.max.wait.ms=50。 吞吐优先场景(如日志采集):fetch.min.bytes=65536, fetch.max.wait.ms=1000。
注意:fetch.min.bytes 的默认值是 1(字节),不是 1KB。很多文章写错了,实际是 1 字节,意味着只要有数据就立即返回,不会为了凑够最小值而额外等待。
总结
| 事故类型 | 快速定位命令 | 紧急止血 | 根治措施 |
|---|---|---|---|
| 发消息超时 | kafka-topics.sh --describe 看 ISR;iostat 看磁盘 IO | 调大 request.timeout.ms + retries | 磁盘告警、min.insync.replicas=2、副本冗余 |
| 消息堆积 | kafka-consumer-groups.sh --describe 看 LAG | 增加消费者(分区够时)、扩容分区 | 消费与处理分离、线程池隔离 |
| 消费延迟 | 端到端时间戳埋点 | 调小 fetch.min.bytes=1 和 fetch.max.wait.ms=50 | 监控埋点 + 延迟告警阈值 |
生产避坑要点:
- 不要相信
acks=all就是安全的——在min.insync.replicas=1时它就是纸老虎 - 消费线程池不要用无界队列,慢调用会拖死整个消费组,触发 Rebalance 恶性循环
- 端到端延迟监控比 Lag 监控更重要,Lag 大不一定代表消息延迟高,Lag 为 0 也不代表延迟一定低
fetch.min.bytes默认值是 1(字节),不是 1KB,网上很多文章写错了- 事故复盘要落到监控告警上,而不是"下次注意"
- Rebalance 期间消费暂停,会造成延迟抖动。频繁 Rebalance 说明消费者稳定性有问题,优先排查 session.timeout 和 max.poll.interval 配置
参考:《Kafka 权威指南(第 2 版)》;Apache Kafka 官方文档;LinkedIn Burrow 文档