跳至主要內容

程序员小富大约 5 分钟

大家好,我是小富。

《十万个why》系列持续更新中

线上一个消费者服务发版重启,起来之后开始疯狂消费消息,监控告警一片红。拉出来一看,好几万条消息被重复消费了。

排查 offset 提交记录,明确显示重启前 offset 已经成功提交到了最新位置。既然提交了,为什么重启后又从很早的位置开始消费?

这个问题坑就坑在:offset 确实提交成功了,但提交的不是你以为的那个 offset。

先理解自动提交的时机

大部分项目用的是 Kafka 默认的自动提交模式:

enable.auto.commit=true
auto.commit.interval.ms=5000  # 默认 5 秒提交一次

自动提交不是"每消费一条就提交一次",而是每隔 5 秒,在下一次 poll() 调用时,顺便把上一次 poll 到的最大 offset 提交掉

注意这个时序:

T0: poll() 拉到 offset 100 ~ 200 的消息
T1: 开始处理这 100 条消息...
T5: 5 秒到了,但还没处理完,也没有再次 poll()
T8: 处理完了,调用 poll() 拉下一批
    → 此时自动提交触发,提交 offset = 200
T8: poll() 拉到 offset 200 ~ 350 的消息
T9: 开始处理...
T10: 服务重启了!offset 350 还没提交
     → 重启后从 offset 200 开始消费,200~350 全部重复

看到了吗?offset 200 确实提交成功了,但你已经处理到 350 了,200 ~ 350 之间的 offset 还没来得及提交。

Rebalance 才是真正的杀手

上面是最理想的场景——干净地重启。更常见的情况是 Rebalance 导致的重复消费。

消费者被 Kafka 判定为"死亡"并触发 Rebalance 的三种情况:

  1. 心跳超时session.timeout.ms(默认 45 秒)内没收到心跳
  2. poll 间隔过长:两次 poll() 的间隔超过 max.poll.interval.ms(默认 5 分钟)
  3. 消费者主动离组:调用 close() 或服务关闭

重点说第 2 种,因为它是最隐蔽、最常见的线上事故来源。

假设你的消费者拉了一批消息,然后业务处理特别慢(比如里面有 HTTP 调用或者数据库大批量写入):

T0: poll() 拉到 100 条消息,offset 1000 ~ 1100
T0: 自动提交了上一轮的 offset 1000
T1 ~ T300: 慢慢处理这 100 条消息,每条涉及 RPC 和数据库操作...

T300 (5分钟): max.poll.interval.ms 超时了!
    → Kafka Coordinator 判定这个消费者死了
    → 触发 Rebalance,把 partition 分配给其他消费者
    → 其他消费者从 offset 1000 开始消费(因为上次提交的就是 1000)
    → 1000 ~ 1100 的消息全部被重复消费

T310: 原消费者处理完了,试图提交 offset 1100
    → 提交失败:CommitFailedException,因为它已经不再是这个 partition 的 owner 了

这种场景下:

  • offset 1000 确实提交成功了
  • 但 1000 ~ 1100 之间的消息你已经处理完了,只是 offset 还没来得及提交
  • Rebalance 后新消费者从 1000 开始,导致重复

更隐蔽的坑:GC 导致 poll 间隔超时

你可能觉得 5 分钟够长了,业务处理不至于那么慢。但你忘了一个东西:JVM GC。

一次 Full GC 的 STW 停顿可能就几十秒甚至几分钟(大堆内存、CMS/Parallel GC 的情况下)。GC 期间所有线程都暂停了,包括消费者线程。GC 结束后,poll() 间隔已经超过了 max.poll.interval.ms,Rebalance 被触发。

你在日志里看到的是:消费者"正常"地处理完了消息,然后提交 offset 失败。看不出任何业务异常,只有一行不起眼的 CommitFailedException

手动提交也不是万无一失

很多人说"关掉自动提交,改手动提交就没事了"。手动提交确实好一些,但同样有坑:

// 手动同步提交
consumer.commitSync();

手动提交只是把提交的时机从"下次 poll 时"变成了"你想提交就提交"。但核心问题没变:消息处理完毕到 offset 提交成功之间,如果发生 Rebalance、网络超时、进程崩溃,offset 就丢了。

唯一能做到"精确一次"的方式是事务性提交或者消费端幂等

实战解决方案

方案一:缩小批次 + 及时提交

# 每次 poll 少拉一点,处理快一些
max.poll.records=50
# poll 间隔给长一些
max.poll.interval.ms=600000
# 关闭自动提交
enable.auto.commit=false
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
    for (ConsumerRecord<String, String> record : records) {
        processMessage(record);
    }
    // 处理完这批就立即提交
    consumer.commitSync();
}

方案二:逐条提交(性能稍差但最精确)

for (ConsumerRecord<String, String> record : records) {
    processMessage(record);
    // 每处理一条就提交这条的 offset
    consumer.commitSync(Collections.singletonMap(
        new TopicPartition(record.topic(), record.partition()),
        new OffsetAndMetadata(record.offset() + 1)
    ));
}

方案三:消费端幂等兜底(终极方案)

不管 offset 提交得多精确,网络分区、进程崩溃等极端场景下还是可能重复。所以关键业务必须在消费端做幂等:

public void processMessage(ConsumerRecord<String, String> record) {
    String bizId = extractBizId(record);
    // 用唯一业务 ID 判断是否已处理
    if (redis.setIfAbsent("consumed:" + bizId, "1", 24, TimeUnit.HOURS)) {
        // 第一次消费,执行业务逻辑
        doBusinessLogic(record);
    } else {
        // 重复消费,跳过
        log.warn("重复消费,跳过: {}", bizId);
    }
}

总结

场景根因offset 状态
正常重启自动提交有 5 秒间隔,最后一批 offset 没来得及提交上一批已提交,当前批未提交
消费慢触发 Rebalance处理时间超过 max.poll.interval.ms已处理的消息 offset 未提交
GC 导致 poll 超时STW 停顿期间无法 poll,被判定死亡同上
手动提交也重复提交和处理之间有窗口,崩溃时 offset 丢失同上

Kafka 的 offset 提交机制保证的是"至少一次消费",不是"精确一次消费"。 在提交 offset 和处理消息之间永远存在一个时间窗口,这个窗口内发生任何异常都会导致重复消费。

所以 Kafka 消费端的黄金法则:不要依赖 offset 精度来保证不重复,而是让业务逻辑本身具备幂等能力。


我是小富,下期见。

上次编辑于: