Skip to content

第 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 提交一次。注意几个关键点:

  1. 提交时机是下一次 poll,不是某个后台线程定时主动提交(很多人误以为有后台线程)。
  2. 提交的 offset 是「已经 poll 出来交给业务的最大 offset + 1」,而不是「业务真的处理完的 offset」——Kafka 客户端根本不知道你业务有没有处理完!
  3. 因此自动提交的语义本质是 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风险
处理快、间隔短50ms1s重复窗口很小
处理慢、间隔长8s5s「丢」非常严重
处理很快、间隔默认50ms5s「重」很严重

📌 结论:自动提交是「为了懒人方便」的开关。任何一个有钱、有人、有责任心的业务都应该 关闭自动提交,改用手动提交。


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。生产环境推荐显式设置为 earliestnone,避免「上线 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_offsets
  • LOG-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 不见了」的几个常见原因

  1. offsets.retention.minutes(默认 7 天):消费者下线超过这个时间,offset 会被 Compaction 清掉(Kafka 2.0 后改为 7 天,1.x 是 1 天)。表现是「停了一周再启动,发现 group 是新的」。
  2. group.instance.id 没设:滚动发布时实例频繁加入退出,触发 Rebalance 风暴,可能在某次 join 失败后 group 被认为空了,offset 被清。
  3. 手动 --reset-offsets 操作:运维误操作直接重置成 earliest/latest。
  4. __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 的对比

维度KafkaRabbitMQRocketMQ
位点谁管客户端(offset 是 long,集中存 __consumer_offsetsBroker(每条消息一个 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. 小结

  1. 三种语义是「端到端」的:At Most Once / At Least Once / Exactly Once,对应「先提交后处理」「先处理后提交」「业务幂等或 EOS 协议」。
  2. 自动提交两种陷阱:处理慢 → 丢消息(提交了未处理);处理快 → 重消息(处理完未提交就崩)。生产强烈建议关闭
  3. 手动提交模板:处理 → 同步提交 / 异步提交 / 混合(异步 + 关闭同步兜底)。提交的是 next offset,不是「最后处理的 offset」。
  4. 业务幂等四种模式:唯一约束、Redis SETNX、去重表、状态机。99% 的业务靠 At Least Once + 业务幂等就够了,不必上 EOS。
  5. seek 系列让 Kafka 的「可回溯性」落地:按 offset、按时间戳,都能精准定位。线上排障必备。
  6. __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 Onceenable.auto.commit=true + 处理时长 > auto.commit.interval.ms,等价于「先提交后处理」,崩溃时已提交但未处理 → 丢。
  • At Least Onceenable.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. commitSynccommitAsync 区别?生产怎么用?

考察点:API 选型。

标准答案

  • commitSync:阻塞、有重试、保证提交成功才返回;缺点是 RTT 阻塞主循环,吞吐受影响。
  • commitAsync:非阻塞、返回失败回调;不会自动重试(避免重试到旧的 offset 覆盖了新的)。
  • 最佳实践:常态用 commitAsync 提高吞吐,优雅关闭 / 异常退出时用 commitSync 兜底,确保最终位点落到 Broker。

加分项:异步提交回调里如果做重试,要做「offset 单调比较」,避免覆盖更新的提交。


Q4. 业务侧如何保证消费幂等?说几种实现方式。

考察点:工程实战。

标准答案

  1. 业务主键 + DB 唯一约束UNIQUE 索引,重复插入捕获 IntegrityError
  2. Redis SETNXSET key 1 NX EX 86400,TTL 自动清理。
  3. 去重表 + TTL:单独 consumer_dedup 表,记录 msg_id + 处理时间。
  4. 状态机UPDATE ... WHERE status='PENDING',依赖业务状态字段做条件更新。

加分项:指出**「dedupe 标记 + 业务写入」最好放在同一事务**,否则可能「标记成功业务失败」或反过来,导致幂等失效。


Q5. 怎样按时间戳回溯消费?底层原理是什么?

考察点offsets_for_times 的使用 + .timeindex 索引。

标准答案

  1. 调用 consumer.offsetsForTimes({TopicPartition: timestamp_ms}),Broker 返回 ≥ ts 的最早 offset。
  2. 对每个分区调 consumer.seek(TopicPartition(t, p, offset))
  3. 底层 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

排查步骤

  1. 看 Lag 在哪些分区,是单分区飙升还是全局高。单分区飙升通常是 Key 热点 / 那个分区的 Leader 慢。
  2. 看消费者实例是不是在 Rebalance(CURRENT-OFFSET 不更新且 CONSUMER-ID 在变)。
  3. 看消费者业务处理时长是否 > max.poll.interval.ms 触发踢出。
  4. 看下游慢(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 ↗