幂等 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_id | PID | Epoch |
|---|
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 看到