主题
第 12 章 Offset 与「消息语义」:At Most Once / At Least Once / Exactly Once
目标读者:知道
consumer.poll()能拿消息、commit()能提交位点,但说不清「自动提交到底丢消息还是重消息」、不会用seek()回放、不知道__consumer_offsets长什么样的同学。学完你会:能用一张图讲清三种消息语义在工程上的真实含义;知道自动提交的「先消费后提交 / 先提交后消费」陷阱具体怎么触发;会写业务幂等消费者;会用
seek/offsetsForTimes做时间回溯;能直接消费__consumer_offsets看到底层是怎么存的。
0. 导读:消息语义不是「Kafka 一个开关」,而是「端到端」的协作
很多同学看到「Kafka 支持 Exactly Once」,第一反应是「好棒!打开它就行了!」——结果上了生产,发现订单仍然偶尔重复扣款,对运维大喊:「Kafka 不是 EOS 吗?!」
这个误解的根源在于:消息语义是「Producer + Broker + Consumer + 业务系统」四方协作的端到端结果,不是 Kafka 自己一个开关能定的。
举个生活例子。你妈让你去超市买酱油(一条消息),三种语义对应:
- At Most Once(最多一次):你出门时口袋装着钱,但没记是不是真的买了,回家空手就空手,妈咪不会让你再跑一趟。结果:要么买到(一次),要么没买(零次),不会重复花钱。对应「先提交 offset,再处理消息」——只要 offset 提交了,消息就当作处理过了,哪怕你处理到一半挂了也不补。
- At Least Once(至少一次):你买了酱油回来,但忘了告诉妈,她以为你没去,再让你去一次。结果:可能买一瓶,可能买两瓶,但绝不会一瓶都没有。对应「先处理消息,再提交 offset」——处理成功后才提交,处理完没来得及提交就挂了,下次会重新处理。
- Exactly Once(恰好一次):你买了酱油,并且妈妈那里有「今天已经买过」的记录,无论你跑几趟,钱只扣一次。对应「业务侧幂等 + offset 提交原子化」——要么靠业务主键去重,要么靠 Kafka 事务把「写下游 + 提交 offset」绑在一起。
所以本章先把「消息语义」的工程含义讲透,然后聚焦在「Consumer 端如何选择和正确实现」。Producer 端的幂等与事务(真·EOS 协议)放在第 13 章细讲。
1. 三种语义的精确定义
1.1 一张图看明白
┌──────────────────────────────────────────────────────────────────┐
│ At Most Once (最多一次) │
│ 消息流 : ──Msg──▶ Consumer.poll │
│ 时间轴 : │ commit offset │ process(Msg) │
│ ↑ 先提交 ↑ 后处理 │
│ 崩溃点 : 在「提交后/处理前」崩溃 → 这条 Msg 永远丢了 │
│ 语义 : 0 或 1 次(可能丢,绝不重) │
│ 适合 : 实时监控指标采样、热度统计、可容忍少量丢失的日志 │
├──────────────────────────────────────────────────────────────────┤
│ At Least Once (至少一次) │
│ 消息流 : ──Msg──▶ Consumer.poll │
│ 时间轴 : │ process(Msg) │ commit offset │
│ ↑ 先处理 ↑ 后提交 │
│ 崩溃点 : 在「处理后/提交前」崩溃 → 重启后重新拉取,再处理一次 │
│ 语义 : ≥1 次(不会丢,可能重) │
│ 适合 : 99% 的业务场景(配合业务幂等就行了) │
├──────────────────────────────────────────────────────────────────┤
│ Exactly Once (恰好一次) │
│ 方案 A : At Least Once + 业务侧幂等(去重表 / 唯一约束) │
│ 方案 B : Kafka 事务 EOS(consume-process-produce 端到端原子) │
│ 语义 : 业务上观察到的效果 = 1 次 │
│ 适合 : 计费扣款、库存扣减、积分变更 │
└──────────────────────────────────────────────────────────────────┘1.2 银行转账类比
把一条消息当作「给客户 A 转账 100 元」:
- At Most Once:转账请求发给柜台,柜台员工先在小本本上记「已处理」(提交 offset),再去操作系统转账。万一在记完小本本、还没点确认的时候停电了——客户 A 没收到钱,但银行系统认为「这条请求已处理」,下次重启不会再转。钱丢了。
- At Least Once:柜台员工先操作系统转账,转完成功了再在小本本上记「已处理」。万一转完账、小本本还没记好就停电了——重启后柜台看到这条请求还没记,再转一次。客户 A 收到 200 元。
- Exactly Once:柜台在本地有一张「转账流水」表,每条请求带一个唯一流水号。每次转账前先查流水:流水号已存在 → 直接告诉调用方「成功」,不再操作;流水号不存在 → 转账 + 写流水表(同一个事务)。这样无论柜台失败重试多少次,客户 A 都只会收到 100 元。这就是「业务幂等」实现 EOS。
📌 生产经验:90% 的场景用「At Least Once + 业务幂等」就够了,性能高、实现简单。Kafka 的 EOS 协议(PID + 事务)更适合「Kafka → 计算 → Kafka」的流处理链路(Kafka Streams / Flink),让消费 + 处理 + 写出三步原子化。本章主讲前者,后者放第 13 章。
2. 自动提交 enable.auto.commit=true 的两种陷阱
2.1 自动提交在做什么?
confluent-kafka / Java kafka-clients 默认开启自动提交:
properties
enable.auto.commit=true
auto.commit.interval.ms=5000意思是:消费者在每次 poll() 时,如果距离上次提交超过 auto.commit.interval.ms(默认 5s),就把当前已经 poll 出来的最大 offset 提交一次。注意几个关键点:
- 提交时机是下一次 poll,不是某个后台线程定时主动提交(很多人误以为有后台线程)。
- 提交的 offset 是「已经 poll 出来交给业务的最大 offset + 1」,而不是「业务真的处理完的 offset」——Kafka 客户端根本不知道你业务有没有处理完!
- 因此自动提交的语义本质是 At Most Once(处理失败也算处理过了),但因为 5s 周期里你可能 poll 了多次,又会出现 At Least Once(处理成功了但还没到 5s 就重启)。两种悲剧并存。
2.2 陷阱一:丢消息(处理还没完,offset 已经提交)
时间轴:
t=0.0s poll() → 拿到 Msg[100..200]
t=0.1s 开始处理 Msg[100],写 DB / 调下游 ……(业务很慢,要 8s)
t=5.0s poll() ← 这里触发了自动提交,把 200+1 提交了!
「客户端觉得 100..200 都消费了」
t=6.0s 消费者进程 OOM 崩溃,被 K8s 重启
t=10s 重启后从 offset=201 开始消费 → Msg[100..200] 永远丢了真实场景:业务里调了一个慢的第三方接口(比如发短信),单条处理 10 秒,而 auto.commit.interval.ms=5000。一旦在第二次 poll() 之前业务进程崩了,那批未处理完的消息就这样消失在了茫茫数据流里。
排查现场你会看到:
- 业务日志最后一行是「正在调用短信接口...」
- Kafka
__consumer_offsets里这个组的 offset 已经推进到 201 - 数据库里没有 Msg[150] 的处理记录
- 老板:「为什么用户没收到验证码?!」
2.3 陷阱二:重消息(处理完了,但还没到 commit 周期)
时间轴:
t=0.0s poll() → 拿到 Msg[100..200]
t=0.1s 快速处理完 Msg[100..200],全部写 DB 成功
t=2.0s 消费者崩溃(还没到 5s commit 周期)
t=5.0s 重启后从上次提交的 offset=100 开始消费
→ Msg[100..200] 又处理了一遍 → 数据库里 100~200 行重复如果业务没有幂等保护,那 200 条订单变 400 条订单、200 次扣款变 400 次扣款,运营那边电话已经打爆。
2.4 自动提交的「时间窗口」可视化
process 时间 = 0.5s, commit interval = 5s
─────────────────────────────────────────────────
t=0.0 poll M1..M100
t=0.0 commit (距上次>5s, 提交 offset=100)
t=0.5 done M1..M100
t=0.5 poll M101..M200
t=1.0 done M101..M200
t=1.0 poll M201..M300
... (4 轮 poll 都在 commit 周期内)
t=5.0 poll M501..M600
t=5.0 commit 提交 offset=600 ← 距上次提交 5.0s
t=5.5 done M501..M600
─────────────────────────────────────────────────
风险窗口:t=2.0 崩溃 → offset 还停在 100 → M101..M400 全部重复| 场景 | process 时长 | commit interval | 风险 |
|---|---|---|---|
| 处理快、间隔短 | 50ms | 1s | 重复窗口很小 |
| 处理慢、间隔长 | 8s | 5s | 「丢」非常严重 |
| 处理很快、间隔默认 | 50ms | 5s | 「重」很严重 |
📌 结论:自动提交是「为了懒人方便」的开关。任何一个有钱、有人、有责任心的业务都应该 关闭自动提交,改用手动提交。
3. 手动提交的正确姿势
关闭自动提交:
properties
enable.auto.commit=false接下来你要决定提交什么、什么时候提交、用同步还是异步。
3.1 三种 API 对照
confluent-kafka Python 客户端提供:
| API | 含义 | 用法 |
|---|---|---|
consumer.commit() | 同步提交,阻塞直到 Broker 确认 | 简单可靠,但有 RTT 延迟 |
consumer.commit(asynchronous=True) | 异步提交,立即返回 | 吞吐高,失败不会重试 |
consumer.commit(message=msg) | 提交单条消息的 offset+1 | 精细化控制 |
consumer.commit(offsets=[TopicPartition(...)]) | 提交一组 (topic, partition, offset) | 批量、跨分区精细化 |
Java KafkaConsumer 几乎一一对应:commitSync() / commitAsync() / commitSync(Map<TopicPartition, OffsetAndMetadata>)。
3.2 「先处理后提交」模板(At Least Once)
python
from confluent_kafka import Consumer
c = Consumer({
"bootstrap.servers": "127.0.0.1:9092",
"group.id": "order-processor-v1",
"enable.auto.commit": False,
"auto.offset.reset": "earliest",
})
c.subscribe(["learn.12.orders"])
try:
while True:
msg = c.poll(timeout=1.0)
if msg is None:
continue
if msg.error():
print("err:", msg.error()); continue
try:
process(msg) # ← 业务处理(要幂等!)
c.commit(message=msg, asynchronous=False) # 处理完才提交
except Exception as e:
# 处理失败:不提交,下次 poll 还会拿到(重试)
log.exception("process failed, will retry")
finally:
c.close()关键点:
process(msg)必须业务幂等——At Least Once 一定会重,靠业务保证最终一次。commit(asynchronous=False)同步提交,确保「真的提交进了 Broker」再继续下一轮。- 异常时不要捕获后跳过!否则等价于自动提交丢消息。让它停在那条消息上反复重试(或者 N 次后投到死信队列)。
3.3 「批量提交」性能优化
每条消息都同步提交一次太慢(每条都有 RTT)。生产实践是:业务处理一批,再统一提交一次最大 offset。
python
buf = []
for _ in range(N_INNER):
msg = c.poll(1.0)
if msg and not msg.error():
buf.append(msg)
if len(buf) >= 100:
break
if buf:
process_batch(buf) # 批量处理(事务/批量插入)
last = buf[-1] # 取最后一条
c.commit(message=last, asynchronous=False)⚠️ 注意:批量处理失败时要么整批回滚、要么记录失败的精确 offset,不能简单 commit(last) 然后丢失中间几条。
3.4 异步 + 同步混合模式(最佳实践)
python
try:
while True:
msg = c.poll(1.0)
if msg is None: continue
process(msg)
c.commit(message=msg, asynchronous=True) # 平时异步(高吞吐)
except KeyboardInterrupt:
pass
finally:
try:
c.commit(asynchronous=False) # 关闭前同步一次(确保最终位点落盘)
finally:
c.close()平时用异步提交(不阻塞处理),优雅关闭时用同步提交「兜底」,是 Kafka 文档推荐的模板。
3.5 精细化:只提交某个分区的某个 offset
业务有时希望「这一批消息中,A 分区前 50 条处理成功,B 分区只成功了 30 条」,分开提交:
python
from confluent_kafka import TopicPartition
c.commit(offsets=[
TopicPartition("learn.12.orders", 0, 51), # 注意:要提交「下次开始消费的 offset」= 已处理 + 1
TopicPartition("learn.12.orders", 1, 31),
], asynchronous=False)🔥 常见 BUG:很多人误以为提交的是「最后处理的 offset」,结果漏提交一条或重复提交一条。Kafka 提交的是 next offset to consume,即「已处理最大 offset + 1」。
4. 业务幂等设计模式
At Least Once 一定会重复,业务幂等是必备技能。下面是四种主流模式。
4.1 模式一:业务主键 + DB 唯一约束(推荐)
最简单也最可靠:消息携带业务唯一 ID(订单号、流水号),DB 表给该字段加 UNIQUE 索引,重复插入直接报错忽略即可。
sql
CREATE TABLE order_event (
order_id VARCHAR(64) PRIMARY KEY, -- 业务主键
user_id BIGINT,
amount DECIMAL(18,2),
status TINYINT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);python
def process(msg):
payload = json.loads(msg.value())
try:
cursor.execute(
"INSERT INTO order_event(order_id, user_id, amount, status) VALUES (%s,%s,%s,%s)",
(payload["order_id"], payload["user_id"], payload["amount"], 1),
)
conn.commit()
except IntegrityError:
# 已经处理过,直接忽略
log.info(f"dup: {payload['order_id']}")优点:靠 DB 原子性,无须额外组件。缺点:必须有「天然业务主键」,且对 DB 写入吞吐有压力。
4.2 模式二:Redis SETNX 去重
适合「业务主键不便落 DB」或「DB 太慢」的场景。
python
def process(msg):
payload = json.loads(msg.value())
key = f"dedup:order:{payload['order_id']}"
# SET key value NX EX 86400 → 只在 key 不存在时设置,TTL 1 天
if not redis.set(key, 1, nx=True, ex=86400):
log.info("dup, skip"); return
# 真正的业务处理
do_business(payload)优点:吞吐高,TTL 自动清理。缺点:① TTL 后再来重复消息会再次处理 ② Redis 与下游 DB 不在一个事务,理论上可能「Redis 标记成功但 DB 失败」(建议把 Redis SET 放在业务最后做「成功标记」)。
4.3 模式三:去重表 + TTL(事件溯源风格)
去重表与业务表分离,可记录处理时间、来源等审计信息,便于排查。
sql
CREATE TABLE consumer_dedup (
msg_id VARCHAR(128) PRIMARY KEY,
consumer VARCHAR(64),
processed_at TIMESTAMP,
KEY idx_processed_at (processed_at)
);
-- 定期任务清理 7 天前的记录python
def process(msg):
msg_id = compute_msg_id(msg) # 优先用业务 ID,没有就用 (topic,partition,offset)
try:
with conn.begin():
cursor.execute(
"INSERT INTO consumer_dedup(msg_id, consumer, processed_at) VALUES (%s,%s,NOW())",
(msg_id, "order-processor-v1"),
)
do_business(json.loads(msg.value()))
except IntegrityError:
log.info(f"dup: {msg_id}")优点:审计友好。缺点:去重表会膨胀,要 TTL 清理。
4.4 模式四:状态机限制(业务领域天然幂等)
很多业务流程有状态机,比如订单:待支付 → 已支付 → 已发货 → 已完成。「已支付」消息只能让订单从「待支付」变「已支付」,重复执行不会让状态再前进。
python
def process_paid(msg):
payload = json.loads(msg.value())
rows = cursor.execute(
"UPDATE orders SET status='PAID' WHERE order_id=%s AND status='PENDING'",
(payload["order_id"],),
)
if rows == 0:
log.info("status mismatch or already paid, skip")UPDATE ... WHERE status='PENDING' 自带幂等:第二次执行时 status 已是 PAID,匹配 0 行,无副作用。
📌 选型建议:
- 有「天然业务主键」→ 模式一(最稳)
- 高吞吐、可容忍极小概率漏判 → 模式二
- 需要审计 → 模式三
- 业务有状态机 → 模式四(甚至不用额外去重组件)
实际项目里常常模式一 + 模式四组合,双重保险。
5. seek 系列:从任意位置开始消费
Kafka 的杀手级特性之一就是「消息可以反复回放」。Offset 是开放的,你可以把消费位点拨到任何地方。
5.1 四类 seek 方法
| 方法 | 含义 |
|---|---|
seek(TopicPartition(t, p, offset)) | 跳到指定 (topic, partition, offset) |
seek_to_beginning(partitions) | 跳到分区最早(earliest) |
seek_to_end(partitions) | 跳到分区最新(latest) |
offsets_for_times({tp: timestamp_ms}) | 按时间戳找到 offset,再 seek 过去 |
⚠️ 重要约束:调用 seek() 之前,必须先 poll() 一次(或者用 assign() 手动指派分区),让 Coordinator 完成 Rebalance、把分区真的分给当前消费者,否则 seek 报「No current assignment for partition」。
5.2 跳到任意 offset
python
from confluent_kafka import Consumer, TopicPartition
c = Consumer({
"bootstrap.servers": "127.0.0.1:9092",
"group.id": "replay-debug",
"enable.auto.commit": False,
})
c.subscribe(["learn.12.orders"])
c.poll(2.0) # 先触发 Rebalance
c.seek(TopicPartition("learn.12.orders", 0, 1234)) # 0 号分区从 offset=1234 重新消费
while True:
msg = c.poll(1.0)
if msg and not msg.error():
print(msg.offset(), msg.value()[:80])5.3 按时间戳定位(线上排障神器)
业务说「线上昨天 14:00~14:05 有一批订单异常」,你只需要:
python
import time
from confluent_kafka import Consumer, TopicPartition
c = Consumer({
"bootstrap.servers": "127.0.0.1:9092",
"group.id": "debug-by-time",
"enable.auto.commit": False,
})
c.subscribe(["learn.12.orders"])
c.poll(2.0)
# 找出所有分区
md = c.list_topics("learn.12.orders").topics["learn.12.orders"]
ts = int(time.mktime(time.strptime("2026-04-16 14:00:00", "%Y-%m-%d %H:%M:%S")) * 1000)
tps = [TopicPartition("learn.12.orders", p, ts) for p in md.partitions]
offsets = c.offsets_for_times(tps, timeout=5.0) # → 每个分区返回 (tp, offset_at_or_after_ts)
for tp in offsets:
print(f"P{tp.partition}: seek to offset {tp.offset}")
c.seek(tp)
end_ts = ts + 5 * 60 * 1000 # 14:05
while True:
msg = c.poll(1.0)
if msg is None: continue
if msg.timestamp()[1] > end_ts: break
print(msg.offset(), msg.value()[:120])底层是 Broker 利用 .timeindex 文件(每条 (timestamp, offset) 索引)做二分查找,详见第 7 章存储结构。
5.4 auto.offset.reset 与 seek 的关系
auto.offset.reset只在「当前消费组在该分区没有任何已提交 offset」时才生效。- 一旦提交过、或者你
seek过,下次启动会从已提交位点继续,不再走auto.offset.reset。
auto.offset.reset 值 | 行为 |
|---|---|
latest(默认) | 从最新位点开始(之前的消息全跳过) |
earliest | 从最早位点开始(重放全部) |
none | 直接抛异常,不允许「无 offset」消费 |
踩坑:很多人第一次跑消费者,发现「之前生产的消息没收到」,原因就是默认 latest。生产环境推荐显式设置为 earliest 或 none,避免「上线 bug 导致丢历史」。
6. __consumer_offsets 内部 Topic 的真相
Kafka 0.9 之前,offset 存在 ZK 里(一秒几十次写 ZK 把 ZK 写死)。0.9 之后引入了 __consumer_offsets 这个特殊 Topic,offset 全部靠 Kafka 自己存。
6.1 基本属性
- 名字:
__consumer_offsets(双下划线开头,是「内部 Topic」,普通 list-topics 默认看不到) - 默认 50 个分区(
offsets.topic.num.partitions) - 默认副本因子
offsets.topic.replication.factor=3(生产至少 3) cleanup.policy=compact:只保留每个 Key 的最新 Value,详见第 14 章
每个消费组的 offset 落到哪个分区?哈希算:
partition = hash(group_id) % num_partitions也就是说,同一个 group 的所有 offset 都集中在 __consumer_offsets 的同一个分区里。一个 group 的 Coordinator(GroupCoordinator)就是这个分区的 Leader 所在 Broker。
6.2 Key / Value 的二进制结构
__consumer_offsets 里每条消息长这样:
Key:
┌─────────────────────────────────────────────────┐
│ short version │
│ string group_id │
│ string topic │
│ int32 partition │
└─────────────────────────────────────────────────┘
→ 用 GroupMetadataManager 的 OffsetKey schema 序列化
Value:
┌─────────────────────────────────────────────────┐
│ short version │
│ int64 offset │
│ int32 leader_epoch (v3+) │
│ string metadata (用户自定义字符串) │
│ int64 commit_timestamp │
│ int64 expire_timestamp (旧版才有) │
└─────────────────────────────────────────────────┘
→ OffsetValue schema每次 commit() 都会向 __consumer_offsets 追加一条新消息(不是「原地更新」)。但因为 cleanup.policy=compact,旧的 (group, topic, partition) Key 的旧 Value 会被压缩掉,最终只剩最新一条。
📌 这个设计的妙处:offset 写入也是「追加日志 + Compaction」——Kafka 自己吃自己的狗粮。整个集群没有任何「随机写」,全都是顺序追加。
6.3 还存了什么?不只是 offset
实际上 __consumer_offsets 里还存了消费组元数据:
- OffsetCommit 消息:上面的 (group, topic, partition) → offset
- GroupMetadata 消息:(group) → 当前 generation / member 列表 / 分配方案 / protocol 等
后者是 Coordinator 用来做 Rebalance 的「快照」,挂掉后新 Coordinator 通过重放 __consumer_offsets 恢复状态。
6.4 命令行:查看一个 group 的 offset
bash
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \
--describe --group order-processor-v1输出示例:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
order-processor-v1 learn.12.orders 0 12345 12350 5 consumer-1-3a7c-... /10.0.0.7 consumer-1
order-processor-v1 learn.12.orders 1 11200 11200 0 consumer-1-3a7c-... /10.0.0.7 consumer-1
order-processor-v1 learn.12.orders 2 9800 9800 0 - - -CURRENT-OFFSET来自__consumer_offsetsLOG-END-OFFSET来自分区末尾LAG= LEO - CURRENT,线上最重要的报警指标
6.5 直接消费 __consumer_offsets(高级)
把它当普通 Topic 消费,需要用 Kafka 的 OffsetsMessageFormatter:
bash
kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \
--topic __consumer_offsets \
--formatter "kafka.coordinator.group.GroupMetadataManager\$OffsetsMessageFormatter" \
--from-beginning输出:
[order-processor-v1,learn.12.orders,0]::OffsetAndMetadata(offset=12345, leaderEpoch=Optional[7], metadata=, commitTimestamp=1745134567890, expireTimestamp=None)
[order-processor-v1,learn.12.orders,1]::OffsetAndMetadata(offset=11200, leaderEpoch=Optional[7], metadata=, commitTimestamp=1745134567892, expireTimestamp=None)可以直观看到「Key 是 [group,topic,partition],Value 是 OffsetAndMetadata」。
第 7 节有 Python 代码版本(inspect_consumer_offsets.py),不依赖命令行 formatter。
6.6 「offset 不见了」的几个常见原因
offsets.retention.minutes(默认 7 天):消费者下线超过这个时间,offset 会被 Compaction 清掉(Kafka 2.0 后改为 7 天,1.x 是 1 天)。表现是「停了一周再启动,发现 group 是新的」。group.instance.id没设:滚动发布时实例频繁加入退出,触发 Rebalance 风暴,可能在某次 join 失败后 group 被认为空了,offset 被清。- 手动
--reset-offsets操作:运维误操作直接重置成 earliest/latest。 __consumer_offsets副本不足:单副本时 Broker 挂了直接丢。生产必须 3 副本 + min.insync.replicas=2。
7. 配套代码(code/ 目录)
本章配套代码:
auto_commit_pitfall.py:演示自动提交丢/重消息的两种场景manual_commit_correct.py:手动提交模板(同步、异步、混合三种)seek_by_timestamp.py:按时间戳回溯消费idempotent_consumer.py:用 sqlite 模拟去重表的幂等消费者inspect_consumer_offsets.py:用 Python 直接消费__consumer_offsets,解析二进制结构
每个脚本独立可运行,命令在脚本头部注释里。
8. 与其他 MQ 的对比
| 维度 | Kafka | RabbitMQ | RocketMQ |
|---|---|---|---|
| 位点谁管 | 客户端(offset 是 long,集中存 __consumer_offsets) | Broker(每条消息一个 ACK,用完即删) | Broker(CommitLog + ConsumeQueue + offset store) |
| 重放 | ✅ 任意 offset / 时间戳 | ❌ 删除即不能重放 | ✅ 重置消费位点 |
| 提交粒度 | 分区级(一次提交一个 (topic,partition,offset)) | 单条 ACK / 批量 ACK | 单条 ACK(实际批量提交) |
| 自动确认风险 | 自动提交丢/重并存 | auto-ack 直接消息丢 | 「同步消费 + 异步提交」也类似 |
| EOS | 幂等 Producer + 事务(第 13 章) | 无原生 EOS | 事务消息(Half + Commit/Rollback) |
📌 本质差异:Kafka 把 offset 暴露给消费者,所以消息「被消费」≠「被删除」。这才有了重放、回溯、流处理。RabbitMQ 是「消息→分发→消费→删除」的一次性管道,不能回头。
9. 小结
- 三种语义是「端到端」的:At Most Once / At Least Once / Exactly Once,对应「先提交后处理」「先处理后提交」「业务幂等或 EOS 协议」。
- 自动提交两种陷阱:处理慢 → 丢消息(提交了未处理);处理快 → 重消息(处理完未提交就崩)。生产强烈建议关闭。
- 手动提交模板:处理 → 同步提交 / 异步提交 / 混合(异步 + 关闭同步兜底)。提交的是
next offset,不是「最后处理的 offset」。 - 业务幂等四种模式:唯一约束、Redis SETNX、去重表、状态机。99% 的业务靠 At Least Once + 业务幂等就够了,不必上 EOS。
seek系列让 Kafka 的「可回溯性」落地:按 offset、按时间戳,都能精准定位。线上排障必备。__consumer_offsets是个 50 分区、Compaction 策略的内部 Topic,每个 group 的 offset 集中在一个分区,Coordinator 就是这个分区的 Leader。
下一章我们进入 Kafka 真正的「灵魂」:幂等 Producer + 事务 + EOS 协议——这才是 Kafka 区别于其他 MQ 的杀手锏。
10. 面试高频题
Q1. At Most Once / At Least Once / Exactly Once 三种语义在 Kafka 里如何实现?
考察点:消息语义的工程含义 + Kafka 的实现机制。
标准答案:
- At Most Once:
enable.auto.commit=true+ 处理时长 >auto.commit.interval.ms,等价于「先提交后处理」,崩溃时已提交但未处理 → 丢。 - At Least Once:
enable.auto.commit=false+ 业务处理成功后再commitSync()/commitAsync(),崩溃时未提交 → 重新拉取处理 → 可能重复。 - Exactly Once:① At Least Once + 业务幂等(业务侧实现);② 幂等 Producer + 事务 Producer + Read Committed Consumer(Kafka 协议层 EOS,第 13 章)。
加分项:明确指出 EOS 是「端到端」的协议,单独打开 enable.idempotence=true 只能保证 Producer 单 Session 单分区 不重复,不能让消费侧 EOS。
Q2. 自动提交 offset 会丢消息还是重消息?为什么?
考察点:自动提交的真实行为。
标准答案:两种都会。
- 丢:业务处理时长 >=
auto.commit.interval.ms→ 第二次poll()触发了自动提交,把还没处理完的 offset 提交上去;进程崩溃后从已提交位点继续,前面那批消息丢了。 - 重:业务处理很快但还没到 commit 周期 → 进程崩溃 → 重启后从上次提交位点拉,已经处理过的再处理一遍。
易错点:很多人以为「自动提交是后台线程定时提交」,其实是在 poll() 时检查时间戳触发,所以在两次 poll() 中间崩溃,commit 不会发生。
Q3. commitSync 与 commitAsync 区别?生产怎么用?
考察点:API 选型。
标准答案:
commitSync:阻塞、有重试、保证提交成功才返回;缺点是 RTT 阻塞主循环,吞吐受影响。commitAsync:非阻塞、返回失败回调;不会自动重试(避免重试到旧的 offset 覆盖了新的)。- 最佳实践:常态用
commitAsync提高吞吐,优雅关闭 / 异常退出时用commitSync兜底,确保最终位点落到 Broker。
加分项:异步提交回调里如果做重试,要做「offset 单调比较」,避免覆盖更新的提交。
Q4. 业务侧如何保证消费幂等?说几种实现方式。
考察点:工程实战。
标准答案:
- 业务主键 + DB 唯一约束:
UNIQUE索引,重复插入捕获IntegrityError。 - Redis SETNX:
SET key 1 NX EX 86400,TTL 自动清理。 - 去重表 + TTL:单独
consumer_dedup表,记录 msg_id + 处理时间。 - 状态机:
UPDATE ... WHERE status='PENDING',依赖业务状态字段做条件更新。
加分项:指出**「dedupe 标记 + 业务写入」最好放在同一事务**,否则可能「标记成功业务失败」或反过来,导致幂等失效。
Q5. 怎样按时间戳回溯消费?底层原理是什么?
考察点:offsets_for_times 的使用 + .timeindex 索引。
标准答案:
- 调用
consumer.offsetsForTimes({TopicPartition: timestamp_ms}),Broker 返回≥ ts的最早 offset。 - 对每个分区调
consumer.seek(TopicPartition(t, p, offset))。 - 底层 Broker 用
.timeindex文件(每条 (timestamp, offset) 索引项),按 timestamp 二分查找定位到 segment,再在 segment 里精确定位 offset。
易错点:offsets_for_times 调用前必须确保已被分配该分区(先 poll() 一次或 assign()),否则报 「No current assignment」。
Q6. __consumer_offsets 是什么?里面存了什么?
考察点:内部 Topic 的设计。
标准答案:
- Kafka 0.9 引入的内部 Topic,存放消费者组的位点和元数据。
- 默认 50 个分区、副本因子 3、
cleanup.policy=compact。 - Key = (group_id, topic, partition);Value = (offset, leader_epoch, metadata, commit_timestamp)。
- 每次
commit()都是一次追加,旧记录通过 Compaction 清理掉。 partition = hash(group_id) % 50决定该 group 的位点落在哪个分区,分区 Leader 所在 Broker 就是该 group 的 Coordinator。
加分项:指出还存有 GroupMetadata(成员列表、generation、分配方案),用于 Coordinator 故障转移时重建。
Q7. auto.offset.reset 三个值的区别?什么时候生效?
考察点:默认行为陷阱。
标准答案:
latest(默认):从最新 offset 开始消费。earliest:从最早 offset 开始消费。none:没有已提交 offset 时直接抛异常。
只在「该 group 在该分区没有已提交 offset」时生效。一旦提交过,下次启动从已提交位点继续,与该参数无关。
易错点:上线第一次跑消费者,没有提交过 offset,又是默认 latest,结果「之前生产的消息全跳过了」——这是新人最常见的踩坑。
Q8. Lag 报警怎么定义?怎么排查?
考察点:可观测性 + 排障。
标准答案:
- Lag = LEO(Log End Offset,分区末尾) - CURRENT-OFFSET(已提交位点)。
- JMX 指标:
kafka.consumer:type=consumer-fetch-manager-metrics,client-id=*,topic=*,partition=*的records-lag-max。 - 命令行:
kafka-consumer-groups.sh --describe --group X。
排查步骤:
- 看 Lag 在哪些分区,是单分区飙升还是全局高。单分区飙升通常是 Key 热点 / 那个分区的 Leader 慢。
- 看消费者实例是不是在 Rebalance(CURRENT-OFFSET 不更新且 CONSUMER-ID 在变)。
- 看消费者业务处理时长是否 >
max.poll.interval.ms触发踢出。 - 看下游慢(DB / 第三方 API),可以加并发或异步处理。
🎯 小练习:写一个消费者,开启自动提交,让
process()函数time.sleep(8),把auto.commit.interval.ms设为 5000,跑一会儿后kill -9进程,重启,看kafka-consumer-groups.sh的 Lag 和 DB 里到底处理了多少条 —— 你会亲眼看到「丢消息」的发生。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
#!/usr/bin/env python3
"""
auto_commit_pitfall.py
======================
演示 enable.auto.commit=true 的两种陷阱:
1) 处理慢 -> 已提交 offset 跑到「真处理位点」前 -> 进程崩溃后丢消息
2) 处理快 -> 未到 commit 周期就崩溃 -> 重启后已处理消息会被重复
用法(本地需要 Kafka 在 127.0.0.1:9092):
pip install confluent-kafka
# 终端 A —— 先生产 50 条 0..49
python auto_commit_pitfall.py produce
# 终端 B —— 演示「丢」(处理 8s,commit interval 5s)
python auto_commit_pitfall.py lose
# 终端 C —— 演示「重」(处理 0.05s,commit interval 5s)
python auto_commit_pitfall.py dup
每个 demo 模式都会在「关键时刻」抛 SystemExit 模拟进程崩溃,
然后请重新跑一次 lose/dup(同一 group),观察 CURRENT-OFFSET 与
真实处理日志的差距:
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \\
--describe --group demo-pitfall-lose
"""
import os
import sys
import time
import json
import logging
from datetime import datetime
from confluent_kafka import Producer, Consumer
from confluent_kafka.admin import AdminClient, NewTopic
BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
TOPIC = "learn.12.pitfall"
N = 50
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s",
datefmt="%H:%M:%S",
)
log = logging.getLogger("pitfall")
def ensure_topic():
admin = AdminClient({"bootstrap.servers": BOOTSTRAP})
md = admin.list_topics(timeout=5).topics
if TOPIC not in md:
log.info(f"creating topic {TOPIC} (3 partitions)")
fs = admin.create_topics([NewTopic(TOPIC, num_partitions=3, replication_factor=1)])
for t, f in fs.items():
try: f.result()
except Exception as e: log.warning(f"create topic: {e}")
def produce():
ensure_topic()
p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 5})
for i in range(N):
body = json.dumps({"id": i, "ts": datetime.utcnow().isoformat()})
p.produce(TOPIC, key=str(i % 3), value=body.encode())
p.flush(10)
log.info(f"produced {N} msgs to {TOPIC}")
def consume(group, process_seconds, commit_interval_ms, crash_after_seconds):
ensure_topic()
c = Consumer({
"bootstrap.servers": BOOTSTRAP,
"group.id": group,
"enable.auto.commit": True,
"auto.commit.interval.ms": commit_interval_ms,
"auto.offset.reset": "earliest",
# 不打开手动 fetch.min/max 等,保持默认让效果更明显
})
c.subscribe([TOPIC])
started = time.time()
handled_ids = []
log.info(
f"consumer group={group} process={process_seconds}s commit_int={commit_interval_ms}ms "
f"crash_after={crash_after_seconds}s"
)
try:
while True:
msg = c.poll(timeout=1.0)
now = time.time()
if now - started >= crash_after_seconds:
# 模拟进程被 kill -9:直接 os._exit,不让 client flush/commit
log.warning("💥 SIMULATED CRASH (os._exit) — last commit may NOT be the truth!")
os._exit(137)
if msg is None:
continue
if msg.error():
log.error(msg.error())
continue
payload = json.loads(msg.value())
log.info(
f" poll p{msg.partition()}@offset={msg.offset()} id={payload['id']} "
f"-> processing for {process_seconds}s"
)
time.sleep(process_seconds) # 模拟业务
handled_ids.append(payload["id"])
log.info(f" done id={payload['id']}, total handled={len(handled_ids)}")
finally:
c.close()
def main():
if len(sys.argv) < 2:
print(__doc__); return
cmd = sys.argv[1]
if cmd == "produce":
produce()
elif cmd == "lose":
# 处理 8s + commit interval 5s + 第 6 秒崩溃
# → poll 第一条 -> commit 立即提交 offset -> 处理 8s 还没完 -> 6s 时崩
# → 重启从「已提交+1」继续,第 0 条永远没处理(实际上更糟:第 1~K 条全部丢)
consume(
group="demo-pitfall-lose",
process_seconds=8.0,
commit_interval_ms=5000,
crash_after_seconds=6.0,
)
elif cmd == "dup":
# 处理 0.05s + commit interval 5s + 第 3 秒崩溃
# → 3s 内已处理几十条,但 commit 还没触发 -> 重启后从 0 重新拉
consume(
group="demo-pitfall-dup",
process_seconds=0.05,
commit_interval_ms=5000,
crash_after_seconds=3.0,
)
else:
print("unknown:", cmd); print(__doc__)
if __name__ == "__main__":
main()python
#!/usr/bin/env python3
"""
idempotent_consumer.py
======================
业务侧幂等示例:用 sqlite 做「去重表 + 业务表」,模拟「订单已支付」事件的
At Least Once 安全消费。
关键设计:
1) order_event 表:业务表,order_id 是 PRIMARY KEY,重复插入直接 IntegrityError
2) consumer_dedup 表:审计表,记录每条消息的 (msg_id, processed_at)
3) 业务处理 + 去重写入放在同一个 sqlite 事务里,保证「要么都成功,要么都失败」
4) Kafka commit 仅在事务成功后才执行
无论 Kafka 重投递多少次,业务表里只会有一条记录。
跑两遍可以亲眼看到「dup, skip」日志。
用法:
pip install confluent-kafka
# 先生产几条订单
python idempotent_consumer.py produce 5
# 启动幂等消费者
python idempotent_consumer.py consume
# ⚠️ 把消费者杀掉再重启 / 重置 offset:依然只有一条订单记录
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \
--group demo-idempotent --reset-offsets --to-earliest --execute --topic learn.12.idem
"""
import os
import sys
import json
import time
import sqlite3
import logging
import hashlib
from datetime import datetime
from confluent_kafka import Producer, Consumer
BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
TOPIC = "learn.12.idem"
DB = os.environ.get("IDEM_DB", "/tmp/learn_kafka_idempotent.sqlite")
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s", datefmt="%H:%M:%S")
log = logging.getLogger("idem")
def init_db():
conn = sqlite3.connect(DB)
conn.execute("""
CREATE TABLE IF NOT EXISTS order_event (
order_id TEXT PRIMARY KEY,
user_id INTEGER,
amount REAL,
status TEXT,
created_at TEXT
)
""")
conn.execute("""
CREATE TABLE IF NOT EXISTS consumer_dedup (
msg_id TEXT PRIMARY KEY,
consumer TEXT,
processed_at TEXT
)
""")
conn.commit()
return conn
def produce(n):
p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 5})
for i in range(n):
oid = f"ORD-2026-{i:04d}"
body = json.dumps({
"order_id": oid,
"user_id": 1000 + i,
"amount": round(99.9 + i, 2),
"ts": datetime.utcnow().isoformat(),
})
p.produce(TOPIC, key=oid, value=body.encode())
p.flush(10)
log.info(f"produced {n} order events")
def compute_msg_id(msg, payload):
# 优先用业务主键;没有的话 fallback 到 (topic, partition, offset)
if "order_id" in payload:
raw = f"{payload['order_id']}".encode()
else:
raw = f"{msg.topic()}|{msg.partition()}|{msg.offset()}".encode()
return hashlib.sha1(raw).hexdigest()[:16]
def consume():
conn = init_db()
c = Consumer({
"bootstrap.servers": BOOTSTRAP,
"group.id": "demo-idempotent",
"enable.auto.commit": False,
"auto.offset.reset": "earliest",
})
c.subscribe([TOPIC])
try:
while True:
msg = c.poll(1.0)
if msg is None: continue
if msg.error():
log.error(msg.error()); continue
try:
payload = json.loads(msg.value())
except Exception:
log.exception("bad msg, skip + commit")
c.commit(message=msg, asynchronous=False)
continue
msg_id = compute_msg_id(msg, payload)
try:
with conn: # sqlite context manager = 一个事务
conn.execute(
"INSERT INTO consumer_dedup(msg_id, consumer, processed_at) VALUES (?,?,?)",
(msg_id, "demo-idempotent", datetime.utcnow().isoformat()),
)
conn.execute(
"INSERT INTO order_event(order_id,user_id,amount,status,created_at) "
"VALUES (?,?,?,?,?)",
(payload["order_id"], payload["user_id"], payload["amount"],
"PAID", datetime.utcnow().isoformat()),
)
log.info(f" ✅ first time, processed order={payload['order_id']} msg_id={msg_id}")
except sqlite3.IntegrityError:
# 唯一约束冲突 -> 已处理过,幂等跳过
log.info(f" ⏭️ dup, skip order={payload['order_id']} msg_id={msg_id}")
# 不管首处理还是跳过都要 commit offset
c.commit(message=msg, asynchronous=False)
finally:
c.close()
# 给个最终统计
n_orders = conn.execute("SELECT COUNT(*) FROM order_event").fetchone()[0]
n_dedup = conn.execute("SELECT COUNT(*) FROM consumer_dedup").fetchone()[0]
log.info(f"FINAL: {n_orders} unique orders in business table; "
f"{n_dedup} dedup records (= 真实处理过的次数)")
def main():
if len(sys.argv) < 2:
print(__doc__); return
cmd = sys.argv[1]
if cmd == "produce":
n = int(sys.argv[2]) if len(sys.argv) > 2 else 5
produce(n)
elif cmd == "consume":
consume()
else:
print(__doc__)
if __name__ == "__main__":
main()python
#!/usr/bin/env python3
"""
inspect_consumer_offsets.py
===========================
直接消费 Kafka 内部 Topic `__consumer_offsets`,解析其二进制 Key/Value 结构,
打印出每条 OffsetCommit / GroupMetadata 记录。
这是「Kafka 自己吃自己的狗粮」最直观的一个演示:
- Key 描述「这是哪个 group 的哪个 (topic, partition) 的位点」
- Value 描述「offset / leader_epoch / metadata / commit_timestamp」
- 整个 topic 用 cleanup.policy=compact 保留每个 Key 的最新 Value
二进制结构(OffsetCommit, version 3+):
Key:
int16 version
string group_id (int16 length + bytes)
string topic
int32 partition
Value:
int16 version
int64 offset
int32 leader_epoch (v3+,无则 -1)
string metadata
int64 commit_timestamp
int64 expire_timestamp (旧版才有)
用法:
pip install confluent-kafka
python inspect_consumer_offsets.py
# 然后另一个终端跑任意消费者并 commit,回到这里看到对应输出。
"""
import os
import struct
import logging
from confluent_kafka import Consumer
BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s", datefmt="%H:%M:%S")
log = logging.getLogger("inspect")
class Reader:
def __init__(self, buf):
self.buf = buf; self.pos = 0
def i16(self):
v = struct.unpack(">h", self.buf[self.pos:self.pos+2])[0]; self.pos += 2; return v
def i32(self):
v = struct.unpack(">i", self.buf[self.pos:self.pos+4])[0]; self.pos += 4; return v
def i64(self):
v = struct.unpack(">q", self.buf[self.pos:self.pos+8])[0]; self.pos += 8; return v
def str(self):
n = self.i16()
if n < 0: return None
v = self.buf[self.pos:self.pos+n].decode("utf-8", errors="replace"); self.pos += n; return v
def remain(self): return len(self.buf) - self.pos
def parse_offset_key(key):
r = Reader(key)
version = r.i16()
# version 0 / 1 = OffsetCommit
# version 2 = GroupMetadata
if version not in (0, 1):
return None
group = r.str()
topic = r.str()
partition = r.i32()
return {"kind": "OffsetCommit", "version": version,
"group": group, "topic": topic, "partition": partition}
def parse_group_meta_key(key):
r = Reader(key)
version = r.i16()
if version != 2:
return None
group = r.str()
return {"kind": "GroupMetadata", "version": version, "group": group}
def parse_offset_value(val):
if val is None:
return {"tombstone": True}
r = Reader(val)
try:
version = r.i16()
offset = r.i64()
leader_epoch = r.i32() if version >= 3 else -1
metadata = r.str()
commit_ts = r.i64()
expire_ts = r.i64() if version <= 1 else None
return {"version": version, "offset": offset, "leader_epoch": leader_epoch,
"metadata": metadata, "commit_timestamp": commit_ts,
"expire_timestamp": expire_ts}
except Exception as e:
return {"raw_len": len(val), "parse_error": str(e)}
def main():
c = Consumer({
"bootstrap.servers": BOOTSTRAP,
"group.id": "inspect-internal-" + str(os.getpid()),
"enable.auto.commit": False,
"auto.offset.reset": "latest", # 只看新提交,避免历史刷屏
})
c.subscribe(["__consumer_offsets"])
log.info("listening to __consumer_offsets ... 在另一个终端跑任意 commit 即可看到输出")
try:
while True:
msg = c.poll(1.0)
if msg is None: continue
if msg.error():
log.error(msg.error()); continue
key = msg.key() or b""
val = msg.value()
if len(key) < 2:
continue
# 先尝试当作 OffsetCommit
parsed_key = parse_offset_key(key)
if parsed_key:
parsed_val = parse_offset_value(val)
log.info(
f"📌 OffsetCommit partition_in_internal=p{msg.partition()} "
f"group={parsed_key['group']!r} topic={parsed_key['topic']!r} "
f"part={parsed_key['partition']} → offset={parsed_val.get('offset')} "
f"leader_epoch={parsed_val.get('leader_epoch')} "
f"commit_ts={parsed_val.get('commit_timestamp')}"
f"{' [TOMBSTONE]' if parsed_val.get('tombstone') else ''}"
)
continue
# 否则尝试 GroupMetadata
gm = parse_group_meta_key(key)
if gm:
size = len(val) if val else 0
log.info(
f"📁 GroupMetadata partition_in_internal=p{msg.partition()} "
f"group={gm['group']!r} value_size={size} bytes "
f"{'(tombstone)' if val is None else '(成员/分配方案 二进制)'}"
)
continue
log.info(f"unknown record key_version={struct.unpack('>h', key[:2])[0]}")
except KeyboardInterrupt:
pass
finally:
c.close()
if __name__ == "__main__":
main()python
#!/usr/bin/env python3
"""
manual_commit_correct.py
========================
手动提交三种正确姿势:
1) commit_sync_each_msg : 每条消息处理完同步提交(最稳,吞吐最低)
2) commit_async_each_msg: 每条消息异步提交(高吞吐,关闭时同步兜底)
3) commit_batch : 批量处理 + 批量提交(生产推荐)
用法:
pip install confluent-kafka
# 先用 auto_commit_pitfall.py 的 produce 命令灌些数据
python manual_commit_correct.py sync
python manual_commit_correct.py async
python manual_commit_correct.py batch
每次跑都用同一个 group(demo-manual),可以观察连续运行
不会丢/重——相比 auto 模式有显著区别。
关键 API:
- consumer.commit(message=msg, asynchronous=False) # 同步、提交单条
- consumer.commit(asynchronous=True) # 异步、提交当前 highwater
- consumer.commit(offsets=[TopicPartition(t, p, off)]) # 精细化
"""
import os
import sys
import time
import json
import logging
from confluent_kafka import Consumer, TopicPartition
BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
TOPIC = "learn.12.pitfall"
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s", datefmt="%H:%M:%S")
log = logging.getLogger("manual")
def make_consumer(group):
return Consumer({
"bootstrap.servers": BOOTSTRAP,
"group.id": group,
"enable.auto.commit": False, # 关键:关闭自动提交
"auto.offset.reset": "earliest",
"session.timeout.ms": 10000,
"max.poll.interval.ms": 60000,
})
def commit_sync_each_msg():
c = make_consumer("demo-manual-sync")
c.subscribe([TOPIC])
try:
while True:
msg = c.poll(1.0)
if msg is None: continue
if msg.error():
log.error(msg.error()); continue
payload = json.loads(msg.value())
try:
process(payload)
except Exception:
log.exception("process failed, will retry next round")
continue
c.commit(message=msg, asynchronous=False) # 同步提交,确认成功
log.info(f" committed p{msg.partition()}@{msg.offset() + 1}")
finally:
c.close()
def commit_async_each_msg():
c = make_consumer("demo-manual-async")
c.subscribe([TOPIC])
def on_commit(err, partitions):
if err:
log.error(f"async commit failed: {err}")
else:
for tp in partitions:
log.debug(f"async commit ok p{tp.partition}@{tp.offset}")
try:
while True:
msg = c.poll(1.0)
if msg is None: continue
if msg.error():
log.error(msg.error()); continue
payload = json.loads(msg.value())
try:
process(payload)
except Exception:
log.exception("process failed, skip commit")
continue
c.commit(message=msg, asynchronous=True)
log.info(f" asynchronously committed p{msg.partition()}@{msg.offset() + 1}")
except KeyboardInterrupt:
pass
finally:
# 关闭前同步提交一次「兜底」
try:
c.commit(asynchronous=False)
log.info("final sync commit done")
finally:
c.close()
def commit_batch():
c = make_consumer("demo-manual-batch")
c.subscribe([TOPIC])
BATCH = 50
try:
while True:
buf = []
t0 = time.time()
while len(buf) < BATCH and time.time() - t0 < 1.0:
msg = c.poll(0.5)
if msg is None: continue
if msg.error():
log.error(msg.error()); continue
buf.append(msg)
if not buf:
continue
# 业务批处理:通常是一次性 INSERT INTO ... VALUES (?,?), ...
try:
process_batch([json.loads(m.value()) for m in buf])
except Exception:
log.exception("batch failed, will retry")
continue
# 取每个分区的最大 offset,统一提交(避免漏提交)
tp_max = {}
for m in buf:
key = (m.topic(), m.partition())
if m.offset() > tp_max.get(key, -1):
tp_max[key] = m.offset()
offsets = [TopicPartition(t, p, off + 1) for (t, p), off in tp_max.items()]
c.commit(offsets=offsets, asynchronous=False)
log.info(f" committed batch: {[(o.partition, o.offset) for o in offsets]}")
finally:
c.close()
def process(payload):
# 业务!务必幂等
time.sleep(0.05)
def process_batch(payloads):
time.sleep(0.05)
log.info(f" process_batch size={len(payloads)}")
def main():
if len(sys.argv) < 2:
print(__doc__); return
cmd = sys.argv[1]
{"sync": commit_sync_each_msg,
"async": commit_async_each_msg,
"batch": commit_batch}[cmd]()
if __name__ == "__main__":
main()python
#!/usr/bin/env python3
"""
seek_by_timestamp.py
====================
按时间戳定位 offset,是线上排障的「时光机」:
> 「线上昨天 14:00~14:05 这一段订单出问题了,能不能只回放这 5 分钟的消息?」
底层原理:Broker 利用 .timeindex 文件(每条 (timestamp, offset) 索引项)
做二分查找,返回 >= 给定 ts 的最早消息 offset。
用法:
pip install confluent-kafka
python seek_by_timestamp.py produce 200 # 灌 200 条带 timestamp 的消息
python seek_by_timestamp.py replay "2026-04-17 12:00:00" "2026-04-17 12:00:05"
# 只回放某 5 秒内的消息
"""
import os
import sys
import time
import json
import logging
from datetime import datetime, timezone
from confluent_kafka import Producer, Consumer, TopicPartition
BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
TOPIC = "learn.12.timeline"
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s", datefmt="%H:%M:%S")
log = logging.getLogger("seek")
def produce(n):
p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 5})
base = int(time.time() * 1000)
for i in range(n):
body = json.dumps({"id": i, "ts": base + i * 1000})
# 显式指定 timestamp,让消费者能按时间戳定位
p.produce(TOPIC, key=str(i % 3), value=body.encode(), timestamp=base + i * 1000)
if i % 50 == 0:
p.poll(0)
p.flush(10)
log.info(f"produced {n} msgs to {TOPIC} (timestamp 1s/条递增)")
def parse_ts(s):
"""把 'YYYY-MM-DD HH:MM:SS' 解析成毫秒时间戳(本地时区)"""
dt = datetime.strptime(s, "%Y-%m-%d %H:%M:%S").astimezone()
return int(dt.timestamp() * 1000)
def replay(start_str, end_str):
start_ts = parse_ts(start_str)
end_ts = parse_ts(end_str)
log.info(f"replay window: {start_str} ({start_ts}) ~ {end_str} ({end_ts})")
c = Consumer({
"bootstrap.servers": BOOTSTRAP,
"group.id": "demo-replay-by-ts",
"enable.auto.commit": False,
"auto.offset.reset": "earliest",
})
c.subscribe([TOPIC])
# 必须先 poll 一次(或者 assign),否则 seek 报「No current assignment」
c.poll(2.0)
md = c.list_topics(TOPIC, timeout=5.0).topics[TOPIC]
tps_query = [TopicPartition(TOPIC, p, start_ts) for p in md.partitions]
log.info(f"querying offsets_for_times for {len(tps_query)} partitions...")
offsets = c.offsets_for_times(tps_query, timeout=5.0)
# offsets 列表里每个 tp.offset 已经是「>= start_ts 的最早 offset」
for tp in offsets:
if tp.offset == -1:
log.info(f" P{tp.partition}: 无消息 >= start_ts (整个分区都比该时间早或没有数据)")
continue
log.info(f" P{tp.partition}: seek to offset {tp.offset}")
c.seek(tp)
received = 0
try:
while True:
msg = c.poll(1.0)
if msg is None:
# 等待 3s 没有新消息就结束
continue
if msg.error():
log.error(msg.error()); continue
ts_type, ts = msg.timestamp()
if ts > end_ts:
log.info(f" P{msg.partition()}@{msg.offset()} ts={ts} > end_ts → 该分区窗口结束")
# 简单起见全部退出(生产里应跟踪每个分区分别结束)
break
received += 1
payload = json.loads(msg.value())
log.info(f" REPLAY p{msg.partition()}@{msg.offset()} ts={ts} id={payload['id']}")
if received >= 1000:
break
finally:
c.close()
log.info(f"window replayed: {received} msgs")
def main():
if len(sys.argv) < 2:
print(__doc__); return
cmd = sys.argv[1]
if cmd == "produce":
n = int(sys.argv[2]) if len(sys.argv) > 2 else 200
produce(n)
elif cmd == "replay":
start = sys.argv[2]; end = sys.argv[3]
replay(start, end)
else:
print("unknown:", cmd); print(__doc__)
if __name__ == "__main__":
main()auto_commit_pitfall.py ↗ · idempotent_consumer.py ↗ · inspect_consumer_offsets.py ↗ · manual_commit_correct.py ↗ · seek_by_timestamp.py ↗