主题
第 13 章 幂等与事务:Exactly Once 的真相
目标读者:听过「Kafka 支持 Exactly Once」却以为「打开一个开关就行」、不知道幂等 Producer 的边界、对事务和 Read Committed 感觉云里雾里的同学。
学完你会:能用一张图讲清 PID + Epoch + Sequence 的去重逻辑;知道幂等的「单 Session、单分区」边界为什么必须存在;能写完整的事务 Producer +
sendOffsetsToTransactionconsume-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 实例做了三件事:
- Producer ID(PID):每个 Producer 实例向 Broker 申请一个全集群唯一的 64-bit ID(
InitProducerIdRequest),打开 Producer 时分配。 - Producer Epoch(int16):同一个 PID 可能被「重启的同名实例」抢用,Epoch 用来区分先后;后启动的实例 Epoch + 1,旧的会被 Broker 拒绝(详见 1.5 节 Fencing)。
- 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三件事必须原子地发生:
- 写入 output topic 的所有消息(可能跨多个 Partition、多个 Topic)
- 把 input topic 的 offset 提交到
__consumer_offsets - 要么全成功,要么全失败
如果只是幂等:
- 写 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() 时:
- Producer 把
transactional.id发给 Transaction Coordinator - Coordinator 查
__transaction_state:- 如果该
transactional.id没记录 → 分配新 PID + Epoch=0 - 如果有记录 → 复用同一个 PID,把 Epoch +1
- 如果该
- 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 发
WriteTxnMarkerRPC,每个 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 ───┘整个事务里包含三件事:
- 写
output topic A - 写
output topic B - 提交
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),单事务覆盖多 taskv2 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 对比
| 维度 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 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. 小结
- 幂等 Producer:PID + Epoch + Sequence 三件套,单 Session 单分区去重。生产环境无脑打开。
- 事务 Producer:
transactional.id让 PID 跨 Session 复用,配合 Epoch 实现 Fencing;通过 Transaction Coordinator +__transaction_state做两阶段提交,把多 Topic 写入 + offset 提交绑成原子。 - Read Committed Consumer:
isolation.level=read_committed+ LSO 实现「只读到稳定结果」。 - consume-process-produce 是 EOS 的核心使用模式,对应 Kafka Streams 的
exactly_once_v2。 - EOS 不能延伸到外部系统:Kafka → MySQL 的端到端 EOS 还是要靠业务幂等。
- 性能代价:幂等几乎免费,事务(批量 ≥ 100 条)≤ 10% 损失,EOS_v2 让 Streams 接近 at_least_once 性能。
10. 面试高频题
Q1. Kafka 的 Exactly Once 是怎么做到的?
考察点:EOS 的全貌。
标准答案(分点):
- 幂等 Producer:通过 PID + Epoch + Sequence 在 Broker 端去重,保证单 Producer 实例对单分区不重复。
- 事务 Producer:通过
transactional.id跨 Session 复用 PID + Epoch +1,配合 Transaction Coordinator +__transaction_state做两阶段提交,把多 Topic 写 + offset 提交绑成原子操作。 - Read Committed Consumer:
isolation.level=read_committed配合 LSO 过滤掉未提交事务的消息。 - 三件套合在一起,覆盖 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 挂了怎么办?
考察点:分布式事务协议细节。
标准答案:
- Producer 调
initTransactions(),Coordinator 分配 PID/Epoch 并写__transaction_state。 beginTransaction()后每次produce都通过AddPartitionsToTxn通知 Coordinator「这个事务用到了 (topic, partition)」。commitTransaction()触发:- 阶段 1:Coordinator 把状态写成
PrepareCommit到__transaction_state(这条日志成功后事务就保证最终被 commit,是关键的 commit point)。 - 阶段 2:Coordinator 给每个参与分区的 Leader 发
WriteTxnMarkerRPC,写一条 Control Batch(COMMIT marker)。 - 全部成功后状态写成
CompleteCommit。
- 阶段 1:Coordinator 把状态写成
- 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_AVAILABLE | Transaction Coordinator 还没选出来 |
COORDINATOR_LOAD_IN_PROGRESS | Coordinator 正在加载状态 |
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_STATE | Producer 状态机错乱(比如重复 begin) |
TRANSACTION_TIMED_OUT | 事务超过 transaction.timeout.ms 被 Coordinator 强制 abort |
OUT_OF_ORDER_SEQUENCE_NUMBER (在事务内) | 幂等 seq 乱了 |
UNKNOWN_PRODUCER_ID | Broker 端 PID 状态被清理(长期未活动 + 超时回收) |
11.3 致命错误(fatal)
ke.args[0].is_fatal() == True:必须关掉 Producer + 退出进程,重启后重新 init。
| 错误码 | 含义 |
|---|---|
INVALID_PRODUCER_EPOCH | 被新实例 fence 掉了(PID Fencing) |
PRODUCER_FENCED | 同上的别名 |
INVALID_TRANSACTION_TIMEOUT | timeout 配置非法 |
TRANSACTIONAL_ID_AUTHORIZATION_FAILED | ACL 权限不足 |
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}"); raise12. 附录:性能调优清单
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=2transaction.state.log.num.partitions:默认 50,超大集群可以加到 100transaction.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 协议差异
| 阶段 | Kafka | RocketMQ |
|---|---|---|
| 事务开始 | beginTransaction() | sendMessageInTransaction(Half message) |
| 业务执行 | produce(...) 多次 | 本地 DB 事务(应用控制) |
| 事务提交 | commitTransaction() → 二阶段写 marker | LocalTransactionState.COMMIT_MESSAGE |
| 失败回滚 | abortTransaction() → 二阶段写 ABORT marker | LocalTransactionState.ROLLBACK_MESSAGE |
| 不确定状态 | 自动超时 abort(60s) | Broker 主动回查 Producer 询问状态 |
| 消费者过滤 | isolation.level=read_committed | Half 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 ↗