Skip to content

第 13 章 幂等与事务:Exactly Once 的真相

目标读者:听过「Kafka 支持 Exactly Once」却以为「打开一个开关就行」、不知道幂等 Producer 的边界、对事务和 Read Committed 感觉云里雾里的同学。

学完你会:能用一张图讲清 PID + Epoch + Sequence 的去重逻辑;知道幂等的「单 Session、单分区」边界为什么必须存在;能写完整的事务 Producer + sendOffsetsToTransaction consume-process-produce 模板;能解释 Read Committed 与 LSO 的关系;能从 __transaction_state 看到 Coordinator 真实的两阶段提交日志;面试官问 EOS 时不再支支吾吾。


0. 导读:「Exactly Once」是 Kafka 的招牌也是误会重灾区

2017 年 Kafka 0.11 发布时,宣传口号是:

"Kafka now supports Exactly Once Semantics!"

整个社区都炸了——「消息中间件居然能 EOS?!」。但接下来几年,无数公司因为误用 EOS 进了坑:

  • enable.idempotence=true 当 EOS 万能药 → 跨 Session 还是会重
  • 用了事务 Producer,但消费者用 read_uncommitted → 读到了未 commit 的数据
  • 业务跨多个 Topic,但只在单 Topic 里搞事务 → 跨 Topic 的「读 → 写」不一致
  • 把 EOS 用在「Kafka → MySQL」链路上 → MySQL 不在 Kafka 事务里,根本没原子性

本章我们把 Kafka 真正的 EOS 协议层 拆透:幂等 Producer + 事务 Producer + __transaction_state Coordinator + read_committed Consumer,四件套缺一不可。学完你能精准说出:「Kafka EOS 适用于 Kafka 内部 consume-process-produce 链路(Streams 的世界),不适用于 Kafka → 外部系统的写入;那种场景需要业务侧幂等(第 12 章)」。


1. 幂等 Producer:单 Session 单分区的去重

1.1 没有幂等之前会发生什么

Producer 发送消息时,常见的失败重试链路是:

Producer --send--> Broker
Broker  --写入磁盘成功 + 复制到 ISR-->
Broker  --回 ACK--> 网络抖动丢了
Producer --超时--> 重发同一条消息
Broker  --又写入磁盘一份--> 同一条消息存了两份!

这就是「At Least Once Producer」的代价:消息写两份

业务表现:在线下单接口超时重试,结果 Kafka 里同一个订单事件出现两次,下游消费者吃两次。

1.2 PID + Epoch + Sequence 三件套

Kafka 0.11 引入 幂等 Producer,开关是:

properties
enable.idempotence=true     # 默认 true (3.0+),老版本要手动开
acks=all                    # 必须
max.in.flight.requests.per.connection=5  # ≤ 5(默认)
retries=Integer.MAX_VALUE   # 内置重试到底

打开后,Kafka 给每个 Producer 实例做了三件事:

  1. Producer ID(PID):每个 Producer 实例向 Broker 申请一个全集群唯一的 64-bit ID(InitProducerIdRequest),打开 Producer 时分配。
  2. Producer Epoch(int16):同一个 PID 可能被「重启的同名实例」抢用,Epoch 用来区分先后;后启动的实例 Epoch + 1,旧的会被 Broker 拒绝(详见 1.5 节 Fencing)。
  3. Sequence Number(int32):Producer 给每个分区单独维护一个递增序号,每条消息带上 (PID, Partition, Sequence)

Broker 端为每个 (PID, Partition) 维护「最近见过的 Sequence」:

  • 收到 seq = last + 1:正常,写入 + 更新 last = seq
  • 收到 seq <= last重复消息,直接丢弃,但 ACK 仍然返回(Producer 看起来发成功了)。
  • 收到 seq > last + 1乱序,返回 OUT_OF_ORDER_SEQUENCE_NUMBER,Producer 必须按顺序重发。

1.3 一张图讲清流程

Producer (PID=42, Epoch=0)        Broker (Topic-Partition: orders-0)
  │                                  │  state[42, orders-0].lastSeq = -1
  │                                  │
  │ ─send seq=0─────────────────────▶│  ✅ write seq=0, lastSeq=0
  │ ◀────ACK seq=0 (但网络丢)──── ✗  │
  │                                  │
  │ ─resend seq=0───────────────────▶│  🔁 seq=0 ≤ lastSeq=0
  │                                  │     ❌ 不再写入磁盘
  │ ◀────ACK seq=0 (本次成功)────────│     ✅ 仍返回 ACK
  │                                  │
  │ ─send seq=1─────────────────────▶│  ✅ write seq=1, lastSeq=1
  │ ◀────ACK seq=1───────────────────│
  │                                  │
  │ ─send seq=3 (跳号!)──────────────▶│  ❌ OUT_OF_ORDER_SEQUENCE
  │ ◀──err: seq must be 2─────────── │
  │                                  │
  │ ─resend seq=2 then seq=3─────────▶│  ✅ both written

注意几个细节:

  • Sequence 是按分区的:发往不同分区的消息有各自独立的 seq。这就解释了为什么幂等只能保证「单分区」内不重——跨分区的去重 Broker 没法判断。
  • max.in.flight.requests.per.connection ≤ 5:Broker 端只缓存最近 5 个 batch 的 seq,超过 5 就没法判断顺序了。3.0+ 默认 5。
  • 不能让 batch 改顺序:Broker 强制要求顺序写入,否则丢 OUT_OF_ORDER,Producer 自动按序重发。

1.4 幂等的边界(务必牢记!)

幂等 Producer 只能保证:

✅ 同一个 Producer 实例(同一个 PID + Epoch) ✅ 同一个分区 ✅ 在 Broker 内存里能记住的「最近 5 个 batch」窗口内 ✅ 不会有重复消息

幂等 Producer 不能保证:

❌ Producer 进程重启(新 PID 申请,旧 PID 状态被 Broker 清理)→ 跨 Session 重复 ❌ 同一条消息发到不同分区(key 改了 / 分区器变了)→ 跨分区无法去重 ❌ 多个 Producer 都发了同一条业务消息 → 不同 PID 互不知情 ❌ 真正的「业务级 Exactly Once」(业务幂等还是要做)

📌 生活类比:幂等 Producer 像「同一个柜员发同一张工单的去重」——这位柜员有自己的工号(PID),他工作期间发的每张单都有递增编号(seq),重复发同一张直接被自己识别。但换班之后(新 PID)、换柜台(新分区)、多个柜员同时发,都没法判断重复。要彻底解决,还得靠后面的「事务 + 业务唯一约束」。

1.5 PID Fencing(跨 Session 的「老实例屏蔽」)

幂等 Producer 重启会拿一个新 PID,但事务 Producer(下一节)会做一件特殊的事:通过固定的 transactional.id 让重启后的实例复用同一个 PID,并把 Epoch + 1。这样:

Producer 实例 A (transactional.id=tx-order-1) → PID=42, Epoch=7

                              进程崩溃

