Skip to content

使用 Pebble 替换 oplog 本地落盘队列 - #987

Merged
zhongli-james merged 7 commits into
alibaba:developfrom
SisyphusSQ:develop
Jul 27, 2026
Merged

使用 Pebble 替换 oplog 本地落盘队列#987
zhongli-james merged 7 commits into
alibaba:developfrom
SisyphusSQ:develop

Conversation

@SisyphusSQ

@SisyphusSQ SisyphusSQ commented Jun 26, 2026

Copy link
Copy Markdown
Contributor

背景

关联 issue:#982

当前 full_sync.reader.oplog_store_disk=true 路径依赖 go-diskqueue 作为本地 oplog 临时队列。本 PR 将该路径替换为 Pebble-backed spool,保持原有用户配置入口不变,并补齐 checkpoint、异常重启恢复、指标和测试覆盖。

主要变更

  • 新增 collector/spool 内部接口和基于 Pebble v1.1.5 的实现,支持 FIFO 批量写入、顺序读取、进程内推进、深度统计、关闭重开、删除和旧格式 fail fast。
  • oplog data、write_seqlast_write_ts 使用同一个 batch,并通过 Commit(pebble.Sync) 原子同步到可恢复持久点。
  • read_seq 仅作为进程内回放游标;spool 重开后从 seq 1 开始 at-least-once replay,并使用 durable last_write_ts 恢复 source reader。
  • 改造 Persister 的磁盘落盘/回放状态机:全量期间批量写入 Pebble spool,全量完成后顺序回放到 pending queue,worker checkpoint 追平后再进入清理。
  • 非空 spool 采用“持久化 finish marker → 删除本地 spool → 清除远端 queue metadata”的两阶段清理,Phase B 失败可以重试。
  • checkpoint 更新采用 copy-on-success,HTTP 非 200 视为错误;容量统计改用 db.Metrics().DiskSpaceUsage(),不再逐批遍历目录。
  • 新增 raw BSON timestamp helper,并在 oplog reader 发送 raw bytes 前 clone 数据,避免 cursor buffer 复用影响 spool timestamp 解析。
  • /persist 增加 spool_depthspool_read_seqspool_write_seq;Prometheus 增加 spool_depthspool_write_seqspool_read_seqspool_write_totalspool_read_totalspool_errors_total
  • 限制 full_sync.reader.oplog_store_disk 仅用于 oplog fetch method,避免 change stream 路径被误用。

兼容性与升级说明

  • 继续沿用 full_sync.reader.oplog_store_disk_max_size,但语义从旧 go-diskqueue 的单个队列文件字节上限调整为整个 Pebble spool 的 MiB 软上限。默认值 256000 表示约 250 GiB,并在批次写入前检查;配置注释和启动日志会显示单位、字节换算、作用范围和 spool 路径。
  • 不迁移旧 go-diskqueue 文件或 PR 早期迭代产生的 Pebble schema v1;检测到不兼容格式时 fail fast,并要求先确认 source oplog 覆盖范围再处理遗留目录和 checkpoint。
  • 当未完成的 spool 仍然存在,并在 disk replay 期间发生重启时,spool 会从 seq 1 整体重放,提供 at-least-once 而非 exactly-once 语义。已经投递或应用的 oplog 可能再次 apply;正确性依赖 MongoShake 现有的幂等 apply 和重复处理语义。

持久化与性能取舍

  • oplog 先按现有 buffer 容量或大小阈值聚合,每批执行一次 Commit(pebble.Sync),而不是每条 oplog 执行一次 fsync。
  • 相比原始 NoSync 方案,这会增加批次提交延迟,并可能降低峰值写吞吐;本 PR 接受这一取舍,以保证成功返回的批次中 data、write_seqlast_write_ts 处于同一个可恢复持久点。

不包含

验证

本地验证

  • 相关 focused tests、race tests、go vet、collector/receiver build 和 git diff --check 均通过。
  • 本次说明性注释提交后重新执行 go test ./collector/spool -count=1、格式检查和 git diff --check,结果通过。
  • 旧的外部 Mongo fixture 依赖用例不纳入本轮本地门禁,不以 go test ./... 作为通过标准。

