ADR-0012 FanoutCoordinator 不使用 Kafka 事务¶
状态¶
已接受(2026-08-25)。取代 ADR-0006 中「Fanout 使用 consume-transform-produce 事务」的实现选择;
日志中间件仍为 Redpanda,不变。
一句话结论:Kafka 事务从来不是逻辑唯一性的来源,唯一性一直由确定性 dispatch_id
与 MailboxStore 的持久去重保证。去掉事务后端到端 P50 降 19%、P99 降 80%,
且消除了长期存在的 P99 门禁不稳定。
背景¶
事务提供了什么¶
publish_dispatches_and_ack_batch 原实现:
begin_transaction()
produce_dispatches() // 产出该批全部 dispatch
send_offsets_to_transaction() // 把 Outbox 消费位点并入事务
commit_transaction()
它保证「dispatch 记录产出」与「Outbox 位点推进」原子。
但唯一性从来不靠它¶
CLAUDE.md 中已冻结的决策写明:
Kafka
enable.idempotence只覆盖同一 producer 的协议重试,不提供应用级 publish 去重。 Fanout 的确定性dispatch_id与 MailboxStore 的dispatch_id + payload_digest持久收敛才是逻辑唯一性边界,文档和指标不得宣称跨崩溃 exactly-once。
即:重复的 dispatch 记录本就允许存在,由下游去重收敛。事务只是减少崩溃后的重复量。
成本¶
实测(隔离环境,60 秒窗口,split Redis):
| 指标 | 2000 msg/s | 5000 msg/s |
|---|---|---|
| 事务提交段耗时 | 13.34 ms | 17.62 ms |
| Outbox lag | 11.41 ms | 24.77 ms |
每批一次两阶段提交,需依次完成 AddPartitionsToTxn → AddOffsetsToTxn →
TxnOffsetCommit → EndTxn 及各参与分区的 marker 写入。
任一跳抖动即成为端到端长尾。
决策¶
不使用 Kafka 事务。 改为:
produce_dispatches() // 产出该批全部 dispatch
-> 等待全部 durable delivery 确认(acks=all)
consumer.commit(offsets, Sync) // 全部确认后才推进 Outbox 位点
producer 去掉 transactional.id 与 init_transactions,保留 enable.idempotence=true
与 acks=all。消费端隔离级别不变。
后续(ADR-0017,2026-09-26):dispatch 生产者改为默认关闭幂等,以解除每连接 5 个在途请求的上限;
acks=all与“先 durable、后提交位点”的顺序不变。
正确性论证¶
去掉事务后只可能出现两类失败:
| 失败 | 后果 | 是否可接受 |
|---|---|---|
| A:dispatch 已产出,位点提交失败 | 重放该批 → 重复产出 dispatch → 由 dispatch_id + payload_digest 去重收敛 |
可接受(本就是既有语义) |
| B:位点已提交,dispatch 产出失败 | 该批 source 被跳过 → 消息永久丢失 | 不可接受 |
因此唯一的不变量是:位点提交必须严格晚于该批全部 dispatch 的 durable delivery 确认。 只要该顺序成立,B 不可能发生。
该不变量已写入 publish_dispatches_and_ack_batch 的注释,与
「lane_outcome 必须与 lane 下标无关」(ADR-0011)同级,
属于改写代码前必须先确认仍成立的一类约束。
位点提交失败时返回错误、不推进调用方的 after,本批整体重放——落入 A,安全。
与既有语义的关系¶
本决策不改变任何持久事实:
dispatch_id的确定性构造不变- MailboxStore 的
dispatch_id + payload_digest去重不变 - DispatchProgress 的 Fresh/Existing 语义不变
- 连续水位
W的推进条件不变 - Outbox / dispatch topic 的 retention 与 recovery window 不变
变化的只是「崩溃后重复 dispatch 的期望数量」——从「几乎为零」变为「最多一批」。
批上限由 QIM_FANOUT_BATCH_SIZE 界定(默认 128,合法 1..=1024)。
实测结果¶
隔离环境,60 秒窗口,split Redis,与改动前同口径。
分段指标¶
| 2000 msg/s | 5000 msg/s | |||
|---|---|---|---|---|
| 有事务 | 去事务 | 有事务 | 去事务 | |
| 提交段 | 13.34 ms | 9.15 ms | 17.62 ms | 10.79 ms |
| Outbox lag | 11.41 ms | 9.09 ms | 24.77 ms | 19.77 ms |
端到端(2000 msg/s)¶
| P50 | P95 | P99 | max | |
|---|---|---|---|---|
| 有事务 | 77.5 ms | 101.0 ms | 455.7 ms | 696.9 ms |
| 去事务 | 62.4 ms | 81.4 ms | 91.5 ms | 155.3 ms |
P50 −19%,P95 −19%,P99 −80%,max −78%。
5000 msg/s 下端到端无显著变化(12.96 s vs 13.02 s),符合预期——
该速率下瓶颈是 Redis 单线程物化(见 ADR-0011),fanout 本就不是限制项。
附带发现:P99 门禁不稳定的根因¶
ADR-0011 记录过一个未解问题:2000 msg/s 下端到端 P99 在多次重复测量中
剧烈波动(125.5 / 222.3 / 351.6 / 351.9 / 455.7 ms),而同批 P50、P95 稳定。
按 CLAUDE.md 纪律,这类「会随机红」的门禁最终必然被当成「本来就飘」而失效。
去掉事务后三次重复测量:
| P50 | P95 | P99 | max | |
|---|---|---|---|---|
| 第 1 次 | 62.4 ms | 81.4 ms | 91.5 ms | 155.3 ms |
| 第 2 次 | 62.7 ms | 81.9 ms | 92.7 ms | 143.5 ms |
| 第 3 次 | 62.5 ms | 81.4 ms | 91.5 ms | 173.7 ms |
P99 三次方差不足 1.5 ms。
结论:该门禁的不稳定根因即 Kafka 事务两阶段提交的长尾,现已消除。
ADR-0011 中「需改为多次取中位数或将判据移至 P95」的建议随之作废——
当前 P99 已是可靠判据。
验证¶
qim-commit-log/qim-fanout/qim-store/qim-mailbox单元与集成测试全通过- 真实 Redpanda 门禁 4/4 通过,含
redpanda_事务初始化并连续消费_outbox_backlog(覆盖新的 「produce 全部 → 等 durable → 提交位点」顺序)、redpanda_dispatch_consumer_验证全部静态分区、 两项 retention/recovery window 失败闭合用例 cargo fmt/clippy -D warnings/check-contract.sh通过
后果¶
收益¶
- 端到端 P50/P95 各降约 19%,P99 降 80%,尾部分布显著收紧
- 消除 P99 门禁的随机红,恢复其作为发布判据的有效性
- 去掉两阶段提交后,fanout 提交路径只剩 produce 与 offset commit,
代码与失败分支均减少(
abort_transaction_after_failure一并删除)
风险与遗留¶
- 顺序不变量是唯一防线。 若未来有人把位点提交移到 produce 确认之前, 将从「安全重复」直接变为「永久丢失」,且无任何测试能自动捕获顺序颠倒 (两种顺序在无崩溃时行为完全一致)。该约束必须靠注释与评审保证。
- 崩溃后重复 dispatch 的期望量从「几乎为零」升至「最多一批」(默认 ≤128 条)。 下游去重成本随之小幅上升,但属既有语义内。
- 本决策未触及 Outbox 单分区与 fanout 串行循环这两项结构限制, 它们仍是 fanout 不可水平扩展的原因;但实测表明 fanout 当前不是瓶颈, 该限制暂不构成问题。
参考¶
ADR-0011(邮箱热路径的成本归因与实现层裁决)ADR-0006(服务端语言与日志选型;Redpanda 选定,本 ADR 不改变该选择)- 实现:
crates/qim-commit-log/src/lib.rs的publish_dispatches_and_ack_batch