Producer 实例 B (transactional.id=tx-order-1) → 复用 PID=42, Epoch=8

                              意外:旧实例 A 的网络包还在路上

Broker 收到 A 发的请求 (PID=42, Epoch=7)
       → 比对当前 Epoch=8,发现「这是个被淘汰的旧实例!」
       → 返回 INVALID_PRODUCER_EPOCH
       → A 抛 ProducerFencedException → A 必须退出

这叫 Producer Fencing,是事务 Producer 在故障转移场景下保证「不会有两个实例同时活着」的关键。一会儿在事务那一节会用到。


2. 事务 Producer:把多条消息绑成一个原子单元

2.1 我们到底要解决什么问题

光有幂等还不够。考虑「Kafka Streams 的核心场景」——consume-process-produce

读 input topic → 处理(map/filter/aggregate)→ 写 output topic + 提交 input offset

三件事必须原子地发生:

  1. 写入 output topic 的所有消息(可能跨多个 Partition、多个 Topic)
  2. 把 input topic 的 offset 提交到 __consumer_offsets
  3. 要么全成功,要么全失败

如果只是幂等:

  • 写 output 成功了,offset 提交失败 → 重启后再处理一次 → output 多写
  • offset 提交成功了,写 output 失败 → 重启后跳过 → output 少写

要把这堆动作绑在一起,就需要 事务——Kafka 的事务 Producer 提供了一个迷你版「分布式 ACID」。

2.2 事务 Producer 的完整生命周期

python
from confluent_kafka import Producer
p = Producer({
    "bootstrap.servers": "127.0.0.1:9092",
    "enable.idempotence": True,
    "transactional.id": "tx-order-processor-1",  # 必填,跨 Session 唯一
    "transaction.timeout.ms": 60000,
})

p.init_transactions()                      # 1. 一次性初始化,绑定 PID/Epoch/Coordinator
while True:
    msgs = consumer.consume(num_messages=100, timeout=1.0)
    if not msgs: continue
    p.begin_transaction()                  # 2. 开启事务
    try:
        for m in msgs:
            handled = process(m)
            p.produce("output_topic", key=handled.key, value=handled.value)
        # 3. 把 input offset 也包进事务
        p.send_offsets_to_transaction(
            consumer.position(consumer.assignment()),
            consumer.consumer_group_metadata(),
        )
        p.commit_transaction()             # 4. 提交事务
    except Exception:
        p.abort_transaction()              # 5. 回滚
        raise

七个 API(Java 同名):

API作用
init_transactions()一次性:拿 PID/Epoch、找 Transaction Coordinator
begin_transaction()标记本次事务开始
produce(...)发送消息(消息会带上 PID + Epoch + Sequence + tx)
send_offsets_to_transaction()把消费位点也作为本事务的一部分提交
commit_transaction()二阶段提交:写 PREPARE_COMMIT → COMMIT marker
abort_transaction()二阶段回滚:写 PREPARE_ABORT → ABORT marker
(隐式) ProducerFencedException收到该异常 → 必须 close + 退出,进入 fencing

2.3 transactional.id 的精髓:跨 Session 复用 PID

transactional.id 是用户指定的字符串,比如 tx-order-processor-1。Producer 启动调 initTransactions() 时:

  1. Producer 把 transactional.id 发给 Transaction Coordinator
  2. Coordinator 查 __transaction_state
    • 如果该 transactional.id 没记录 → 分配新 PID + Epoch=0
    • 如果有记录 → 复用同一个 PID,把 Epoch +1
  3. Coordinator 把当前 Epoch 写回 __transaction_state,并通知所有 Broker 持久化

效果:同名 Producer 重启后,旧实例发的请求会因 Epoch 不匹配被 Fencing 掉,杜绝「两个实例同时写」的脑裂。

📌 transactional.id 必须全集群唯一。在 Kafka Streams 里,每个 task 自动派生一个 tx-id(基于 application.id + topic-partition),无须手动设置。

2.4 Transaction Coordinator 与 __transaction_state

每个 transactional.id 都映射到一个 Transaction Coordinator(实际上是 __transaction_state 这个内部 Topic 的某个分区的 Leader Broker):

coordinator_partition = hash(transactional.id) % num_partitions(__transaction_state)
                                               ↑ 默认 50

__transaction_state 的属性:

  • 内部 Topic,默认 50 分区,副本 3,cleanup.policy=compact
  • 每个分区由一个 Broker 作为 Leader → 该 Broker 就是这一批 transactional.id 的 Coordinator
  • 存的是「每个 transactional.id 当前事务的状态」

事务状态机:

Empty → Ongoing → PrepareCommit → CompleteCommit → Empty
              ↘ PrepareAbort  → CompleteAbort  ↗

每次状态变更都追加一条记录到 __transaction_state,是事务日志(WAL)的角色。

2.5 两阶段提交流程图

关键点:

  • 阶段 1(PREPARE):Coordinator 把状态写成 PrepareCommit只要这条日志成功,事务就保证最终会被 commit——即便 Coordinator 崩溃,新 Coordinator 重放 __transaction_state 也会继续把 commit marker 写完。
  • 阶段 2(COMMIT/ABORT marker):Coordinator 给所有参与的分区 Leader 发 WriteTxnMarker RPC,每个 Leader 在该分区的日志末尾写一条控制消息(Control Batch),标记「这个事务在我这里结束了」。
  • 控制消息对消费者不可见(被 Consumer 客户端过滤),但是 Consumer 用它判断「事务的 LSO 推进到哪里了」。

2.6 事务消息的「物理布局」

写入磁盘后,分区日志大致长这样:

offset  data                                    isTransactional  PID Epoch Seq
   100  普通消息 X                              false             -   -    -
   101  事务消息 A (tx_id=tx1)                  true              42  8    0
   102  事务消息 B (tx1)                        true              42  8    1
   103  另一个事务的消息 P (tx_id=tx2)          true              17  3    0
   104  事务消息 C (tx1)                        true              42  8    2
   105  Control Batch: COMMIT (PID=42 Epoch=8) [ABORT/COMMIT marker]  ← 来自 Coordinator
   106  Control Batch: ABORT  (PID=17 Epoch=3)
   107  普通消息 Y                              false

事务消息和非事务消息可以混在同一分区!靠 Control Batch 把它们分组。Consumer 端的 read_committed 模式会读到 100 / 105 之后的「已 commit」消息(A/B/C),跳过 P。


3. Read Committed Consumer:消费侧的过滤

3.1 三种消费隔离级别

properties
isolation.level=read_uncommitted   # 默认:所有消息都读到(包括未 commit / 已 abort 的)
isolation.level=read_committed     # 只读到「已 commit」事务的消息 + 所有非事务消息

Java/Python confluent-kafka 客户端都支持这个参数。

3.2 LSO(Last Stable Offset)

普通分区有 HW(High Watermark):消费者最多能读到 HW 之前的消息。事务引入了第二个水位:LSO(Last Stable Offset)

