找回密码
立即注册
搜索
热搜: Java Python Linux Go
发回帖 发新帖

4169

积分

0

好友

539

主题
发表于 5 小时前 | 查看: 4| 回复: 0

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

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

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

先理解自动提交的时机

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

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

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

注意这个时序:

Kafka自动提交机制:poll间隔与offset提交时机

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 调用或者数据库大批量写入):

Kafka消费者处理超时触发Rebalance导致重复消费

这种场景下:

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

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

你想,一次 Full GC 的 STW 停顿可能就几十秒甚至几分钟,这段时间所有线程全卡住了,消费者线程也跑不了。等 GC 结束缓过来,两次 poll() 的间隔早就超过 max.poll.interval.ms 了,Rebalance 直接触发。

最恶心的是日志里看不出啥毛病,消费者还在正常处理消息,就是提交 offset 的时候突然蹦出来一个 CommitFailedException,不仔细翻日志根本注意不到。

手动提交也不是万无一失

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

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

手动提交无非是把提交的主动权拿到自己手里,想什么时候提交就什么时候提交。但核心问题没变——消息处理完到 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);
    }
}

说在最后

回过头来捋一下:正常重启,自动提交有 5 秒间隔,最后一批 offset 来不及提交;消费太慢触发 Rebalance,已经处理完的消息 offset 还没提交就被别的消费者抢走了;GC 一来 STW 停顿让 poll 间隔超时,效果跟消费慢一个德行;就算换成手动提交,处理和提交之间总有个窗口期,窗口期内崩了 offset 照样丢。

归根结底就一件事:消息处理完了不等于 offset 提交完了,中间永远有一个空档。

Kafka 的 offset 机制天生就只保证至少一次,不保证精确一次。

所以写 Kafka 消费端的时候记住:别指望靠 offset 来保证不重复,让你的业务逻辑本身能扛住重复消费才是正道。




上一篇:TCPdump 从入门到排障实战:替代 Wireshark 的服务器网络抓包方案
下一篇:大模型的下半场:拼的是记忆,RAG之后AI原生记忆如何落地?
您需要登录后才可以回帖 登录 | 立即注册

手机版|小黑屋|网站地图|云栈社区 ( 苏ICP备2022046150号-2 )

GMT+8, 2026-8-4 07:22 , Processed in 0.780202 second(s), 41 queries , Gzip On.

Powered by Discuz! X3.5

© 2025-2026 云栈社区.

快速回复 返回顶部 返回列表