环境验证

  • 复制集和 6-shard 环境均完成 open/write/read/retain/cleanup/incremental 生命周期。
  • 两轮均达到 full progress=1,各 shard spool depth=0,checkpoint 持续增长,错误关键字扫描为 0。
  • 两轮均使用 debug executor,不写目标端业务数据。
  • 在 6 个 spool 持续写入期间硬终止测试 collector,确认进程退出后 durable spool 全部保留;使用同一配置重启后,6 个旧 spool 均以 create=false 重新打开,6 个 shard 均从 durable last_write_ts 恢复读取,进程、Metrics 与后续 spool 写入正常,错误扫描为 0。

验收边界

@zhongli-james zhongli-james self-assigned this Jul 6, 2026
@zhongli-james

Copy link
Copy Markdown
Collaborator

感谢贡献!~我尽快评估一下

@zhongli-james

zhongli-james commented Jul 13, 2026

Copy link
Copy Markdown
Collaborator

@SisyphusSQ 非常感谢这个高质量的 PR!整体设计和实现都很扎实,尤其认可几点:用 read_seq/write_seq 单调序列重建 FIFO 并解决了原来 //TODO, there is a bug if MongoShake restarts 的重启恢复问题;cloneOplogRaw 修复了直接投递 cursor 内部 buffer 可能读到脏数据的隐患;以及 ExtractRawTimestamp 避免完整 unmarshal 的优化。测试和指标也很完整。

合并前有几个点想和你确认一下:

  • 配置项语义:full_sync.reader.oplog_store_disk_max_size 原先作为 go-diskqueue 的 maxBytesPerFile(字节),现在被当作整个 spool 的 MB上限(MaxBytesMB),默认值 256000 的实际含义因此从「每文件约 250KB」变成「总量约256GB」。对已经显式配置过这个值的老用户会是行为突变。是否考虑新增独立配置项,或至少在文档和启动日志里明确单位/含义的变化?
  • 持久化语义:目前 Put/Advance/init 全部使用 pebble.NoSync,只有 Close()Flush()。想确认在非优雅退出(kill -9 / OOM/断电)时的数据保证——此时未刷盘的数据和 read_seq 会丢失,重启后与 checkpoint 配合是否可能导致 oplog空洞(丢数据)或重复回放?能否补充一个「非优雅关闭后重开」的恢复用例?(现有测试覆盖的是优雅 Close 后重开。)
  • 写入热点:ensureUnderMaxBytesLocked 每次 Put 都会 directorySize 遍历整个目录。全量高频写入时这是 O(files)
    的开销,建议缓存目录大小或降低检查频率。
  • 依赖引入:引入 Pebble 会带来 go 1.23 及 cockroachdb-errors、sentry-go、gogo-protobuf等一批传递依赖。想了解你在选型时是否评估过更轻量的方案?相比纯 FIFO 队列,完整 LSM 的收益(随机读/事务 batch/成熟度)是否是这里必需的?
  • 另外一个小点:checkpoint.godiskQueueLastTs = -2 // mark -1 so next time won't call 的注释和代码值对不上(沿用了旧注释),顺手修一下即可。

期待你的回复,我们一起把它推进到可合并状态~

@SisyphusSQ

Copy link
Copy Markdown
Contributor Author

逐项确认如下:

1. 配置项语义

这里继续沿用现有配置项:

full_sync.reader.oplog_store_disk_max_size

但会明确说明它在 Pebble spool 下的新语义:

  • go-diskqueue 实现中,表示单个队列文件的最大字节数;
  • Pebble spool 实现中,表示整个 spool 目录的最大容量;
  • 新单位为 MB;
  • 默认值 256000 表示总容量约 250 GiB

由于 Pebble 的 SST、WAL、MANIFEST 和 compaction 文件模型与 go-diskqueue 的分段文件不同,旧的“单文件上限”没有可以直接对应的 Pebble 参数,因此这里选择保留配置入口,将其调整为更有实际意义的 spool 总容量限制。

为了避免升级后出现静默语义变化,我会:

  1. 更新配置文档和示例注释,明确单位是 MB、作用范围是整个 Pebble spool;
  2. collector 启动时输出配置值、换算后的字节数、作用范围和 spool 路径;
  3. 在升级说明中明确该配置从“单文件字节上限”调整为“spool 总容量 MB 上限”,提醒显式配置过该值的用户重新确认容量。

启动日志会类似:

pebble oplog spool max size configured: value=256000 unit=MB scope=total_spool

2. 持久化与异常退出

当前 NoSync 配合 Close() 时执行 Flush(),只能覆盖优雅退出,不能作为 kill -9、OOM 或断电场景下的持久化保证。现有测试也只覆盖了正常关闭后重开,因此目前不能严谨地承诺异常退出时不会产生 oplog 空洞。