定义:LSO = 最早一个仍处于「未完成事务」状态的 offset。也就是说:

  • LSO 之前的消息:所有事务都已经收到 commit / abort marker,结果稳定
  • LSO 及之后的消息:可能属于还未完成的事务,结果未定

read_committed 消费者只能读到 LSO 之前的内容。

分区 orders-0:
  offset   |  100 | 101 | 102 | 103 | 104 | 105 (COMMIT) | 106 | 107 |
            normal  txA   txA   txB    txA   marker(A)     txB   normal
                                                  ↑                   ↑
                                                  LSO (B还未结束)     HW

read_committed Consumer 能看到: 100, 101, 102, 104(A 的消息已 commit)
被跳过 (因为 B 还没完成):       103, 106, 107  ← 等 B commit 后再可见

直到 txB 也 commit/abort:

  • 如果 commit:103, 106 显示,107 显示
  • 如果 abort:103, 106 永远跳过,107 显示

📌 为什么这样设计?因为事务可能跨多个分区,每个分区独立写入,commit marker 是逐个写过去的。为了让消费者「看到的就是 commit 之后的稳定状态」,引入 LSO 把「未稳定」的部分挡在外面。

3.3 事务被 abort 时消费者怎么办?

Consumer 收到 ABORT marker 后,会把该事务已经消费但还没交给业务的消息直接丢弃

  • Fetcher 内部维护「未完成事务的 PID 集合」
  • 拉到一条 (PID=42, Epoch=8) 的消息时,先 buffer
  • 拉到 ABORT marker 时,把所有 buffer 里 PID=42 的消息丢掉
  • 拉到 COMMIT marker 时,把 buffer 里 PID=42 的消息交给应用

业务代码完全无感,跟普通消费一样写 for msg in poll()

3.4 「未提交事务卡住消费」的真实坑

如果一个事务 Producer 启动了事务,没有显式 commit/abort 就崩了,那个事务会处于 Ongoing 状态,LSO 就停在那里推进不了——所有 read_committed 消费者会卡住。

Kafka 用 transaction.timeout.ms(默认 60 秒)来兜底:超过这个时间还没 commit 的事务自动被 Coordinator abort。所以崩溃后的最差等待时间是 60 秒,过完消费者会继续。

如果你设了一个很大的 timeout(比如 1 小时),那卡 1 小时。


4. consume-process-produce 模式(Streams 的灵魂)

4.1 完整的端到端 EOS 链路

                         ┌──────────────────────┐
                         │  input topic (orders) │
                         └──────────┬───────────┘
                                    │ poll

                         ┌──────────────────────┐
                         │   App: process(msg)  │
                         └──────────┬───────────┘

                              一个事务里

              ┌─────────────────────┼─────────────────────────┐
              ▼                     ▼                          ▼
      output topic A         output topic B           __consumer_offsets
      (filtered)             (aggregated)             (input offset)
              ▲                     ▲                          ▲
              └─────────────────────┴─────── COMMIT / ABORT ───┘

整个事务里包含三件事

  1. output topic A
  2. output topic B
  3. 提交 input topic 的 offset

任何一件失败,整个事务都 abort,下次重启重做。下游 read_committed 消费者只看到 commit 后的稳定结果。

4.2 Python 完整模板

python
from confluent_kafka import Consumer, Producer

c = Consumer({
    "bootstrap.servers": "127.0.0.1:9092",
    "group.id": "stream-app-1",
    "enable.auto.commit": False,
    "isolation.level": "read_committed",
    "auto.offset.reset": "earliest",
})
p = Producer({
    "bootstrap.servers": "127.0.0.1:9092",
    "enable.idempotence": True,
    "transactional.id": "stream-app-1-tx",
    "transaction.timeout.ms": 60000,
})

p.init_transactions()
c.subscribe(["input"])

while True:
    msgs = c.consume(num_messages=100, timeout=1.0)
    msgs = [m for m in msgs if m and not m.error()]
    if not msgs: continue

    p.begin_transaction()
    try:
        for m in msgs:
            out = handle(m)
            p.produce("output", key=out.key, value=out.value)

        # input offsets 也作为事务一部分
        p.send_offsets_to_transaction(
            c.position(c.assignment()),
            c.consumer_group_metadata(),
        )
        p.commit_transaction()
    except KafkaException as ke:
        if ke.args[0].retriable():
            p.abort_transaction()
            continue
        elif ke.args[0].txn_requires_abort():
            p.abort_transaction()
            continue
        else:
            raise   # 致命:让进程退出

4.3 Kafka Streams 的 processing.guarantee

Kafka Streams 把上面这套模板封装成了一个开关:

properties
processing.guarantee=at_least_once         # 默认
processing.guarantee=exactly_once          # 旧版(≤ 2.4),每个 task 一个事务
processing.guarantee=exactly_once_v2       # 推荐(≥ 2.6),单事务覆盖多 task

v2 vs v1 的关键差异

  • v1:每个 task(StreamThread × Partition)一个 transactional.id,每个 task 自己开事务、commit 事务。一个 1000 分区的 topic 会产生 1000 个事务,commit 频繁,性能差。
  • v2:基于 Brokers 2.5 引入的「fetch from follower + producer 共享 transactional.id」改进,一个 StreamThread 上的所有 task 复用同一个 producer,事务 commit 次数大幅下降,吞吐接近 at_least_once。

📌 生产经验:能用 v2 就用 v2。v1 已经是「历史包袱」,2.6 之后的 Streams 默认 v2。


5. EOS 的性能代价

社区 benchmark 数据(Confluent 官方):

模式吞吐相对值延迟(p99)
Producer 不开幂等100%5ms
幂等 Producer~98%5ms
事务 Producer,1000 条/事务~93%8ms
事务 Producer,10 条/事务~70%25ms
EOS_v1 Streams~50%100ms
EOS_v2 Streams~95%12ms

