Skip to content

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=allmin.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 的告警。

bash
# 排查 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.ms3000060000生产端超时,太短容易误判瞬时抖动
delivery.timeout.ms120000120000-300000包含重试的总超时,必须 > request.timeout.ms
min.insync.replicas12Topic 级别配置,1 时 acks=all 形同虚设
retriesInteger.MAX_VALUE保持默认重试次数,设为 0 会导致一次超时直接丢消息

踩坑delivery.timeout.ms 必须大于 request.timeout.ms + retry.backoff.ms,否则还没重试完就超时了。我见过一个配置:request.timeout.ms=30000delivery.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_AVAILABLETimeoutExceptionkafka-topics.sh --describe 看到多个 Partition 的 Leader 为 -1。

排查时序:先看 Controller 是否频繁切换(kafka-controller.logNew controller elected 出现次数),再检查 ZooKeeper 连接是否稳定。

bash
# 查看 Controller 状态
zookeeper-shell.sh localhost:2181 get /controller
# 输出:{"version":1,"brokerid":2,"timestamp":"..."}
# brokerid 稳定不变才是正常的,如果频繁变化说明 Controller 在颠簸

消息堆积(Consumer Lag 飙升)的应对策略

消息堆积的本质是消费者的处理速度跟不上生产者的写入速度。排查的第一步是确认 Lag 有多大:

bash
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 个在空转。

bash
# 查看分区数
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分区数消费组实例数能否追上
正常业务5000200066能(理论 12000,有富余)
分区不足100002000612不能(瓶颈在分区数)
扩容后1000020001212能(理论 24000)

场景二:消费逻辑阻塞(最常见)

消费线程池被一个慢调用占满,后续消息排队等待,Lag 飙升。典型场景:消费线程调了外部 API,该 API 超时 30 秒,线程池全部阻塞。

踩坑实况:某风控系统,消费线程里调了第三方黑名单接口,接口偶尔超时 30 秒。max.poll.records=500max.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 线程只负责拉取消息并放入本地内存队列,业务线程池从本地队列取消息处理。

java
@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 监控不管用,需要端到端时间戳埋点。

端到端埋点方案

java
// 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-groupsCOORDINATOR-REBALANCE调大 session.timeout.ms、减少成员变动
事务阻塞延迟 = 事务提交时间查看 Producer 事务日志缩短事务范围、或者用 read_uncommitted
Consumer GC 停顿延迟周期性跳变查看 Consumer GC 日志优化 GC 参数、减少 max.poll.records

fetch.min.bytesfetch.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=1fetch.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 文档

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。