Durable Streams 百万流基数性能修复:CARDINALITY_1M 根因分析与优化实践
2026/9/15 18:40:49 网站建设 项目流程

Durable Streams 百万流基数性能修复:CARDINALITY_1M 根因分析与优化实践

【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric

导读

本文基于 electric 仓库 packages/durable-streams-rust/CARDINALITY_1M.md 展开,系统梳理 Rust 版 Durable Streams 服务在**高流基数(stream cardinality)**场景下吞吐断崖的四个根因、对应的源码级修复手段,以及 20k→1M 流的本地 A/B 与远程 GKE 验证数据。读完本文,你将理解"每条流每检查点周期平均操作数低于 1 时,被摊还的成本如何退化为每次操作成本"这一核心机理,掌握 WAL 检查点、meta sidecar、注册表查找等路径上避免 O(touched)/O(streams) 热路径成本的具体工程手段,并能够复现与观测这些瓶颈(--wal-stats、repro harness)。

背景速览:Durable Streams 是构建在同步之上的开源 agent 平台 electric 中的数据基元协议实现,这是一个无数据库、无 broker的单二进制 Rust 服务,每条流在磁盘上就是一个与线上字节完全一致的连续数据文件加一个.meta旁路文件,写即追加、读即字节区间(详见 README.md 与 ARCHITECTURE.md)。

问题背景:从 500k 到 1M 的"悬崖"

此前的工作(WRITE_BOTTLENECKS_1M.md的瓶颈 #2,以及CONTENTION_INVESTIGATION.md)在perf/combined-t1a-t1c-t2a分支上定位了流基数问题。本文档对应的服务端提交为662b0c845

在修复前,当流数量从 500k 增长到 1M 时,同负载下吞吐出现断崖式下跌;修复后的结果(2026-07-02):

  • 1M 流在 16 vCPUc4d-standard-16-lssd实例上达到 1,114,644 ops/s,且压测阶梯(ladder)尚未饱和;
  • 500k→1M 同负载下的退化仅为−17%(修复前是断崖)。

核心机理:摊还成本在何时退化为逐操作成本

文档给出了一条可复用的判断准则:

在高基数下,每条流每个检查点周期(checkpoint interval)的平均操作数跌破 1,于是所有"每条流每周期摊还一次"的成本都变成了每次操作的成本。

这正是理解下面四个根因的钥匙。检查点(checkpoint)默认每 ~3 秒触发一次(见 wal/shard.rs 中的--wal-checkpoint-interval-ms,默认 3000ms)。当流总数很大而每条流的写入稀疏时,一个检查点周期内大量流只被追加过一次甚至零次,那些本应"一个周期只做一次"的收集、捕获、fsync、sidecar 写入就被摊薄到几乎每次追加上。

根因一:检查点 drain 在持有 dirty 互斥锁时执行 O(touched) 捕获

问题

检查点 drain 需要对每个 touched 流执行shared.read()+Arcclone(捕获逻辑 tail 与 live file 句柄),而这在持有 shard 的dirty互斥锁期间完成。问题在于:这个锁正是每次追加在 epoch 迁移时都要取的同一把锁——在 400k 流时,几乎每次追加都是一次 epoch 迁移(因为每条流每个周期平均被追加次数 < 1)。

证据(--wal-stats观测):

  • WAL_CKPT drain_us达到 25–140 ms/tick;
  • dirty_wait_load在流数从 20k 增到 400k 时从 0.01 涨到 0.30 核;
  • p99 约等于最大 drain 时长。

修复

将临界区降为 O(1):在锁内只做"取走 Vec + 递增 epoch",捕获工作全部放到释放锁之后执行。对应实现见 wal/shard.rs:

let drained: Vec<Arc<StreamState>> = { let mut g = self.dirty.lock().unwrap(); self.dirty_epoch.fetch_add(1, Ordering::AcqRel); std::mem::take(&mut *g) };

同一段代码注释里明确记录了这次修复的动机:高基数下每个周期内几乎所有追加都走在"首次 touch"的迁移路径上,任何 O(touched) 的锁内工作都会让该 shard 上所有 appender 在整个 drain 期间被阻塞(在 400k 流时实测每 tick 25–140 ms)。

同时,register_dirty本身也做了无锁化(Tier-1a):去重不再依赖每次追加在锁内做 HashMap insert,而是用StreamState.dirty_epoch与 shard 的dirty_epoch做比较:

  • 热路径(该流本周期已注册):仅两次 relaxed load + 一个分支,完全不碰dirty锁;
  • 迁移路径(本周期首次 touch):CAS 成功后由唯一获胜者 push 进dirtyVec(每流每周期至多一次),因此 Vec 的Mutex完全离开热路径。

见 wal/shard.rs。dirty_wait_load因此塌缩到接近 0。

根因二:检查点主体跑在异步运行时线程上且各 shard 串行

问题

检查点的三个阶段——capture、累计 tails 文件的"重读 + 重排 + 重写"、recycle——原本运行在 async runtime 线程上,并且各 shard 串行执行。观测数据:每个 shard 每 tick 在 runtime 线程上消耗 capture 31 ms + tails 24 ms。

修复

  • 整个检查点主体收敛到单个spawn_blocking,不再占用异步 worker 线程(见 wal/shard.rs 的tokio::task::spawn_blocking+Self::checkpoint_blocking);
  • tails 映射改为内存驻留:持久化的累计 per-stream durable-tail 映射(<shard_dir>/tails,即 task 11b)在首次需要时从磁盘播种,此后每次检查点直接在内存中合并再序列化,不再每 ~3 秒重读 + 重解析整个文件(400k 流时原先约每 tick 20 ms),见 wal/shard.rs;
  • 各 shard 的检查点并发执行JoinSet),shard 之间不再互相等待。