结论

  • 幂等 Producer 几乎免费(生产环境强制打开
  • 事务 Producer 的开销主要在「commit 的延迟」,批量越大开销越摊薄(控制每个事务里至少 100~1000 条消息)
  • EOS_v2 让 Streams 的 EOS 几乎不亏吞吐,没理由不用

6. 适用与不适用场景

✅ 适合 Kafka EOS

  • Kafka → Kafka 的流处理链路(Streams / Flink Kafka Connector / 自研 consume-process-produce)
  • 多 Topic 写入需要原子(订单写 orders + 写 audit_log
  • 状态聚合(KTable changelog 的写入和源 topic 的 offset 提交需原子)

❌ 不适合 Kafka EOS

  • Kafka → MySQL / ClickHouse / Elasticsearch 的下游写入:DB 不在 Kafka 事务里,两阶段提交无法保证。这种场景做业务幂等(第 12 章)。
  • Kafka → 第三方接口(短信 / 支付):第三方接口幂等设计是必备,靠业务主键 + 唯一性。
  • Producer + 外部数据库的 XA 二阶段提交:Kafka 不支持作为 XA 的 RM 参与外部事务(设计哲学上拒绝)。

📌 一张图总结

                  哪种「Exactly Once」?

         ┌────────────────┴────────────────┐
         ▼                                 ▼
  「Kafka 内部链路」                 「Kafka → 外部系统」
         │                                 │
         ▼                                 ▼
  Kafka EOS 协议                    业务侧幂等
  (idempotent + tx + read_committed)  (DB 唯一约束 / Redis SETNX / 状态机)

         └─ Kafka Streams: processing.guarantee=exactly_once_v2 一键开启

7. 配套代码

  • idempotent_producer.py:开关 enable.idempotence 对比,模拟 ACK 丢失重发
  • transactional_producer.py:完整事务 Producer,多 Topic 一次提交
  • consume_process_produce.py:经典「读 → 处理 → 写 + 提交 offset」事务模式
  • read_committed_consumer.py:read_committed vs read_uncommitted 对照
  • inspect_transaction_state.py:消费 __transaction_state 看事务日志

8. 与其他 MQ 对比

维度KafkaRabbitMQRocketMQPulsar
Producer 幂等✅ PID + Epoch + Seq❌(依赖业务)✅(消息 ID 重复检测,Broker 端去重)✅(Producer Sequence + dedup window)
事务消息✅ 二阶段提交 + Coordinator + __transaction_state❌(仅 publisher confirms)✅ 半消息 + 事务回查(更适合本地事务联动)✅ 类 Kafka,跨 Topic 事务
消费隔离read_committed (LSO 过滤)顺序消费 + 事务标记类 Kafka
使用门槛中(要懂 transactional.id / fencing / send_offsets)
适用场景流处理 (consume-process-produce)RPC 风格队列业务事务(与本地 DB 联动)类 Kafka

📌 RocketMQ 的事务消息和 Kafka 思路不同:RocketMQ 让 Producer 先发「半消息」(消费者不可见),然后业务做本地事务,根据本地事务结果通知 Broker 提交 / 回滚;Coordinator 还会主动回查 Producer 询问状态。更贴近「业务本地事务联动」,但不适合 Kafka 这种「写多 Topic 原子」的流处理场景。


9. 小结

  1. 幂等 Producer:PID + Epoch + Sequence 三件套,单 Session 单分区去重。生产环境无脑打开。
  2. 事务 Producertransactional.id 让 PID 跨 Session 复用,配合 Epoch 实现 Fencing;通过 Transaction Coordinator + __transaction_state 做两阶段提交,把多 Topic 写入 + offset 提交绑成原子。
  3. Read Committed Consumerisolation.level=read_committed + LSO 实现「只读到稳定结果」。
  4. consume-process-produce 是 EOS 的核心使用模式,对应 Kafka Streams 的 exactly_once_v2
  5. EOS 不能延伸到外部系统:Kafka → MySQL 的端到端 EOS 还是要靠业务幂等。
  6. 性能代价:幂等几乎免费,事务(批量 ≥ 100 条)≤ 10% 损失,EOS_v2 让 Streams 接近 at_least_once 性能

10. 面试高频题

Q1. Kafka 的 Exactly Once 是怎么做到的?

考察点:EOS 的全貌。

标准答案(分点):

  1. 幂等 Producer:通过 PID + Epoch + Sequence 在 Broker 端去重,保证单 Producer 实例对单分区不重复。
  2. 事务 Producer:通过 transactional.id 跨 Session 复用 PID + Epoch +1,配合 Transaction Coordinator + __transaction_state 做两阶段提交,把多 Topic 写 + offset 提交绑成原子操作。
  3. Read Committed Consumerisolation.level=read_committed 配合 LSO 过滤掉未提交事务的消息。
  4. 三件套合在一起,覆盖 consume-process-produce 流处理链路的端到端 EOS。

加分项:明确说出适用范围——Kafka 内部链路(Streams 的世界),不能保证 Kafka → 外部系统的 EOS。

易错点:把 enable.idempotence=true 当 EOS。它只是 EOS 的第一块拼图。


Q2. PID + Epoch + Sequence 是怎么去重的?

考察点:幂等 Producer 协议。

标准答案

  • Producer 启动时向 Broker 申请一个 64-bit 的 PID
  • 同一个 PID 可能被「重启的同名实例」抢用,Epoch 用于区分先后,新 Epoch 顶替旧的。
  • 每个分区单独维护一个递增 Sequence,每条消息附带 (PID, Partition, Sequence)
  • Broker 内存里维护「(PID, Partition) → 最近见过的 Sequence」:
    • 收到 seq = lastSeq + 1:写入 + 更新
    • 收到 seq ≤ lastSeq:丢弃,但回 ACK(让 Producer 当作成功)
    • 收到 seq > lastSeq + 1:返回 OUT_OF_ORDER 让 Producer 重发
  • 最多缓存 5 个 batch(受 max.in.flight.requests.per.connection ≤ 5 限制)。

Q3. transactional.id 的作用是什么?为什么必须全集群唯一?

考察点:事务 Producer 的命名机制 + Fencing。

标准答案

  • transactional.id 是用户指定的逻辑 Producer 标识(比如 tx-order-1),在 Producer 重启后用来复用 PID 并把 Epoch + 1。
  • 旧实例的请求带着小 Epoch 发给 Broker 时,Broker 比对当前 Epoch 后会拒绝并返回 INVALID_PRODUCER_EPOCH,触发 ProducerFencedException,老实例必须立即退出。
  • 这就是 Producer Fencing——保证「不会有两个同名 Producer 同时在写」,避免脑裂。
  • 必须全集群唯一是因为它是 Coordinator 的查找键,否则多个 Producer 互相 fence 会乱成一锅粥。

Q4. Read Committed 是怎么实现的?LSO 是什么?

考察点:消费侧过滤机制。

标准答案

  • isolation.level=read_committed 让消费者只能读到「已 commit 的事务消息」+「所有非事务消息」。
  • 实现关键是 LSO(Last Stable Offset):分区里最早一个仍处于 Ongoing 状态的事务起点 offset。Consumer 最多读到 LSO 之前的内容。
  • Fetcher 客户端拉到事务消息时先 buffer,看到 COMMIT marker 才交给业务,看到 ABORT marker 直接丢弃。
  • 这样跨分区的事务也能保证「要么所有分区都看到,要么都看不到」。

Q5. 事务两阶段提交的流程是怎样的?Coordinator 挂了怎么办?

考察点:分布式事务协议细节。

标准答案

  1. Producer 调 initTransactions(),Coordinator 分配 PID/Epoch 并写 __transaction_state
  2. beginTransaction() 后每次 produce 都通过 AddPartitionsToTxn 通知 Coordinator「这个事务用到了 (topic, partition)」。
  3. commitTransaction() 触发:
    • 阶段 1:Coordinator 把状态写成 PrepareCommit__transaction_state这条日志成功后事务就保证最终被 commit,是关键的 commit point)。
    • 阶段 2:Coordinator 给每个参与分区的 Leader 发 WriteTxnMarker RPC,写一条 Control Batch(COMMIT marker)。
    • 全部成功后状态写成 CompleteCommit
  4. Coordinator 挂了:新 Coordinator 接管该 __transaction_state 分区时会重放日志,看到 PrepareCommit 就继续执行阶段 2;看到 Ongoing 且超时(默认 60s)就自动 abort。

加分项:说出 Coordinator 实际就是 __transaction_state 某分区 Leader 所在 Broker,故障转移就是该分区 Leader 切换。


Q6. processing.guarantee=exactly_once_v2 比 v1 强在哪?

考察点:Streams EOS 的演进。

标准答案

  • v1:每个 StreamTask 一个 transactional.id,1000 个分区就 1000 个事务,commit RPC 大量增加,吞吐打折严重(约 50%)。
  • v2(Kafka 2.5 引入,2.6+ 默认):基于 Producer 协议的改进——一个 StreamThread 上的所有 task 共用一个 Producer + 一个 transactional.id,事务 commit 次数从 N(任务数)降到 1(线程数)。
  • 性能从 ~50% at_least_once 提升到 ~95%,几乎免费。
  • 配套的 broker 改进:fetch_request 携带 transactional 元数据,让 producer 复用更安全。

Q7. EOS 能保证「Kafka → MySQL」端到端不重复吗?为什么?

考察点:EOS 的边界与适用范围。

标准答案

  • 不能。Kafka 的事务只能覆盖:
    • Kafka 内部多 Topic 的写入
    • 消费者 offset 的提交(写到 __consumer_offsets
    • 这些都是 Kafka 自己的「事务参与方」
  • MySQL 不是 Kafka 事务的参与方。「写 MySQL → 提交 Kafka 事务」中间任何故障都会导致状态不一致。
  • 正确做法:业务侧幂等(业务主键 + DB UNIQUE / Redis SETNX / 状态机)。先 At Least Once 拉消息,业务侧靠唯一约束保证最终一次。
  • 进阶方案:Outbox Pattern——业务写 MySQL 时同时往 outbox 表写一条事件,Debezium CDC 把 outbox 变化发到 Kafka,整个链路都靠 DB 事务 + Kafka 事务。

🎯 小练习:用 transactional_producer.py 启动两个相同 transactional.id 的实例,看第二个启动时第一个会抛 ProducerFencedException 立刻挂掉,这就是 Fencing 的真实效果。


11. 附录:常见错误码速查

理解事务 Producer 的错误处理对生产环境至关重要。Kafka 把事务相关错误分成几类:

11.1 可重试错误(retriable)

ke.args[0].retriable() == True:网络抖动、Coordinator 临时不可达。 策略:直接 continue 重试 begin/commit。

错误码含义
COORDINATOR_NOT_AVAILABLETransaction Coordinator 还没选出来
COORDINATOR_LOAD_IN_PROGRESSCoordinator 正在加载状态
NOT_COORDINATOR该 Broker 不是这个 tx_id 的 Coordinator
REQUEST_TIMED_OUT网络超时
CONCURRENT_TRANSACTIONS同 tx_id 的事务正在收尾,稍等再 begin

11.2 需要 abort 的错误(txn_requires_abort)

ke.args[0].txn_requires_abort() == True:当前事务搞砸了,必须 abort 后才能 begin 下一个。 策略:捕获后调 abort_transaction() 然后 continue

错误码含义
INVALID_TXN_STATEProducer 状态机错乱(比如重复 begin)
TRANSACTION_TIMED_OUT事务超过 transaction.timeout.ms 被 Coordinator 强制 abort
OUT_OF_ORDER_SEQUENCE_NUMBER (在事务内)幂等 seq 乱了
UNKNOWN_PRODUCER_IDBroker 端 PID 状态被清理(长期未活动 + 超时回收)

11.3 致命错误(fatal)

ke.args[0].is_fatal() == True:必须关掉 Producer + 退出进程,重启后重新 init。

错误码含义
INVALID_PRODUCER_EPOCH被新实例 fence 掉了(PID Fencing)
PRODUCER_FENCED同上的别名
INVALID_TRANSACTION_TIMEOUTtimeout 配置非法
TRANSACTIONAL_ID_AUTHORIZATION_FAILEDACL 权限不足
CLUSTER_AUTHORIZATION_FAILED没有 IdempotentWrite 权限

11.4 标准错误处理模板

python
from confluent_kafka import KafkaException, KafkaError

try:
    p.commit_transaction()
except KafkaException as ke:
    err = ke.args[0]
    if err.is_fatal():
        log.error(f"FATAL: {err}, exiting"); raise
    elif err.txn_requires_abort():
        log.warning(f"abort: {err}")
        p.abort_transaction()
    elif err.retriable():
        log.info(f"retry: {err}")
        # 通常什么都不用做,下个循环 begin_transaction 就行
    else:
        log.error(f"unknown error class: {err}"); raise

12. 附录:性能调优清单

12.1 Producer 端

  • 批大小:每个事务里至少 100~1000 条消息,让 commit RTT 摊薄
  • linger.ms=20~100:换批量化,吞吐有明显提升
  • compression.type=lz4 / zstd:事务也支持压缩
  • max.in.flight.requests.per.connection=5:保持默认,不要降到 1

12.2 Coordinator / Broker 端

  • __transaction_state 副本:生产至少 3 副本 + min.insync.replicas=2
  • transaction.state.log.num.partitions:默认 50,超大集群可以加到 100
  • transaction.max.timeout.ms:默认 15min,限制单个事务最大 timeout,防止滥用
  • transaction.remove.expired.transaction.cleanup.interval.ms:默认 1h,过期事务清理频率

12.3 Consumer 端

  • isolation.level=read_committed:必加(否则等于没用 EOS)
  • fetch.min.bytes / fetch.max.wait.ms:read_committed 模式下 broker 要等事务收尾,可以适当调大让 batch 更高效
  • 不要把超长事务的下游 Consumer 跟普通 Consumer 共用 group:read_committed 下游会被未收尾事务阻塞

13. 附录:与 RocketMQ 事务消息的深度对比

很多面试官会顺手问「Kafka 事务和 RocketMQ 事务区别」,这里专门展开:

13.1 解决的问题不同

  • Kafka 事务:解决「Kafka 内部多 Topic 写入 + offset 提交的原子性」(流处理场景)
  • RocketMQ 事务消息:解决「Kafka 写入 + 本地 DB 事务的原子性」(业务事务联动场景)

13.2 协议差异

阶段KafkaRocketMQ
事务开始beginTransaction()sendMessageInTransaction(Half message)
业务执行produce(...) 多次本地 DB 事务(应用控制)
事务提交commitTransaction() → 二阶段写 markerLocalTransactionState.COMMIT_MESSAGE
失败回滚abortTransaction() → 二阶段写 ABORT markerLocalTransactionState.ROLLBACK_MESSAGE
不确定状态自动超时 abort(60s)Broker 主动回查 Producer 询问状态
消费者过滤isolation.level=read_committedHalf message 对消费者不可见,commit 后变可见

13.3 「主动回查」是 RocketMQ 的杀手锏

RocketMQ 的 Half Message 提交后,如果 Producer 长时间不告诉 Broker「commit 还是 rollback」(因为 Producer 也可能崩了),Broker 会主动调用 Producer 的 checkLocalTransactionState() 接口询问该事务最终结果。

这在「Kafka + 本地 DB」场景里非常好用:

java
// RocketMQ 风格
public class TxListener implements TransactionListener {
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            db.insert(...);     // 本地 DB 事务
            return COMMIT_MESSAGE;
        } catch (Exception e) {
            return ROLLBACK_MESSAGE;
        }
    }
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        // Broker 回查时,看 DB 里这条记录到底有没有
        return db.exists(msg.getKey()) ? COMMIT_MESSAGE : ROLLBACK_MESSAGE;
    }
}