这里需要同时处理两个边界:

  1. Put 返回成功前,oplog data、write_seq 和 last-write meta 是否已经到达可恢复的持久化点;
  2. read_seq 的推进和 spool 删除,是否早于 worker checkpoint 能够确认的安全位置。

我会重新梳理写入同步和 spool 生命周期策略,避免只依赖正常关闭时的 Flush()。同时补充以下恢复测试:

  • 子进程写入后不调用 Close() 直接退出,再重新打开 spool;
  • 写入过程中异常退出,检查 data 和 write_seq 是否保持一致;
  • read_seq 推进后、worker checkpoint 尚未追平时异常退出;
  • spool 已完成内存投递、但 checkpoint 尚未清理 queue metadata 时异常退出。

最终会明确恢复保证:不能产生 oplog 空洞;如果某些异常窗口只能保证至少一次投递,也会把可能重复回放的边界说明并测试清楚。

3. 写入热路径目录扫描

这个问题成立。当前每次 Put 都通过 filepath.WalkDir 统计整个 spool 目录大小,随着 SST、WAL 和 compaction 文件增加,会形成与文件数量相关的热路径开销。

这里会改为使用 Pebble 自身维护的:

db.Metrics().DiskSpaceUsage()

该指标已经包含 WAL、live/obsolete table、MANIFEST 和 compaction 中的本地文件空间,可以避免每次写入遍历目录。Stats() 和容量限制也会统一使用这一统计口径。

4. Pebble 选型和依赖

当前 spool 的业务语义确实是 FIFO,现有调用没有依赖通用随机读能力。Pebble 在这里的主要价值不是随机读,而是:

  • 使用 batch 原子维护 oplog data、write_seq 和 last-write meta;
  • 使用有序 key iteration 实现 FIFO 恢复与回放;
  • 使用成熟的 WAL、损坏检测和恢复机制;
  • 避免继续维护自定义 segment queue 的文件轮转、尾部截断、异常恢复和清理逻辑。

因此,完整 LSM 提供的随机读能力并不是必需条件。这里的主要取舍,是用相对较重的依赖换取成熟的本地持久化与恢复实现。原方案没有形成与轻量方案之间的量化对比,这部分说明确实不足。

我对比了 Pebble v1.1.5 和当前使用的 v2.1.6

  • v1.1.5 要求 Go 1.20,因此 MongoShake 可以继续保持 Go 1.21;
  • 能减少 axisdscrlibswissminlz 等部分 v2 新增依赖;
  • cockroachdb/errorssentry-gogogo-protobuf 等依赖在 v1 中仍然存在。

因此,切换到 v1 可以解决 Go 1.23 的升级问题,并适度缩小依赖树,但不能把 Pebble 变成轻量依赖。

我会优先验证切换到 Pebble v1.1.5 后的接口兼容性、目标包测试和 collector 构建。如果验证通过,本 PR 将改用 v1,并补充依赖变化说明。同时会补充与 bbolt、segment FIFO 等方案在依赖规模、写入模型和异常恢复复杂度方面的取舍,说明继续选择 Pebble 的原因。

5. 旧注释

checkpoint.go 中:

diskQueueLastTs = -2 // mark -1 so next time won't call

注释确实与代码值不一致,是沿用旧逻辑后没有同步更新,会修改为准确描述 -2 状态含义的注释。

以上调整完成后,我会补充对应测试结果、构建结果和依赖变化,再更新 PR。

@zhongli-james

Copy link
Copy Markdown
Collaborator

好的,感谢澄清。

@SisyphusSQ

SisyphusSQ commented Jul 20, 2026

Copy link
Copy Markdown
Contributor Author

针对前面 review 关注点,本轮主要调整如下:

  • 容量单位与范围:配置语义明确为整个 Pebble spool 的 MiB 软上限,在批次写入前检查。
  • 异常退出持久化:数据与 write_seq/last_write_ts 使用 PutBatch + pebble.Sync 同批提交;进入回放前强制 flush 尾批。
  • 恢复语义:schema 升级到 v2;重开后从 seq 1 做 at-least-once replay,并从 durable last_write_ts 恢复 source reader,避免已持久化数据产生缺口。
  • 热路径扫描:容量统计改用 Pebble DiskSpaceUsage(),移除逐批目录遍历。
  • 依赖兼容性:Pebble 调整为 v1.1.5,保持 Go 1.21 兼容。
  • checkpoint 与清理:checkpoint 改为 copy-on-success;新 spool 元数据先落远端;非空 spool 使用“完成标记 → 本地删除 → 清除远端队列元数据”的两阶段清理,Phase B 失败可重试。
  • 错误边界:HTTP checkpoint 非 200、非法路径、legacy 目录和旧 schema 均 fail-fast;同步修正了旧注释和状态语义。