checkpoint_blocking内部保留了严格的硬顺序:捕获 tails → fdatasync/syncfs 每条 touched 流文件 → 持久化 tails 映射 → 持久化checkpoint_lsn→ recycle WAL 段,并且"先 barrier 后 recycle"是硬约束(见 wal/shard.rs)。ack 从不 gate 在检查点上(检查点只是限制保留 WAL 大小 = 崩溃重放时间的阀门),因此把检查点整体移出热路径是零 ack 代价的。

根因三:每次追加都做 meta sidecar 同步刷盘(数据目录 inode rwsem 争用)

问题

这是 38–46% 的全服务器 CPU消耗来源,而且与流基数无关,任何基数下都成立。原实现中,当两次追加的间隔超过 100 ms 的 debounce(高基数下几乎必然如此),每次 producer 追加都会执行 JSON 序列化 +File::create(.meta.tmp)+rename,导致所有 worker 都自旋在数据目录 inode 的 rwsem上(perf 显示osq_lock+rwsem_spin_on_owner出现在write_meta_sync之下)。

修复

  • WAL 模式下的追加只标记meta_dirty,不再同步刷盘;
  • sidecar 的写入由检查点在 drain 完成后的 recycle 阶段统一执行
  • memory 模式保留 debounced flush(经由 store 级周期 sweeper 批量冲刷,对应mark_meta_dirty队列,见 main.rs 与 handlers.rs)。

producer/access 状态的陈旧度上界从 100 ms debounce 变为检查点周期(约 3 秒)——协议本就允许这种滞后(contract already allows lag),因此这是被明确接受的权衡。源码注释还给出了精确的成本记录:"meta 的File::create+rename(及其父目录 rwsem,在写饱和下实测约占服务器 CPU 的 40%)以及 timer task 都移出了每次追加路径"(见 handlers.rs)。

关键细节——这个修复是带门控(gated)的:只有追加真的改变了 sidecar 必须持久化的状态时才标记meta_dirty,即producer/seq 幂等状态或滑动 TTLmeta_persist_needed = producer.is_some() || seq_header.is_some() || st.config.ttl_seconds.is_some(),见 handlers.rs)。普通非 TTL 流的追加不需要任何 sidecar 重写:

  • durable_tail由检查点 per-shardtails映射记录(该映射而非 sidecar 才是恢复时对账的权威 durable-tail 证据,见 handlers.rs);
  • last_access仅用于 TTL 门控;
  • memory 模式下恢复时 tail 从数据文件长度重新推导(Store::new_with_tier)。

配套测试覆盖了这两条语义,见 handlers.rs:memory_append_defers_sidecar_to_store_sweep(memory 模式追加不得绕过 store 级 sweeper 立即刷 sidecar)与memory_plain_append_skips_sidecar_flush(普通非 TTL 追加不得排队 sidecar 刷写)。

根因四:每次追加两次注册表查找

问题

每次追加会做两次 registry 查找(handle_append的 metric label 一次 +_inner一次),即 2× SipHash + 在 1M key 的冷 DashMap 上遍历。此根因主要通过代码审查确认。

修复

_inner现在直接返回is_json,把一次查找合并进既有路径,去掉了每次追加的多余 DashMap 遍历。

顺带说明:流本身存放在DashMap(streams 的注册表)中,每个流有独立的 per-stream appender mutex——唯一的串行化点且按流隔离,不同流之间从不互相争用;读走"短暂快照 + 定位 read",从不阻塞写者。这套无全局锁的设计是架构级前提(见 ARCHITECTURE.md)。

本地 A/B:Linux harness 验证(6 srv 核,conn=256,shards=6)