Kafka 没有这个机制——它假定「Producer 自己知道是 commit 还是 abort」,崩溃就只能等超时强制 abort。所以 Kafka 不太适合做「业务 DB 事务联动」,更适合「纯 Kafka 内部流处理」。

13.4 选型建议

场景推荐
Streams 流处理(consume-process-produce)Kafka EOS
业务下单 + 发消息(DB 事务联动)RocketMQ 事务消息 / Outbox Pattern
跨多个 Kafka Topic 原子写入Kafka EOS
Kafka → MySQL 端到端不丢不重业务幂等(DB 唯一约束)

14. 附录:完整 EOS 配置 cheat sheet

把分散在前面章节的关键配置整理一下,照抄即可:

properties
# ===== Producer =====
bootstrap.servers=127.0.0.1:9092
enable.idempotence=true
acks=all
transactional.id=stream-app-1-tx     # 必填,集群唯一,跨重启不变
transaction.timeout.ms=60000
max.in.flight.requests.per.connection=5
retries=2147483647
linger.ms=20
compression.type=lz4

# ===== Consumer =====
bootstrap.servers=127.0.0.1:9092
group.id=stream-app-1-group
enable.auto.commit=false              # 必关
isolation.level=read_committed        # 必开
auto.offset.reset=earliest