测试结果(已脱敏):

  • 相关 focused tests、race tests、go vet、collector/receiver build 均通过。
  • 复制集和 6-shard 环境均完成 open/write/read/retain/cleanup/incremental 生命周期。
  • 两轮均达到 full progress=1,各 shard spool depth=0,checkpoint 持续增长,错误关键字扫描为 0。
  • 两轮均使用 debug executor,不写目标端业务数据。
  • 补充完成远端 kill -9 恢复验证:在 6 个 spool 持续写入期间硬终止 collector,确认进程退出后 durable spool 全部保留;使用同一配置重启后,6 个旧 spool 均以 create=false 重新打开,6 个 shard 均从 durable last_write_ts 恢复读取,进程、Metrics 与后续 spool 写入正常,错误扫描为 0。

验收边界:

  • 本次远端 kill -9 验证范围为 collector 硬崩溃后的 spool 保留、旧 spool 重开和续读恢复;未重复等待该故障轮完成整个 full/cleanup/incremental 生命周期,完整生命周期由前述双集群验收覆盖。
  • 未验证目标端 count/checksum。
  • 旧的外部 Mongo fixture 依赖用例不纳入本轮本地门禁。

@zhongli-james

Copy link
Copy Markdown
Collaborator

@SisyphusSQ 感谢这一轮扎实的修改,我逐项对照了代码,确认关注点都已真正落地:

  • 配置语义:conf 注释、MiB 软上限和启动日志已明确单位与作用范围,升级不再静默变语义。👍
  • 崩溃持久化:PutBatchCommit(pebble.Sync) 把 data / write_seq / last_write_ts 原子同步落盘;read_seq 不持久化、load() 固定 readSeq=1 从头 at-least-once 重放,并在重开时用 LastWriteTimestamp() 让源 reader 从 durable 点续读——这套「宁可重复、绝不空洞」的设计和MongoShake 幂等 apply 语义很契合,比之前纠结 read_seq 崩溃一致要干净得多。remote kill -9 恢复验证也补上了,很好。
  • 写入热点:改用 db.Metrics().DiskSpaceUsage(),去掉了逐批目录遍历。👍
  • 依赖:回退到 Pebble v1.1.5 + Go 1.21,依赖树也收窄了。
  • 清理与旧注释:两阶段清理(finish marker → 删本地 → 清远端 queue name)顺序正确、Phase B 可重试,我确认了「删本地后崩溃不会重开已删spool」这个不变量成立;旧注释也修正了。

还有几个纯文档/说明层面的小点,处理完我这边就可以 approve:

  1. 建议在升级说明里显式点明 at-least-once 边界:disk-replay 期间重启会重放整个 spool、可能重复 apply 已应用的 oplog,正确性依赖幂等apply。
  2. 每批 pebble.Sync 相比最初 NoSync 设计牺牲了一点写吞吐——因为已批量化且不在热路径,这个取舍我认为合理,建议在 PR里一句话说明「接受该性能权衡」。
  3. schemaVersion 直接从 2 起(PR 内迭代过盘格式),加一行注释说明即可。
  4. Phase B 清理窗口的崩溃恢复、目标端 count/checksum 属本轮范围外,建议作为 follow-up issue 记录。

整体实现质量很高,感谢你的耐心打磨!补充上述说明后即可合并。

@SisyphusSQ

Copy link
Copy Markdown
Contributor Author

感谢再次逐项确认,这四点已经补齐:

  • PR 描述新增了兼容性与升级说明,明确 disk replay 重启时的 at-least-once 重放和可能的重复 apply 边界;
  • PR 描述补充了每批一次 pebble.Sync 的持久化与性能取舍;
  • schemaVersion 上方说明了版本直接从 2 开始的原因(d0f096b);
  • 新建 follow-up issue 补充 Pebble spool Phase B 崩溃恢复与目标端一致性验证 #994,跟踪 Phase B 两个清理窗口的故障注入和真实 target count/checksum 验证。

麻烦再帮忙看一下,感谢!

@zhongli-james
zhongli-james merged commit 5351444 into alibaba:develop Jul 27, 2026
1 check passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants