第 13 章 · 幂等与事务(EOS)可视化

PID + Epoch + Seq · 事务两阶段提交 · Read Committed · PID Fencing · consume-process-produce

幂等 Producer 重传去重:相同 (PID, Partition, Sequence) 被 Broker 拒绝

每次「发送」都会向 Broker 发一条带 (PID=42, Partition=0, Sequence=N) 的消息;点击「重发」按钮再发同一条 seq,会被 Broker 识别为重复并丢弃,但 ACK 仍返回。点击「乱序发送」模拟 seq 跳号,Broker 返回 OUT_OF_ORDER。

Producer (PID=42, Epoch=0)

就绪

Broker (orders-0)

分区日志:
已写入 被拒绝(重复)
0
Producer 发送次数
0
实际写入条数
0
被 Broker 去重
0
乱序拒绝

事务 Producer 两阶段提交时序图(可步进)

点「下一步」逐步推进事务从 init → begin → produce → prepare → commit marker → done 的 11 个步骤。每一步都展示 Producer / Coordinator (`__transaction_state`) / 各 Partition Leader 的状态变化。

事务步骤

Coordinator 与各分区状态

__transaction_state
Empty
orders-0 日志
outbox-0 日志
__consumer_offsets-N
就绪

Read Committed 演示:未提交消息对消费者不可见

左侧是分区物理日志(含 commit/abort marker),右侧分别展示 read_uncommitted 与 read_committed 消费者看到的内容。注意 LSO 标记。

orders-0 物理日志

就绪

read_uncommitted Consumer

HW 之前的所有内容都看得到(包括未提交的事务消息)

read_committed Consumer

LSO = 0
只读到 LSO 之前的稳定结果(commit 才显示,abort 直接丢)

PID Fencing 演示:旧 Producer 被 Epoch+1 的新实例顶替

启动新 Producer(同一个 transactional.id),Coordinator 把 Epoch 加一并通知所有分区 Leader。旧实例下次发请求时携带 Epoch=旧值,Broker 拒绝并返回 INVALID_PRODUCER_EPOCH。

Producer 实例池
transactional.id = tx-order-1

Coordinator 状态

tx_idPIDEpoch

Broker 响应日志

就绪

consume-process-produce 端到端流转

Streams 经典模式:从 input 读 → 处理 → 写 output + 提交 input offset,全部在一个事务里。下游 read_committed 消费者只看到 commit 后的稳定结果。

流转控制

就绪

Pipeline

① input topic (raw_orders)
② Stream App 内部状态
未消费
③ output topic (cleaned_orders)
④ input offset (in __consumer_offsets)
commit 之前 = -1
⑤ 下游 read_committed Consumer 看到