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=false3.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=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。
__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
- 成因矩阵:先提交=可能丢,后提交=可能重;主流方案是至少一次+下游幂等
- 开自动提交,poll 后把消息丢进线程池并让线程池抛异常,验证消息丢失;改成处理完手动提交,验证变为重复消费。
- 写一个按分区提交的消费者,在处理到第二个分区时故意崩溃,重启后验证只有第二个分区的消息被重放。
- 用 offsetsForTimes + seek 实现"从 10 分钟前开始重新消费",并解释它和 CLI --to-datetime 重置的适用场景差别。