问题描述
分布式事务(prepared)与大事务(partialTxn)的数据可能在 MongoShake 重启时静默丢失,全程无任何报错。
时序:
t1 prepare/partialTxn oplog (ts=100) collector 仅缓存在内存 TxnBuffer,不下发 worker
t2 普通 DML oplog (ts=101) 被下发并 ack,checkpoint 持久化到 101
t3 进程崩溃(kill -9 / OOM / 断电) TxnBuffer 内存数据丢失
t4 重启,oplog 拉取查询为 {ts: {$gt: 101}} ts=100 的 prepare 永远不会被重新拉取
t5 commitTransaction (ts=102) 到达
t5 时当前代码不会报错:handleTransaction 先调 AddOp,为该未知事务新建一个空 state,GetTxnStream 正常流出 0 条,事务被当作「空事务」完成——静默丢数据。
根因(已代码级验证)
prepare/partialTxn oplog 只存在于内存 TxnBuffer(oplog/txn_buffer.go,注释自证 "state is currently kept in memory"),不下发 worker(collector/batcher.go handleTransaction 对非 final 条目提前返回)。
calculateWorkerLowestCheckpoint(collector/checkpoint.go)只看 worker ack/unack,完全不感知 TxnBuffer。上游 mongo-tools 的 TxnBuffer.OldestOpTime() 正是为此设计,但生产代码中零调用(仅单测使用)。
- 重启拉取用严格
$gt checkpoint 查询(collector/reader/oplog_reader.go),checkpoint 一旦越过 buffered 条目即永久跳过。
- 第二窗口:commit 处理中
PurgeTxn 先于内层 op(暂存内存 barrierOplogs)被下发/ack,该窗口内崩溃同样丢数据。
复现方式
真实分片集群的 2PC prepare→commit 窗口为毫秒级,无法稳定复现。已用以下方式端到端验证(集成测试见修复分支):
- 2-shard 集群 + 跨 shard 事务;
docker pause shard2 将 2PC 冻结在 prepared 状态(shard1 oplog 已出现 prepare 条目);
- shard1 上普通 DML 同步并 ack;
- 未修复版本此时 checkpoint 恰好追平 prepare ts(丢数据窗口成立);
- kill -9 collector 并重启;unpause 完成 commit → 事务数据丢失。
修复方案(已在分支 worktree-fix-txn-checkpoint-pin 实现并验证)
- 事务 pin:batcher 为每个事务登记 pin(startTs = 首条 buffered 条目 ts);abort 立即释放;commit 后待 worker ack 达到 commit ts 才释放(覆盖 PurgeTxn 先于下发的窗口)。
- checkpoint clamp:持久化 checkpoint 封顶为
oldest pin - 1,保证重启必能重拉 buffered 条目。
- barrier 两轨:
checkCheckpointUpdate 的等待增加 worker-ack 进度作为等价判定,避免与 clamp 死锁。
- 孤儿 commit 防御:
GetTxnStream 发现某 state 从未收到数据条目(只有 final commit)时明确报错,拒绝静默空事务。
验证:单元测试红绿闭环;集成测试 integration/txn_pin_test.go(build tag integration)在修复版通过(checkpoint 钉在 prepare-1、两 shard 数据完整),在 develop 基线失败(checkpoint 追平 prepare)。
相关:修复 commit 3cafd3e、集成测试 commit 2fcc387(分支 worktree-fix-txn-checkpoint-pin,PR 待提)。
问题描述
分布式事务(prepared)与大事务(partialTxn)的数据可能在 MongoShake 重启时静默丢失,全程无任何报错。
时序:
t5 时当前代码不会报错:
handleTransaction先调AddOp,为该未知事务新建一个空 state,GetTxnStream正常流出 0 条,事务被当作「空事务」完成——静默丢数据。根因(已代码级验证)
prepare/partialTxnoplog 只存在于内存TxnBuffer(oplog/txn_buffer.go,注释自证 "state is currently kept in memory"),不下发 worker(collector/batcher.gohandleTransaction对非 final 条目提前返回)。calculateWorkerLowestCheckpoint(collector/checkpoint.go)只看 worker ack/unack,完全不感知 TxnBuffer。上游 mongo-tools 的TxnBuffer.OldestOpTime()正是为此设计,但生产代码中零调用(仅单测使用)。$gt checkpoint查询(collector/reader/oplog_reader.go),checkpoint 一旦越过 buffered 条目即永久跳过。PurgeTxn先于内层 op(暂存内存barrierOplogs)被下发/ack,该窗口内崩溃同样丢数据。复现方式
真实分片集群的 2PC prepare→commit 窗口为毫秒级,无法稳定复现。已用以下方式端到端验证(集成测试见修复分支):
docker pause shard2将 2PC 冻结在 prepared 状态(shard1 oplog 已出现 prepare 条目);修复方案(已在分支 worktree-fix-txn-checkpoint-pin 实现并验证)
oldest pin - 1,保证重启必能重拉 buffered 条目。checkCheckpointUpdate的等待增加 worker-ack 进度作为等价判定,避免与 clamp 死锁。GetTxnStream发现某 state 从未收到数据条目(只有 final commit)时明确报错,拒绝静默空事务。验证:单元测试红绿闭环;集成测试
integration/txn_pin_test.go(build tagintegration)在修复版通过(checkpoint 钉在 prepare-1、两 shard 数据完整),在 develop 基线失败(checkpoint 追平 prepare)。相关:修复 commit 3cafd3e、集成测试 commit 2fcc387(分支 worktree-fix-txn-checkpoint-pin,PR 待提)。