Learn
Kafka/08-offsets

Offset 管理

Offset 提交策略决定了你的系统在"至少一次"和"至多一次"之间站在哪边。提交早了会丢消息,提交晚了会重复消费——没有免费的午餐,但你必须清楚自己选了什么。本章讲清 offset 存哪、怎么提交、怎么重置。

1. Offset 存在哪:__consumer_offsets

消费进度存储在 Kafka 内部 topic __consumer_offsets(默认 50 分区、compact 清理策略)中:

key:   (group.id, topic, partition)
value: (offset, 元数据, 提交时间戳)
 
group=inventory, topic=orders, p0 -> offset=1500
group=inventory, topic=orders, p1 -> offset=980
group=report,    topic=orders, p0 -> offset=52000
  • 每次提交就是往这个 topic 写一条消息;compact 策略保证每个 key 只留最新值。
  • 提交的 offset 语义是"下一条要消费的位置":处理完 offset=5 应提交 6。
  • 消费组的进度按 group.id 隔离,所以不同组互不影响。

注意区分两个易混概念:

概念含义
消费位置(position)消费者内存里"下次 poll 从哪读",poll 自动推进
提交位移(committed offset)持久化到 __consumer_offsets 的进度,重启/Rebalance 后从这里恢复

两者不同步是所有重复消费问题的根源。

2. 自动提交的隐患

enable.auto.commit=true       # 默认开启
auto.commit.interval.ms=5000  # 每 5 秒在 poll 内部顺带提交

自动提交在 poll 调用时 检查距上次提交是否超过间隔,超过就把上一次 poll 返回的所有消息的 offset 提交掉。两个隐患:

隐患 1: 重复消费
  poll 拉到 100 条 -> 处理到第 60 条时进程崩溃
  -> 崩溃前没到提交时机, offset 还停在批次开头
  -> 重启后这 100 条全部重来, 前 60 条被处理两次
 
隐患 2: 消息丢失 (更隐蔽)
  poll 拉到 100 条 -> 交给另一个线程池异步处理 -> 立刻进入下一轮 poll
  -> 下一轮 poll 时自动提交了这 100 条的 offset
  -> 线程池处理失败/进程崩溃 -> 这批消息再也不会被消费 -> 丢了

结论:只要"拉取"与"处理完成"不同步,自动提交就可能丢消息。业务系统建议一律手动提交。

3. 手动提交

enable.auto.commit=false

3.1 同步与异步提交

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
    for (ConsumerRecord<String, String> r : records) {
        process(r);
    }
    // 同步: 阻塞直到成功, 失败自动重试 —— 简单可靠, 吞吐略降
    consumer.commitSync();
}
// 异步: 不阻塞, 高吞吐; 失败不重试(重试可能把新 offset 覆盖成旧的)
consumer.commitAsync((offsets, ex) -> {
    if (ex != null) log.warn("commit failed: {}", offsets, ex);
});

生产常用组合模式:循环内 commitAsync 保吞吐,关闭前 commitSync 兜底:

try {
    while (running) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
        records.forEach(this::process);
        consumer.commitAsync();
    }
} finally {
    try {
        consumer.commitSync();    // 最后一次务必同步提交
    } finally {
        consumer.close();
    }
}

3.2 按分区精细提交

一批消息里不同分区处理进度不同时,可以按分区提交,缩小重复范围:

for (TopicPartition tp : records.partitions()) {
    List<ConsumerRecord<String, String>> partRecords = records.records(tp);
    partRecords.forEach(this::process);
    long lastOffset = partRecords.get(partRecords.size() - 1).offset();
    // 提交 "已处理的最后一条 + 1"
    consumer.commitSync(Map.of(tp, new OffsetAndMetadata(lastOffset + 1)));
}

Python 对照(按消息提交):

while True:
    msg = consumer.poll(0.5)
    if msg is None or msg.error():
        continue
    process(msg)
    consumer.commit(message=msg, asynchronous=False)  # 提交 msg.offset()+1
⚠️提交的是 offset+1

提交语义是"下一条要读的位置"。处理完 offset=100 要提交 101。手动构造 OffsetAndMetadata 时忘了 +1,每次重启都会重复消费最后一条——这是隐蔽度极高的经典 bug。

4. auto.offset.reset 与 seek

4.1 auto.offset.reset

只在找不到已提交 offset时生效(新组、offset 过期被清理、数据被删):

值行为
latest(默认)从最新开始,历史消息不读
earliest从最早开始,读全部历史
none抛异常,交给应用决策

新上线的消费组用默认 latest 会跳过所有历史数据;要补历史必须 earliest。

4.2 seek:程序内重置

// 必须在分区分配完成后调用
consumer.subscribe(List.of("orders"), new ConsumerRebalanceListener() {
    public void onPartitionsAssigned(Collection<TopicPartition> parts) {
        // 回到 1 小时前
        long ts = System.currentTimeMillis() - 3600_000;
        Map<TopicPartition, Long> query = new HashMap<>();
        parts.forEach(tp -> query.put(tp, ts));
        consumer.offsetsForTimes(query).forEach((tp, ot) -> {
            if (ot != null) consumer.seek(tp, ot.offset());
        });
    }
    public void onPartitionsRevoked(Collection<TopicPartition> parts) {}
});

运维场景更常用 CLI 重置(需先停消费者,见第 4 章 --reset-offsets)。

5. 重复与丢失的成因矩阵

提交时机语义崩溃后果
处理前提交(或自动提交+异步处理)至多一次 at-most-once丢消息:已提交未处理
处理后提交(手动 commitSync)至少一次 at-least-once重复消费:已处理未提交
处理与提交原子化(事务/外部存储)精确一次 exactly-once不丢不重,成本最高

工程实践中的主流选择是至少一次 + 下游幂等:

  • 数据库写入带唯一键(如 order_id),重复消息触发唯一约束直接跳过;
  • 或用 INSERT ... ON DUPLICATE KEY UPDATE / upsert;
  • Redis 用 SETNX 记录已处理的消息 ID。
💡offset 过期清理

__consumer_offsets 中的记录在消费组空置超过 offsets.retention.minutes(默认 7 天)后会被清理。一个下线两周的服务重新上线时,已提交进度已不存在,会按 auto.offset.reset 走——latest 则跳过两周数据,earliest 则全量重刷。长期停用的组要心里有数。

小结

  • Offset 存在内部 topic __consumer_offsets,按 (group, topic, partition) 记录"下一条位置"
  • 自动提交在拉取与处理解耦时会丢消息,业务系统用手动提交
  • 生产模式:循环 commitAsync + 关闭前 commitSync;精细化用按分区提交
  • auto.offset.reset 只在无已提交 offset 时生效;程序内用 seek,运维用 CLI reset
  • 成因矩阵:先提交=可能丢,后提交=可能重;主流方案是至少一次+下游幂等
🎯练习
  1. 开自动提交,poll 后把消息丢进线程池并让线程池抛异常,验证消息丢失;改成处理完手动提交,验证变为重复消费。
  2. 写一个按分区提交的消费者,在处理到第二个分区时故意崩溃,重启后验证只有第二个分区的消息被重放。
  3. 用 offsetsForTimes + seek 实现"从 10 分钟前开始重新消费",并解释它和 CLI --to-datetime 重置的适用场景差别。