# ===== Broker (集群级) =====
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
offsets.topic.replication.factor=3
offsets.topic.min.isr=2
transaction.state.log.num.partitions=50
transaction.max.timeout.ms=900000     # 限制单事务最大 timeout (15min)

这套配置对绝大多数 Streams 应用都直接可用。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
#!/usr/bin/env python3
"""
consume_process_produce.py
==========================
Streams 灵魂模式:从 input topic 读 → 处理 → 写 output topic + 提交 input offset,
全部在一个事务里。这是 Kafka EOS 协议层最重要的使用场景。

用法:
  pip install confluent-kafka

  # 1) 灌一些原始数据
  python consume_process_produce.py seed 200

  # 2) 启动事务式 stream app(持续运行,Ctrl+C 退出)
  python consume_process_produce.py run

  # 3) 用 read_committed_consumer.py 看 output topic
  python read_committed_consumer.py learn.13.cpp_out committed
"""

import os
import sys
import time
import json

from confluent_kafka import Producer, Consumer, KafkaException, TopicPartition
from confluent_kafka.admin import AdminClient, NewTopic

BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
T_IN = "learn.13.cpp_in"
T_OUT = "learn.13.cpp_out"
TX_ID = "tx-cpp-stream-1"


def ensure():
    a = AdminClient({"bootstrap.servers": BOOTSTRAP})
    existing = a.list_topics(timeout=5).topics
    new = [NewTopic(t, num_partitions=3, replication_factor=1)
           for t in (T_IN, T_OUT) if t not in existing]
    if new:
        a.create_topics(new); time.sleep(1)


def seed(n):
    ensure()
    p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 5})
    for i in range(n):
        p.produce(T_IN, key=str(i % 5), value=json.dumps({"raw": i}).encode())
    p.flush(10)
    print(f"seeded {n} raw events to {T_IN}")


def run():
    ensure()
    consumer = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "cpp-stream-group",
        "enable.auto.commit": False,
        "isolation.level": "read_committed",
        "auto.offset.reset": "earliest",
    })
    consumer.subscribe([T_IN])

    producer = Producer({
        "bootstrap.servers": BOOTSTRAP,
        "enable.idempotence": True,
        "transactional.id": TX_ID,
        "transaction.timeout.ms": 60000,
        "linger.ms": 5,
    })
    producer.init_transactions()

    print(f"running consume-process-produce: {T_IN}{T_OUT}, tx_id={TX_ID}")
    try:
        while True:
            msgs = consumer.consume(num_messages=50, timeout=1.0)
            msgs = [m for m in msgs if m and not m.error()]
            if not msgs:
                continue

            producer.begin_transaction()
            try:
                for m in msgs:
                    raw = json.loads(m.value())
                    out = {"clean": raw["raw"] * 2, "src_off": m.offset(), "src_p": m.partition()}
                    producer.produce(T_OUT, key=m.key(), value=json.dumps(out).encode())

                # 把 input offsets 包进同一个事务
                producer.send_offsets_to_transaction(
                    consumer.position(consumer.assignment()),
                    consumer.consumer_group_metadata(),
                )
                producer.commit_transaction()
                print(f"✅ tx commit: processed {len(msgs)} msgs (offsets推进 input topic)")
            except KafkaException as e:
                err = e.args[0]
                if err.txn_requires_abort():
                    producer.abort_transaction()
                    print(f"⚠️ tx aborted: {err}")
                    continue
                else:
                    raise
    except KeyboardInterrupt:
        pass
    finally:
        consumer.close()


if __name__ == "__main__":
    if len(sys.argv) < 2:
        print(__doc__); sys.exit(0)
    if sys.argv[1] == "seed":
        seed(int(sys.argv[2]) if len(sys.argv) > 2 else 100)
    elif sys.argv[1] == "run":
        run()
    else:
        print(__doc__)
python
#!/usr/bin/env python3
"""
idempotent_producer.py
======================
对比开 / 关 enable.idempotence 两种 Producer 在「重发」时的差别。

为了模拟「ACK 丢失 → Producer 重发」的场景:
  - 我们故意把同一条 (key, value) 在循环里 send 两次
  - 不开幂等:Broker 老老实实写两份;下游 consumer 看到两条
  - 开幂等:第二条因 (PID, partition, seq) 重复被 Broker 静默去重

⚠️ 注意:这里的「重发」是模拟。真实的 librdkafka 自动重试发生在 Producer 内部,
对应用层透明,看不到第二次 send 调用——但物理上 Broker 收到了两次相同的 batch。
本脚本通过显式调用两次 produce 来制造同样的物理效果。

用法:
  pip install confluent-kafka
  python idempotent_producer.py off    # 不开幂等,下游会读到 2N 条
  python idempotent_producer.py on     # 开幂等,下游只有 N 条 (但实际只对真实重试有效,
                                       # 应用层显式 produce 仍然算两条不同消息;
                                       # 见下方说明)
"""

import os
import sys
import time
import json

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.13.idem"
N = 10


def ensure_topic():
    a = AdminClient({"bootstrap.servers": BOOTSTRAP})
    if TOPIC not in a.list_topics(timeout=5).topics:
        a.create_topics([NewTopic(TOPIC, num_partitions=1, replication_factor=1)])
        time.sleep(1)


def run(idempotent: bool):
    ensure_topic()
    cfg = {
        "bootstrap.servers": BOOTSTRAP,
        "enable.idempotence": idempotent,
        "acks": "all",
        # 故意调小 in-flight 让重传场景更易复现
        "max.in.flight.requests.per.connection": 5,
        # 触发 librdkafka 自动重试:把 message.send.max.retries 调大
        "retries": 5,
        "linger.ms": 5,
    }
    p = Producer(cfg)
    print(f"=== Producer enable.idempotence={idempotent} ===")

    # 强制制造「Broker 物理上收到两次相同 batch」的效果:
    # 我们直接发同一个 key/value 两次。如果开了幂等,librdkafka 会复用同一个
    # PID + Sequence(因为是同一次 send 的 retry 才会复用 seq;显式两次 send
    # 仍是两条不同消息,seq 是 N 和 N+1,Broker 都会收下)。
    # 因此真实的去重场景请用 inject_dup_via_socket() 模拟(涉及到 Producer 内部
    # 状态较复杂),这里我们直接打印 Producer 看到的元数据,让读者直观体会
    # 「PID/Epoch 是否被分配」。
    def cb(err, msg):
        if err:
            print(f"  ❌ delivery err: {err}")
        else:
            print(f"  ✅ delivered offset={msg.offset()}")

    for i in range(N):
        body = json.dumps({"i": i, "ts": time.time()})
        p.produce(TOPIC, key=str(i), value=body.encode(), callback=cb)
    p.flush(10)

    # 打印实际生效的 PID(confluent-kafka 不直接暴露,但可以从 log/统计 推断)
    print()
    print("→ 注意:开了 enable.idempotence 后,librdkafka 内部重试会带相同 (PID,seq),")
    print("       Broker 通过 (PID, partition, lastSeq) 比对自动去重;")
    print("       关闭幂等时网络抖动重试会让 Broker 写入两份相同消息(外人无感知)。")
    print(f"→ 已发送 {N} 条到 {TOPIC}")
    print(f"→ 用如下命令统计实际写入条数:")
    print(f"   python idempotent_producer.py count")


def count():
    c = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "count-" + str(os.getpid()),
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
    })
    c.subscribe([TOPIC])
    n = 0
    deadline = time.time() + 5
    while time.time() < deadline:
        msg = c.poll(0.5)
        if msg and not msg.error(): n += 1
    c.close()
    print(f"读到 {n} 条")


if __name__ == "__main__":
    if len(sys.argv) < 2:
        print(__doc__); sys.exit(0)
    if sys.argv[1] == "off":
        run(False)
    elif sys.argv[1] == "on":
        run(True)
    elif sys.argv[1] == "count":
        count()
    else:
        print(__doc__)
python
#!/usr/bin/env python3
"""
inspect_transaction_state.py
============================
直接消费内部 Topic `__transaction_state`,观察事务的状态变迁日志。

每个事务从 Empty → Ongoing → PrepareCommit/PrepareAbort →
CompleteCommit/CompleteAbort 的过程,都会作为一条消息追加到 `__transaction_state`,
Key 是 transactional.id,Value 是当前的事务元数据(PID, Epoch, state, partitions, ...)。

由于 Value 是 Kafka 内部的二进制 schema,比 __consumer_offsets 还复杂(带 partition
list、多个版本),完整解析超出脚本范围;这里我们只解析 Key(transactional.id),
并用启发式打印 Value 大小、是否 tombstone,让你看到「事务状态变化触发的写入流」。

用法:
  pip install confluent-kafka
  python inspect_transaction_state.py

  # 然后另一个终端跑 transactional_producer.py commit/abort,
  # 这里会实时打印每个事务状态变更
"""

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("txstate")


def parse_key(key):
    """
    TransactionLogKey schema:
      int16  version
      string transactional_id
    """
    if len(key) < 2: return None
    version = struct.unpack(">h", key[:2])[0]
    pos = 2
    if pos + 2 > len(key): return None
    n = struct.unpack(">h", key[pos:pos+2])[0]; pos += 2
    if n < 0: return {"version": version, "tx_id": None}
    tx_id = key[pos:pos+n].decode("utf-8", errors="replace")
    return {"version": version, "tx_id": tx_id}


def heuristic_state(value):
    """
    Value 的 schema(v0 简化):
      int16  version
      int64  producer_id
      int16  producer_epoch
      int32  txn_timeout_ms
      int8   txn_state    (0=Empty, 1=Ongoing, 2=PrepCommit, 3=PrepAbort,
                           4=CompleteCommit, 5=CompleteAbort, 6=Dead, 7=PrepareEpochFence)
      int32  txn_start_ts (some versions)
      ...    partitions array (复杂)
    """
    if value is None:
        return {"tombstone": True}
    if len(value) < 15: return {"len": len(value), "parse": "too short"}
    try:
        version = struct.unpack(">h", value[0:2])[0]
        pid = struct.unpack(">q", value[2:10])[0]
        epoch = struct.unpack(">h", value[10:12])[0]
        timeout = struct.unpack(">i", value[12:16])[0]
        state = value[16] if len(value) > 16 else None
        STATE_NAMES = {0: "Empty", 1: "Ongoing", 2: "PrepareCommit", 3: "PrepareAbort",
                       4: "CompleteCommit", 5: "CompleteAbort", 6: "Dead",
                       7: "PrepareEpochFence"}
        return {
            "version": version, "PID": pid, "Epoch": epoch,
            "timeout_ms": timeout, "state": STATE_NAMES.get(state, f"unknown({state})"),
            "value_size": len(value),
        }
    except Exception as e:
        return {"len": len(value), "err": str(e)}


def main():
    c = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "inspect-txstate-" + str(os.getpid()),
        "enable.auto.commit": False,
        "auto.offset.reset": "latest",
    })
    c.subscribe(["__transaction_state"])
    log.info("listening to __transaction_state ...")
    log.info("(在另一个终端运行 transactional_producer.py 即可看到事务状态变迁)")
    try:
        while True:
            msg = c.poll(1.0)
            if msg is None: continue
            if msg.error():
                log.error(msg.error()); continue
            k = parse_key(msg.key() or b"")
            v = heuristic_state(msg.value())
            tx_id = k["tx_id"] if k else "?"
            if v.get("tombstone"):
                log.info(f"🗑️  tx_id={tx_id!r} → TOMBSTONE (Coordinator 清理记录)")
            elif "err" in v:
                log.info(f"📝 tx_id={tx_id!r} → value parse error: {v['err']} (size={v['len']})")
            else:
                log.info(
                    f"📝 tx_id={tx_id!r}  PID={v['PID']}  Epoch={v['Epoch']}  "
                    f"state={v['state']}  timeout={v['timeout_ms']}ms  "
                    f"size={v['value_size']}B"
                )
    except KeyboardInterrupt:
        pass
    finally:
        c.close()


if __name__ == "__main__":
    main()
python
#!/usr/bin/env python3
"""
read_committed_consumer.py
==========================
对比 read_uncommitted vs read_committed 的实际效果。

用法:
  pip install confluent-kafka

  # 终端 A:跑 transactional_producer.py commit/abort 演示
  # 终端 B:用 uncommitted 模式看,会立刻看到所有事务消息(包括 abort 掉的)
  python read_committed_consumer.py learn.13.tx_orders uncommitted

  # 终端 C:用 committed 模式看,只能看到 commit 后的,abort 的永远看不到
  python read_committed_consumer.py learn.13.tx_orders committed
"""

import os
import sys

from confluent_kafka import Consumer

BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")


def main():
    if len(sys.argv) < 3:
        print(__doc__); sys.exit(0)
    topic = sys.argv[1]
    level = sys.argv[2]   # uncommitted | committed

    if level not in ("uncommitted", "committed"):
        print("level must be 'uncommitted' or 'committed'"); sys.exit(1)

    iso = "read_uncommitted" if level == "uncommitted" else "read_committed"

    c = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": f"reader-{level}-{os.getpid()}",
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
        "isolation.level": iso,
    })
    c.subscribe([topic])

    print(f"== isolation.level={iso} → topic={topic} ==")
    print("(事务正在进行中的消息:")
    print("  read_uncommitted 会立刻看到(包括 abort 后还展示);")
    print("  read_committed 只在 commit 后才看到,abort 的永远看不到)")
    print()

    n = 0
    try:
        while True:
            msg = c.poll(1.0)
            if msg is None: continue
            if msg.error():
                print(f"err: {msg.error()}"); continue
            n += 1
            print(f"  [{n:4d}] p{msg.partition()}@{msg.offset()} key={msg.key()!r} "
                  f"value={msg.value()[:80]!r}")
    except KeyboardInterrupt:
        print(f"-- received {n} msgs --")
    finally:
        c.close()


if __name__ == "__main__":
    main()
python
#!/usr/bin/env python3
"""
transactional_producer.py
=========================
完整的事务 Producer 演示:往两个 Topic 写消息,成功则 commit,失败则 abort。
观察 Read Committed 消费者的行为:commit 之前看不到,commit 之后才出现;abort
的消息永远看不到。

用法:
  pip install confluent-kafka
  python transactional_producer.py commit   # 一个事务,最后 commit
  python transactional_producer.py abort    # 一个事务,写一半 abort
  python transactional_producer.py loop     # 100 个事务循环跑(看吞吐)

  # 同时另一个终端跑 read_committed_consumer.py 看效果

⚠️ 启动两个相同 transactional.id 的实例会触发 ProducerFencedException —— 这是
PID Fencing 的真实效果。
"""

import os
import sys
import time
import json

from confluent_kafka import Producer, KafkaException, KafkaError
from confluent_kafka.admin import AdminClient, NewTopic

BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
TX_ID = os.environ.get("TX_ID", "tx-order-1")
T1 = "learn.13.tx_orders"
T2 = "learn.13.tx_outbox"


def ensure_topics():
    a = AdminClient({"bootstrap.servers": BOOTSTRAP})
    existing = a.list_topics(timeout=5).topics
    new_topics = [NewTopic(t, num_partitions=3, replication_factor=1)
                  for t in (T1, T2) if t not in existing]
    if new_topics:
        a.create_topics(new_topics)
        time.sleep(1)


def make_tx_producer():
    p = Producer({
        "bootstrap.servers": BOOTSTRAP,
        "enable.idempotence": True,        # 事务必须开幂等
        "acks": "all",
        "transactional.id": TX_ID,
        "transaction.timeout.ms": 60000,
        "linger.ms": 5,
    })
    p.init_transactions()                  # 一次性初始化,向 Coordinator 申请 PID/Epoch
    return p


def commit_demo():
    ensure_topics()
    p = make_tx_producer()
    p.begin_transaction()
    print(f"== commit demo == tx_id={TX_ID}")
    try:
        for i in range(5):
            payload = json.dumps({"order_id": f"O-{i}", "amount": 99 + i, "ts": time.time()})
            p.produce(T1, key=f"O-{i}", value=payload.encode())
            p.produce(T2, key=f"O-{i}", value=("audit-" + payload).encode())
            print(f"  → produced {i} to {T1} & {T2}  (still pending, read_committed 看不到)")
            time.sleep(0.3)
        p.commit_transaction()
        print("✅ commit_transaction() done. read_committed Consumer 现在能看到这 10 条")
    except KafkaException as e:
        print(f"❌ tx error: {e} → abort")
        p.abort_transaction()


def abort_demo():
    ensure_topics()
    p = make_tx_producer()
    p.begin_transaction()
    print(f"== abort demo == tx_id={TX_ID}")
    try:
        for i in range(3):
            payload = json.dumps({"order_id": f"BAD-{i}"})
            p.produce(T1, key=f"BAD-{i}", value=payload.encode())
            print(f"  → produced BAD-{i}(pending, 即将被 abort)")
        raise RuntimeError("simulated business failure!!!")
    except RuntimeError as e:
        print(f"💥 simulated: {e}")
        p.abort_transaction()
        print("✅ abort_transaction() done. read_committed Consumer 永远看不到这 3 条")


def loop_demo():
    ensure_topics()
    p = make_tx_producer()
    print(f"== loop demo == tx_id={TX_ID}, 100 个事务,每个 100 条")
    t0 = time.time()
    for tx in range(100):
        p.begin_transaction()
        try:
            for i in range(100):
                p.produce(T1, key=f"L-{tx}-{i}", value=f"v{tx}-{i}".encode())
            p.commit_transaction()
        except KafkaException as e:
            err = e.args[0]
            if err.txn_requires_abort():
                p.abort_transaction()
            else:
                raise
        if tx % 10 == 0:
            print(f"  done {tx + 1} txns")
    print(f"== finished in {time.time() - t0:.2f}s ==")


if __name__ == "__main__":
    if len(sys.argv) < 2:
        print(__doc__); sys.exit(0)
    {"commit": commit_demo, "abort": abort_demo, "loop": loop_demo}[sys.argv[1]]()

consume_process_produce.py ↗ · idempotent_producer.py ↗ · inspect_transaction_state.py ↗ · read_committed_consumer.py ↗ · transactional_producer.py ↗