流数修复前修复后p99
20k~43–46k ops/s80.4k41 → 7.7 ms
200k32.3k50.6k53 → 17.9 ms
400k16.0k36–44k144 → ~28 ms

正确性:95 个 crate 单元测试 + 326 个 conformance 测试全部通过(conformance 套件位于 conformance/conformance.test.ts,运行方式见 README.md 的RUST_SERVER_URL=... pnpm exec vitest run --config packages/durable-streams-rust/conformance/vitest.config.ts)。

新增可观测性:

  • WAL_CKPT--wal-stats下每个 shard 的按阶段检查点时间行(capture/fsync/tails/checkpoint/recycle 各阶段一次时钟读取,每 shard 每 ~3 秒一次,完全不在热路径上,见 wal/shard.rs);
  • repro harness 新增--tmpfs/--wal-stats旋钮,并汇总 WAL_CKPT 输出。

远程验证:GKE 16 vCPU,1M 流

  • 套件:ds-bench/suites/run-durable-cpu16-1m-card.json,镜像durable-streams:combined-card@sha256:d74840bd…
  • 完整细节与 caveats 见ds-bench/results/run-durable-cpu16-1m-card/FINDINGS.md(该目录属于仓库外的 ds-bench 独立基准仓库,本仓库仅引用其结果路径)。
流数podsops/sp50 / p99 / max
1M32898,5823.5 / 30.5 /149.6 ms(baseline 862k,3.3 / 32 / 405 ms)
1M641,114,6443.4 / 60.3 / 211 ms —— 仍在爬升(+21%/+16 pods)
500k481,110,2683.1 / 42.7 / 145 ms

要点解读:

  • 与修复前基线(baseline 862k)相比,同样 32 pods 时 p99/max 从 405 ms 降到 149.6 ms;
  • 64 pods 时吞吐还在随资源增长而爬升,说明 16 vCPU 的真实天花板尚未探到;
  • 500k 与 1M 在相近 pod 数下吞吐几乎持平(1,110,268 vs 1,114,644),印证"无基数悬崖"的结论。

复现与观测方法

在 packages/durable-streams-rust 下构建并运行:

cargo build --release # → ./target/release/durable-streams-server cargo test --release # 95 个 crate 测试(单元 + 集成)

启动时用--wal-stats观察检查点各阶段耗时:

./target/release/durable-streams-server --port 4438 --data-dir ./data \ --wal-stats 1

--wal-stats N表示每 N 秒打印一次(解析见 main.rs),会输出 per-shard 的WAL_STATS行与检查点WAL_CKPT行。repro harness 可配合--tmpfs(把数据目录放到 tmpfs 上隔离磁盘因素)与--wal-stats复现 20k→400k 的基数实验。

高基数场景下影响检查点行为的两个旋钮(见 wal/shard.rs):

  • --wal-checkpoint-interval-ms(默认 3000):时间触发间隔;
  • --wal-checkpoint-wal-bytes(默认 0 = 禁用):shard 保留 WAL 超过该字节数即触发检查点,把硬编码定时器变成显式重放时间预算,并让各 shard 按各自写速率自我错峰(避免共享 tick 上同时风暴)。

遗留问题与后续方向

  1. 1M 流下 16 vCPU 的真实天花板:压测阶梯还未越过 64 pods,需要更长的 ladder,以及 32 vCPU 上的 1M 流运行;
  2. per-shard producer-state journal(进行中):全基数下 sidecar 写入速率 ≈ ops/s(现已移出热路径,但会拉长检查点周期——本地 400k 流时 meta 阶段约 1.6 s/shard——并限制 producer-state 的陈旧度上界)。计划方向:每 tick 一个累计的 per-shard 文件,恢复时按max(epoch, seq)叠加 producer;
  3. 1M 流下的读写混合:目前所有数据均为纯写;
  4. NVMe 上的--wal-statscellrun-durable-cpu16-1m-card-stats.json,尚未运行):用于确认真实磁盘上 1M 流的检查点 fsync/meta 阶段行为。

相关文档与源码索引

  • 本文主文档:CARDINALITY_1M.md
  • 架构总览:ARCHITECTURE.md(写路径、durability 模式、I/O 加速技术)
  • 服务使用与全部 flag:README.md
  • 并发与争用背景:CONTENTION_INVESTIGATION.md
  • 崩溃仿真:CRASH_SIM_FINDINGS.md
  • WAL 调优:WAL_TUNING.md
  • 混负载验证:MIXED_WORKLOAD_VALIDATION.md
  • 核心实现:src/wal/shard.rs(检查点、dirty 集、group-commit watermark)、src/handlers.rs(追加路径、sidecar 门控)、src/main.rs(meta sweeper、--wal-stats解析)
  • 协议一致性测试:conformance/conformance.test.ts

【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询