跳转至

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