arjunsood2025/raft-kv-engine-project

GitHub: arjunsood2025/raft-kv-engine-project

一个用 Rust 从零构建的线性一致复制键值存储,通过 sans-IO 架构和确定性模拟测试在两万次随机故障调度中验证了正确性。

Stars: 1 | Forks: 0

# raft-kv 一个线性一致、可复制的键值数据库,使用 Rust 从零开始构建, 涵盖了存储引擎、共识、网络以及证明其正确性所需的 测试基础设施:一个 FoundationDB 风格的确定性模拟器,它已经执行了 **20,000 次随机故障调度**(约 4.2 亿个模拟事件,约 1730 万次 客户端操作),每一次都检查了线性一致性。这 20,000 次调度中,正好有一次(种子 19519)暴露了一个真实的已提交条目丢失 bug ([docs/BUGS.md](docs/BUGS.md));其余的均正常运行。 没有使用任何数据库 crate、共识 crate 或 RPC 框架。依赖 列表为:`tokio`、`serde`/`bincode`、`crc32fast`、`tempfile`(仅用于测试)。 ![节点的分层架构:客户端通过 TCP 上的长度前缀 bincode 与 tokio server 通信; 在 server 之下是 kvsm 复制状态机、无 IO 的 raft 共识核心和 LSM 存储引擎。 部署为三个 Raft 复制的节点。](https://static.pigsec.cn/wp-content/uploads/repos/cas/44/4446844243a15b454801d0c421a662e2e66c8422e26403ebc027684562f9de7f.svg)
纯文本形式的相同架构 ``` ┌─────────────────────────────────────────────────────────────────┐ │ clients (kvctl CLI · loadgen · library) │ │ leader routing · retries w/ jitter · sessions (dedup-safe) │ └──────────────────────────────┬──────────────────────────────────┘ │ length-prefixed bincode / TCP ┌──────────────────────────────┴──────────────────────────────────┐ │ server: tokio host — single-owner event loop, group commit │ │ consistency modes: ReadIndex · leader lease · stale │ ├─────────────────────────────────────────────────────────────────┤ │ kvsm: replicated state machine — client sessions, idempotent │ │ apply (at-least-once delivery + dedup = exactly-once effect) │ ├─────────────────────────────────────────────────────────────────┤ │ raft (sans-IO): pre-vote elections · pipelined replication · │ │ snapshots + InstallSnapshot · membership changes · ReadIndex │ ├─────────────────────────────────────────────────────────────────┤ │ storage: LSM engine — WAL (CRC, torn-write recovery) · │ │ memtable · SSTables (bloom filters, block checksums) · │ │ leveled compaction · MVCC snapshots │ └─────────────────────────────────────────────────────────────────┘ ▲ production: tokio + real disks ▲ tests: deterministic │ (crates/server) │ simulator (crates/sim) ```
raft 核心是 **无 IO 的**:一个没有套接字、时钟 或线程的纯状态机。相同的代码既能在生产环境的 tokio 宿主下运行,也能在 控制所有调度决策的单线程模拟器下运行,这也正是下一节内容得以实现的原因。 ## 确定性模拟测试 `crates/sim` 在虚拟时间上运行整个集群(节点、磁盘、网络和客户端),由 一个带种子的 PRNG 驱动。每次运行都会注入:消息丢失、 重复、重排序、对称和不对称分区、伴随易失状态丢失的崩溃/重启, 以及强制进行快照传输的日志压缩压力。在运行期间,它会持续断言选举安全性 (每个任期最多 1 个领导者)、状态机安全性(在所有地方的相同索引处应用相同的条目) 以及日志匹配。最后,它会检查全集群的收敛情况,并针对完整的 客户端历史记录运行一个从头构建的 **Wing & Gong 线性一致性检查器**。 ![模拟器循环:单个种子驱动一个 PRNG,后者为虚拟时间事件队列提供数据; 每个弹出的事件步进无 IO 的 raft 核心,其输出(发送前持久化) 经过一个注入故障并加入新事件的模拟网络和磁盘。选举安全性、状态机安全性和日志 匹配在每一步都会被断言;Wing & Gong 线性一致性检查在结束时 运行。](https://static.pigsec.cn/wp-content/uploads/repos/cas/b4/b4aaced81fc7f02a9e286ee4e10a6ae9ed5b99d34ed7deadc9d4122b218e309c.svg) 一次运行是其种子的纯函数。一次失败会打印出种子;该种子 可以精确到虚拟微秒地重放该 bug: ``` cargo run --release -p sim --bin simulate -- --seeds 1000 # sweep cargo run --release -p sim --bin simulate -- --start 6 --seeds 1 # replay ``` 吞吐量:在上述基准测试中约为每分钟 5,300 个种子(每个种子大约有 21k 个 模拟事件)。每一个曾经暴露过 bug 的种子都保存在 `crates/sim/regressions/seeds.txt` 中,并在 CI 中重放。 **[docs/BUGS.md](docs/BUGS.md)** 记录了这个技术栈捕获的八个真实 bug, 包含根本原因和重现方法,其中包括一个 **状态机安全性违规** (种子 19519):一个描述日志 *前缀* 的 `InstallSnapshot` 截断了 follower 已确认的条目,随后一个已提交的条目被合法的选举 覆盖。它需要五个罕见条件在某一瞬间同时满足,并且能够逃过任何数量的常规测试。 ## 实测性能 单机上的 3 节点集群:AMD Ryzen 5 2600 (6C/12T)、48 GB、三星 970 EVO Plus NVMe、Windows 11。本地 TCP,100 字节值,zipfian keys (θ=0.99),20 ms ticks,raft 日志上每次提交执行 fsync。所有三个节点 *以及* 负载生成器共享这六个核心,因此这些数据是保守的。 下面每次读取/混合工作负载都是 **3 次测量运行的中位数**,每次 运行前都有一次被丢弃的预热过程;load 行是整个测试套件中 从头开始进行的 10 万次加载运行的中位数(冷插入没有可丢弃的 预热)。数据目录会被擦除,并且运行之间会重启集群。 冷存储测量 SSTable 遍历,热存储测量页面缓存,未擦除的存储 会累积压缩债务,因此该协议控制了这三者。在大多数工作负载上, 每次运行之间的波动在几个百分点以内(在 100% 读取的行上最严格, ≤2%);工作负载 B 的波动最大,约为 15%。可以使用 `chaos/local-cluster.sh start 20` + `loadgen` 复现(具体命令见 *快速开始*)。 | 工作负载 (YCSB) | 一致性 | 吞吐量 | p50 | p99 | p999 | |---|---|---|---|---|---| | Load (100% insert, 16 conns) | — | 1,318 ops/s | 10.7 ms | 18.6 ms | 345 ms | | A (50% read / 50% update, 32 conns) | linearizable reads | 2,308 ops/s | 11.0 ms (reads) | 20.6 ms | 756 ms | | B (95% read / 5% update, 32 conns) | leader-lease reads | 16,187 ops/s | 0.69 ms (reads) | 2.73 ms | 4.0 ms | | C (100% read, 32 conns) | stale reads | 13,317 ops/s | 2.35 ms | 3.20 ms | 4.36 ms | 该表并**没有**对一致性的代价进行标价,这一点值得一提:每一 行都同时改变了工作负载 *和* 一致性模式。尽管 B 是更严格的模式,但 B 的读取 (p50 为 0.69 ms)击败了 C(p50 为 2.35 ms)。 B 的 5% 更新使得 zipfian 热点 key 保留在 memtable 中,因此它的 读取是从内存中提供的,而 C 则需要遍历 SSTables。这是一种披着 一致性外衣的工作负载效应,也是每个以这种方式展示的 YCSB 表格中 都会存在的陷阱。 *确实* 对一致性进行了标价的比较则保持工作负载不变 (100% 读取,6 万次操作,32 个连接),并且只改变 `--consistency`: | 一致性 | 吞吐量 | p50 | 代价 | |---|---|---|---| | stale | 13,517 ops/s | 2.34 ms | 1 RTT | | lease | 13,057 ops/s | 2.42 ms | 1 RTT + 本地租约检查 | | linearizable | 9,043 ops/s | 3.49 ms | quorum RTT | **linearizable 读取的代价约为 1.5 倍**,这就是 ReadIndex quorum 往返, 也是严格读取的真正代价。**stale 和 lease 难以区分** (相差约 3%,在运行间方差范围内);这是客户端带来的结果, 并非偶然,如下所述。从 stale 读取到写入之间的差距 约为 10 倍(1,318 → 13,517 ops/s),这个差距就是 fsync 加上复制。 诚实地衡量它才是关键。 ![linearizable 读取的代价:在工作负载固定为 100% 读取的情况下,stale 和 lease 读取的吞吐量在 13,000–13,500 ops/s 之间,而 linearizable 读取的代价高出约 1.5 倍, 为 9,043 ops/s。](https://static.pigsec.cn/wp-content/uploads/repos/cas/30/30786fdc7053386cc12df86b7570759d050b438832bddf301e00143cd3e79dfb.png) 组提交,即事件循环在执行一次 fsync 之前排空所有排队的提议 (`crates/server/src/core.rs`),是使得写入路径 变得可承受的原因。这里 **没有** 引用分组与未分组的加速比: 这需要构建一个禁用了组提交的第二个版本,但我还没有运行过, 因此没有诚实的数字可以汇报。 ### Stale 读取不会分散负载 `KvClient` 从 `addrs[0]` 开始,并且只有在节点回复 `NotLeader` 时才会移动。stale 读取永远不会触发该回复,因为每个副本都会 提供读取服务,因此 **每个客户端在其整个生命周期内都固定在节点 1 上**,并且 stale 读取永远不会在集群中分散开来。实测结果:工作负载 C 针对 所有三个节点时达到 13,317 ops/s;仅针对节点 1 时,达到 13,268 ops/s, 数字完全相同。它本来就只在使用节点 1。 因此,下表中的“到任意节点的 1 RTT”描述的是协议,而不是这个 客户端的行为,这也是为什么在这里 stale 和 lease 读取的代价相同的原因: 它们遍历了完全相同的路径,且租约检查是本地的。修复它 (在副本之间轮询 stale 读取的目标)是 *我下一步会构建的内容* 下的 第一项。 ### 测量真实的故障转移 在 32 个连接的 lease-read 工作负载下,对领导者执行 `kill -9`: ![负载下的故障转移:吞吐量保持在约 16k ops/s,在领导者于 t≈8.5s 被杀死后 降至零约两秒钟,同时集群选举出新的领导者且客户端重新路由,然后 在 t=12s 时完全恢复到约 18k ops/s,30 万次操作中丢失了 0 次。](https://static.pigsec.cn/wp-content/uploads/repos/cas/12/12dc97f117f3d99694a524d45ad52f735d288b6b9d3835d805a53d968b6aa240.png) 连续 10 次杀死领导者时的客户端观察到的故障转移 (`chaos/kill-leader.sh 10 20`):**最小 2,103 ms / 中位数 2,141 ms / 最大 2,204 ms**。分解(20 ms ticks → 200–400 ms 选举超时):选举 本身不到一秒钟;客户端实际等待的约 2.1 s 主要耗费在领导者 *发现* 上:幸存的 followers 会不断提示已死亡的领导者,直到 新领导者的第一次心跳,因此客户端会进行带退避的乒乓操作。Server 恢复 和客户端收敛是两个不同的指标; 大多数系统只会引用前者。 在杀死进程后约 3.5 秒内吞吐量完全恢复,期间有约 2 秒处于零。 在交接期间进行中的写入会被延迟,而不会丢失:**30 万次操作中有 0 次 失败**,并且更新操作的尾部延迟保持紧凑(p999 = 99 ms)。少数 在选举过程中被捕获的写入只是简单地等待了约 2 秒的停机。 上面的两次测量都使用了混沌测试工具的客户端重试设置(短 退避,最多尝试 60 次:`chaos/kill-leader.sh` 中的 `RAFTKV_BACKOFF_*` / `RAFTKV_MAX_ATTEMPTS`),因此客户端会撑过选举,而不是 放弃。因此,这里的“故障转移”衡量的是集群的恢复时间,而不是 客户端放弃的行为;如果使用默认重试次数,一个操作可能会报错,而不是 等待。 ## 精确的一致性保证 | 模式 | 保证 | 代价 | 允许的异常 | |---|---|---|---| | 写入 / CAS | Linearizable | quorum + 2 fsyncs | 无 | | `linearizable` read | Linearizable(ReadIndex:heartbeat-quorum 确认领导权,在 ≥ 已确认的提交索引处提供服务) | quorum RTT,无 fsync,无日志条目 | 无,在任何时钟行为下 | | `lease` read | **如果时钟漂移是有界的**则为 Linearizable(heartbeat-quorum 租约) | 到领导者的 1 RTT | 如果暂停/漂移的领导者在其租约真正过期后提供服务,则会出现 stale read | | `stale` read | 已提交但可能陈旧的数据(绝不包含未提交的数据) | 到任意节点的 1 RTT(但请参见 *Stale 读取不会分散负载*;此客户端总是询问节点 1) | 任意程度的陈旧;没有 read-your-writes | 重试的写入是安全的:每个客户端会话都带有一个序列号, 并且复制的状态机维护着一个按会话划分的去重表(包含在快照中),因此 在超时后重试会从缓存中给出响应,永远不会 重新执行。这是至少一次交付加上幂等应用,是“精确一次” 的诚实构造(在异步网络中,真正的精确一次 *交付* 是不存在 的)。 ## 目录结构 ``` crates/storage LSM engine: wal, memtable, sstable (+bloom), manifest, merge iterators, leveled compaction, MVCC snapshots crates/raft sans-IO consensus core + scenario test battery crates/kvsm replicated KV state machine w/ sessions crates/proto wire types + length-prefixed bincode framing crates/server tokio host: kvd binary, raft-log persistence on the LSM engine, group commit, consistency modes, Prometheus metrics crates/client smart client library + kvctl CLI crates/bench YCSB-style workloads (zipfian, A–F) + loadgen binary crates/sim deterministic simulator, WGL checker, seed regressions chaos/ local-cluster/kill-leader scripts, docker-compose + netem docs/BUGS.md eight bugs, root causes, reproducing seeds ``` ## 快速开始 ``` cargo test --workspace # 47 tests incl. 30-seed sim sweep cargo build --release chaos/local-cluster.sh start # 3 nodes on localhost target/release/kvctl --cluster 127.0.0.1:6001,127.0.0.1:6002,127.0.0.1:6003 \ put hello world target/release/kvctl --cluster ... get hello --consistency lease chaos/kill-leader.sh 5 # failover distribution, 5 kills ``` ## 设计权衡 - **在 ack 之前总是对 raft 日志执行 fsync。** 丢失一个已确认的条目或 重新投票都会破坏安全性;WAL 将整个事件排空的内容批处理为 一次 fsync(组提交)以使其可承受。 - **状态机 Db 从不执行 fsync。** 它是派生状态,在重启时从 raft 快照 + 日志重放中重建。持久性消耗只支付一次,而不是两次。 - **背压导致的消息丢失。** 出站对等队列在满时会丢弃; raft 的重传就是流控制。死掉的 peer 绝不能阻塞 事件循环。 - **使用长度前缀的 bincode 而不是 gRPC。** 这使得通信格式保持 从零开始的状态,并且构建过程仅使用 `cargo build`;RPC 语义的 设计使得在 proto crate 边界处可以替换为 tonic。 - **单 server 成员变更**(不是 joint consensus):一次只能有一个 未提交的配置变更,配置在追加时生效;这是 更简单的协议,其边缘情况实际上是可测试的。joint consensus 是多变更的推广。 ## 我下一步会构建的内容 在副本之间分散 stale 读取(今天每个客户端都固定在 `addrs[0]`, 因此唯一 *可以* 水平扩展的模式却没有这么做;修复方法是 每个客户端使用一个起始偏移量加上在 stale 路径上轮询);多 group 分片,具有分片映射配置服务和迁移切换;事务 (基于 raft groups 的 2PC,然后是 percolator 风格);基于租约的时钟 (用测量的时钟误差限制 lease-read 异常);存储路径上的 io_uring; 模拟器中覆盖率引导的调度探索。
标签:LSM存储引擎, Raft算法, Rust, 分布式一致性, 分布式数据库, 可视化界面, 测试模拟器, 网络流量审计, 通知系统, 键值存储