【免费下载链接】toydb
Distributed SQL database in Rust, written as an educational project
toyDB 是一个以教学为目的、用 Rust 编写的分布式 SQL 数据库,其目标不是性能与规模,而是用尽量简单的设计完整呈现"分布式数据库到底是怎么造出来的"。本文以仓库内 docs/architecture 架构文档为主干,结合 源码目录 与 配置文件,自底向上逐层讲解 toyDB 的架构:存储引擎、键值编码、MVCC 事务、Raft 共识、SQL 引擎、服务端与客户端,读完即可掌握一个可运行的分布式数据库全链路实现思路,并能直接在本仓库中对照源码、配置与测试脚本进行学习与验证。
架构总览:一个运行在集群上的复制状态机
toyDB 由一组节点(node)组成集群,每个节点都执行 SQL 事务,作用于同一个"复制状态机"(replicated state machine)。客户端可以连接集群中的任意一个节点并提交 SQL 语句;当集群中只有少数派(minority)节点崩溃或断连时集群仍然可用,但一旦多数派节点故障,整个集群就会停止服务——这是分布式共识协议为保证一致性与持久性必须付出的代价。
整体架构可以归纳为以下性质(来自 overview 文档):
- 分布式:运行在由多个节点组成的集群上。
- 高可用:容忍少数派节点故障。
- SQL 兼容:正确支持绝大多数常用 SQL 特性。
- 强一致:已提交的写入对所有读方立即可见(满足线性一致性 linearizability)。
- 事务性:提供 ACID 事务——原子性(一组写入作为一个原子单元提交)、一致性(约束与引用完整性始终被强制)、隔离性(并发事务互不影响,采用快照隔离)、持久性(已提交的写入永不丢失)。
同时文档也毫不避讳地列出了 toyDB 的"非目标",这些恰恰是真实系统中复杂度的重要来源:
- 不可扩展:每个节点都保存全量数据,读写只在一个节点上执行。
- 不可靠:只处理崩溃故障,不处理网络分区、节点卡死等场景。
- 不追求性能:数据处理慢且完全未优化。
- 不高效:整表加载进内存,无压缩、无垃圾回收。
- 功能不完整:仅实现基础 SQL 功能。
- 不向后兼容:数据格式与协议变更会破坏旧数据库。
- 不灵活:运行期间不能增删节点,节点加入集群耗时很长。
- 不安全:没有认证、授权与加密。
核心组件
从内部看,toyDB 由五个主要组件构成:
- 存储引擎(Storage Engine):在磁盘上存储数据并管理事务。
- Raft 共识引擎(Raft Consensus Engine):复制数据、协调集群节点。
- SQL 引擎(SQL Engine):组织 SQL 数据、管理 SQL 会话、执行 SQL 语句。
- 服务端(Server):负责网络通信,既连接 SQL 客户端,也连接 Raft 节点。
- 客户端(Client):提供 SQL 用户界面并与服务端通信。
下图展示了单个 toyDB 节点的内部结构:
下文将按照文档的编排顺序,从最底层的存储引擎开始,逐层向上展开。
存储引擎:嵌入式键值存储
toyDB 使用嵌入式 键值存储 存放数据,键和值都是任意二进制字节串。存储引擎不关心键值内容的含义——SQL 数据模型(表、行)如何映射到键值结构,是更上层的问题。
存储引擎提供单键的 set/get/delete 操作,本身不支持事务(事务构建在它之上)。键以有序方式存储,因此支持区间扫描(range scan):可以迭代两个键之间的所有键值对,也可以按指定键前缀扫描。这一能力被系统其他组件广泛依赖:扫描某张 SQL 表的全部行、扫描 MVCC 键的多个版本、扫描 Raft 日志尾部等。
存储引擎是可插拔的:存在多种实现,用户通过配置文件选择使用哪一种。所有实现都遵循storage::Enginetrait(见 src/storage/engine.rs)。
Memory存储引擎
最简单的实现是storage::Memory,它直接用 Rust 标准库的BTreeMap把数据放在内存里,不落盘,主要用于测试(见 src/storage/memory.rs)。
BitCask存储引擎
主存储引擎是storage::BitCask,它是 Riak 数据库所用 BitCask)。
toyDB 的 BitCask 实现使用单一追加写日志文件存储数据:
- 写入键值对:直接追加到文件末尾。
- 删除键:追加一个特殊的墓碑(tombstone)值。
- 读取键:使用文件中该键的最后一条记录。
每个键值对的文件格式非常简单:
- 键长度,大端序
u32(4 字节)。 - 值长度,大端序
i32(4 字节);如果是墓碑则为-1。 - 键的二进制内容(n 字节)。
- 值的二进制内容(n 字节)。
例如键值对foo=bar写成十六进制就是:
keylen valuelen key value 00000003 00000003 666f6f 626172因为数据文件本身就是追加式日志,所以不需要单独的预写日志(WAL)来做崩溃恢复——数据文件就是 WAL。
为了快速读取,BitCask 在内存中维护一个KeyDir索引,把键映射到该键最新值在文件中的位置。因此所有键都必须能装进内存。这个索引在打开文件时通过完整扫描日志文件重建;写入键时追加记录并更新KeyDir;删除键时同样更新;读取时先在KeyDir中查找文件位置再读取。
KeyDir内部同样使用标准库BTreeMap保持键有序,从而支持按区间有序迭代并逐个加载键值。
随着键不断被更新和删除,日志文件中会积累大量旧版本。为此 BitCask 在启动时执行压缩(compaction):把所有仍存活键值对的最新值写入一个新文件,再替换旧文件;键按有序顺序写出,让后续扫描更快。仓库提供了对应的测试脚本 src/storage/testscripts/bitcask/compact、compact_open、log、status。
键值编码:Bincode 值与 Keycode 键
键值存储使用二进制Vec<u8>作为键和值,因此需要一套编码方案在内存中的 Rust 数据结构与磁盘二进制数据之间转换。encoding模块(见 src/encoding/mod.rs)为键和值提供了两套不同的编码方案。
Bincode值编码
值使用 Rust 生态中常见的 Bincode 二进制编码。Bincode 的便利之处在于可以轻松编码任意 Rust 数据类型;当然也可以换成 JSON、Protobuf、MessagePack 等。
为了让所有编解码使用一致的配置,src/encoding/bincode.rs 提供了辅助函数,统一使用bincode::config::standard()。Bincode 基于 Serde 框架,toyDB 还提供了encoding::Value辅助 trait,为值类型自动增加encode()/decode()方法:
#[derive(serde::Serialize, serde::Deserialize)] struct Dog { name: String, age: u8, good_boy: bool, } impl encoding::Value for Dog {} let pluto = Dog { name: "Pluto".into(), age: 4, good_boy: true }; let bytes = pluto.encode(); println!("{bytes:02x?}"); // 输出 [05, 50, 6c, 75, 74, 6f, 04, 01]: // // * 字符串 "Pluto" 的长度:05。 // * 字符串 "Pluto":50 6c 75 74 6f。 // * 年龄 4:04。 // * good_boy:01(true)。 let pluto = Dog::decode(&bytes)?; // 还原出 PlutoKeycode键编码
键不能像值那样随意编码。存储引擎依赖键的有序性来实现区间扫描,因此键编码必须保持字典序(lexicographical order):二进制字节序列的排序必须与原始值的排序一致。
为什么不能直接用 Bincode?考虑字符串 "house" 和 "key",字典序应该是 "house" 在前。但 Bincode 会把长度前缀写在字符串前面,导致 "key" 反而排到前面:
03 6b 65 79 ← 3 字节:key 05 68 6f 75 73 65 ← 5 字节:house数字同理:小端序会让很大的数排在很小的数前面,符号位会把正数排到负数前面,违背自然数排序;而值序列需要按元素逐个比较,例如 ("a", "xyz") 应排在 ("ab", "cd") 之前,因此也不能简单拼接字符串。
toyDB 在 src/encoding/keycode.rs 提供了一种保序编码,名为Keycode。与 Bincode 一样,它不自描述:二进制数据不携带类型信息,调用方必须提供解码目标类型;它只支持少量基本数据类型,且只需对同类型值排序。Keycode 实现为 Serde(反)序列化器,编码规则如下:
| 类型 | 编码规则 |
|---|---|
bool | false为00,true为01 |
u64 | 大端序二进制编码 |
i64 | 大端序二进制编码,但翻转符号位,使负数排在正数之前 |
f64 | 大端序 IEEE 754 编码,翻转符号位,负数再翻转全部位,以正确排序负数 |
Vec<u8> | 以00 00结尾终止,内部的00转义为00 ff以消除歧义 |
String | 同Vec<u8> |
Vec<T>、[T]、(T,) | 内部值的拼接 |
enum | 变体的数值索引(u8),随后是内部值(如有) |
与encoding::Value类似,还有encoding::Key辅助 trait。不同种类的键通常用枚举表示,例如同时存储汽车和电子游戏:
#[derive(serde::Serialize, serde::Deserialize)] enum Key { Car(String, String, u64), // make, model, year Game(String, u64, Platform), // name, year, platform } #[derive(serde::Serialize, serde::Deserialize)] enum Platform { PC, PS5, Switch, Xbox, } impl encoding::Key for Key {} let returnal = Key::Game("Returnal".into(), 2021, Platform::PS5); let bytes = returnal.encode(); println!("{bytes:02x?}"); // 输出 [01, 52, 65, 74, 75, 72, 6e, 61, 6c, 00, 00, 00, 00, 00, 00, 00, 00, 07, e5, 01]。 // // * Key::Game:01 // * Returnal:52 65 74 75 72 6e 61 6c 00 00 // * 2021:00 00 00 00 00 00 07 e5 // * Platform::PS5:01 let returnal = Key::decode(&bytes)?;由于键按元素顺序排序,可以借此做前缀扫描(例如查出 Returnal 发行过的所有平台)或区间扫描(例如查出 2010–2015 年间 Nissan Altima 的全部车型)。
MVCC 事务:快照隔离的实现
事务是把一组读写(例如对不同键的操作)作为一个整体提交的单元。以转账为例:从账户 A 转 100 到账户 B,就是一组读和写:
a = get(A) b = get(B) if a < 100: error("insufficient balance") set(A, a - 100) set(B, b + 100)toyDB 提供 ACID 事务:
- 原子性(Atomicity):所有写入在提交的那一刻作为一个原子单元同时生效,其他用户绝不会看到其中一部分。
- 一致性(Consistency):数据库约束(引用完整性、唯一性等)绝不被违反。
- 隔离性(Isolation):每个用户仿佛独占整个数据库。两个事务可能冲突,冲突方需要重试;一旦成功,用户就能确信操作没有受到任何干扰,从而消除竞态条件。
- 持久性(Durability):已提交的写入永不丢失(即使系统崩溃)。
为提供这些保证,toyDB 使用常见的多版本并发控制(MVCC)技术,实现在键值存储层 src/storage/mvcc.rs,底层使用一个storage::Engine作为实际数据存储。MVCC 提供的隔离级别是快照隔离(snapshot isolation):事务看到的是其开始那一刻的数据库快照,之后发生的任何变更对它不可见。
实现方式是保存键值对的历史版本。版本号是一个每次新事务都递增的编号;每个事务拥有唯一版本号,写入键值对时把版本号追加到键上(形如Key::Version(&[u8], Version),使用前面讲过的 Keycode 编码)。如果一个键已有旧版本,它带有不同的版本号后缀,因此会作为独立键存储;删除的键则是带特殊墓碑值的版本。此外还需要跟踪当前正在进行(未提交)的事务版本集合,称为活跃集(active set)。
由此可以概括 MVCC 协议的几条简单规则:
- 事务开始时:获取下一个可用版本号;对活跃集(其他未提交事务)拍快照;把自己的版本号加入活跃集。
- 事务读取键时:返回该键在自身版本及以下的最新版本;忽略高于自身版本的版本;忽略活跃集快照中的版本。
- 事务写入键时:查找高于自身版本的键版本,找到则报错;查找活跃集快照中的版本,找到则报错;以自身版本写入键值对。
- 事务提交时:把所有写入刷到磁盘;把自己从活跃集中移除。
关键魔法发生在事务把自己从活跃集移除这一时刻:这是一个单一原子操作,完成后它的所有写入立刻对新事务可见;而正在进行的事务仍然看不到这些写入(版本仍在它们的活跃集快照中,或版本号更新),因此彼此隔离。事务能看到自己的未提交写入而他人不能;若任何写入与其他事务冲突,则报错并需重试。此外这还支持时间旅行查询:选定一个版本号读取,就能查询数据库在过去任意时刻的状态。
文档还提及一些略去的细节:回滚需要跟踪写入并撤销(toyDB 用Key::TxnWrite(version, key)记录已写键);只读查询可避免分配新版本号;出于简洁性,MVCC 不做旧版本垃圾回收。完整的模块文档见 src/storage/mvcc.rs 顶部注释。
代码走读:使用 MVCC API
下面的示例展示了 MVCC 的使用方式——注意使用 API 时完全不需要接触版本号,这是 MVCC 的内部实现细节:
// 在文件 "toy.db" 中打开一个带 MVCC 支持的 BitCask 数据库。 let path = PathBuf::from("toy.db"); let db = MVCC::new(BitCask::new(path)?); // 开始一个新事务。 let txn = db.begin()?; // 读取键 "foo",并用 bincode 把二进制值解码为 u64。 let bytes = txn.get(b"foo")?.expect("foo not found"); let mut value: u64 = bincode::deserialize(&bytes)?; // 删除 "foo"。 txn.delete(b"foo")?; // 给 value 加 1,写回到键 "bar"。 value += 1; let bytes = bincode::serialize(&value); txn.set(b"bar", bytes)?; // 提交事务。 txn.commit()?;各步骤对应的实现:
MVCC::begin()调用Transaction::begin():获取存储在Key::NextVersion中的版本号并递增,再对Key::ActiveSet拍快照并把自己加入其中,返回一个封装了版本号与活跃集的Transaction对象,提供 get/set/delete 主 API。Transaction::get():由于多版本键以Key::Version(key, version)形式存储,而 Keycode 编码保证所有版本有序排列,因此可以从Key::Version(b"foo", self.version)到Key::Version(b"foo", 0)做反向区间扫描,返回第一个对自身可见的最新版本(忽略未来版本与活跃集)。Transaction::delete()/set():都走同一个write_version()方法,普通键值对用Some(value),删除墓碑用None。写入前先区间扫描检查是否存在对自身不可见的版本(即与其他并发事务冲突),无冲突后写入Key::Version(key, version),并写一条Key::TxnWrite记录以便回滚。Transaction::commit():只需把自己从Key::ActiveSet移除并清理Key::TxnWrite即可生效。注释特别说明:这里不需要主动刷盘,因为持久性由 Raft 日志提供,后面会讲到。
仓库提供了大量 MVCC 测试脚本,例如模拟两个并发用户修改一组银行账户的 src/storage/testscripts/mvcc/bank,以及各类异常场景测试(脏读、脏写、模糊读、丢失更新、幻读、读偏斜、写偏斜等,见 src/storage/testscripts/mvcc 目录),还支持begin_as_of(指定版本开始的时间旅行)、begin_readonly(只读事务)等。
Raft 共识:复制与可用性的来源
Raft 是一个分布式共识协议,以一致且持久的方式把数据复制到集群各节点。它本质上与 Paxos、Viewstamped Replication 是同一类协议,但被设计得更简单、更可读、更实用,并被 CockroachDB、TiDB、etcd、Consul 等广泛使用。
简言之:Raft 选举出一个领导者(leader)节点,由它协调写入并复制到跟随者(follower);一旦多数派(>50%)节点确认了某条写入,该写入即被视为持久提交。领导者通常也承担读请求,因为它总是持有最新数据,因而强一致。集群必须保持多数派节点(法定人数 quorum)在线连通才能提供服务,否则不提交写入,以保障一致性与持久性。由于集群中只可能存在一个多数派,这从机制上杜绝了"脑裂"(split brain)——例如网络分区期间不可能同时存在两个活跃领导者并写入冲突数据。
Raft 领导者把写入追加到一条有序命令日志,复制给跟随者;一旦多数派把日志复制到某条记录,该日志前缀即被提交(committed),然后应用到状态机。这保证了所有节点以相同顺序应用相同命令、最终到达相同状态(假设命令是确定性的)。Raft 本身不关心状态机与命令是什么——在 toyDB 中,状态机就是存放在 MVCC 键值存储里的 SQL 表与行。文档中同时点明:Raft 解决的是复制与可用性,不解决扩展性(所有读写都经领导者、每节点存全量数据);真实系统通常靠分片(sharding)把大数据集拆到多个独立 Raft 集群来扩展,这超出 toyDB 范围。toyDB 只实现 Raft 的最小必要部分,省略了状态快照、日志截断、领导者租约等论文中的优化。
协议实现在 src/raft 模块,其模块文档(src/raft/mod.rs 顶部)对整体设计有详细说明;src/raft/testscripts/node 目录还提供了覆盖各种场景的完整测试脚本集。
日志存储
Raft 复制的有序命令日志由raft::Entry组成(见 src/raft/log.rs):index指明日志中的位置,command是应用到状态机的二进制命令,term标识命令被提出时的领导任期(每次举行新的领导者选举都会开启新任期)。日志条目由领导者追加并复制给跟随者;一旦被法定人数确认,该 index 之前的日志即被提交且永不改变,未提交的条目在领导者变更后可能被替换或删除。
raft::Log把日志条目存放在storage::Engine键值存储中,同时保存当前任期、投票、提交索引等元数据。提供的接口包括:Log::append(追加单条,通常由领导者复制新写入时使用)、Log::splice(批量追加,常用于向跟随者复制,也允许替换未提交条目)、Log::commit(标记已提交条目,使其不可变并可被状态机应用)、Log::get/Log::scan(读取单条或按区间迭代)。
状态机接口
Raft 不关心日志命令的内容,也不关心状态机如何处理它们,只负责把raft::Entry从日志取出交给状态机。状态机由raft::Statetrait 表示(见 src/raft/state.rs):Raft 通过State::get_applied_index询问最后应用的条目,通过State::apply喂入新提交的条目,还通过State::read提供读取(后文详述)。
状态机不必在每次状态转换后把状态刷到持久存储:节点崩溃时允许状态机回退,通过重放未应用的日志条目来追平。也可以实现纯内存状态机——toyDB 确实允许用Memory存储引擎运行状态机。关键约束是状态机必须确定性:相同命令以相同顺序应用必须得到相同状态,因此命令不能读取当前时间或生成随机数(这些值必须包含在命令里);非确定性错误(如 IO 错误)必须中止命令应用(toyDB 的做法是直接 panic 崩溃节点)。
节点角色
Raft 节点有三种角色:
- Leader(领导者):向跟随者复制写入、服务客户端请求。
- Follower(跟随者):从领导者复制写入。
- Candidate(候选人):竞选领导权。
在 toyDB 中,节点用raft::Node枚举表示(见 src/raft/node.rs),它包装了raft::RawNode<Role>类型,后者以角色为泛型参数,利用 Rust 的类型状态模式(typestate pattern)在编译期强制执行状态转换与不变式——例如只有RawNode<Candidate>拥有into_leader()方法,因为只有候选人在赢得选举后才能转为领导者。RawNode::role字段保存角色专属状态,各角色状态是实现Role标记 trait 的结构体。
下图来自 Raft 论文,概括了三种角色及其转换关系(选举的细节见下一节):
节点接口与通信
raft::Node有两个驱动节点的主方法:tick()与step(),两者都"消费"当前节点并返回新节点(可能角色已变化)。
tick():推进一个逻辑 tick,用于度量时间流逝,例如触发选举超时或周期性领导者心跳。toyDB 的 tick 间隔为 100 毫秒(raft::TICK_INTERVAL),会以该频率对节点调用tick()。step():处理来自其他节点或客户端的一条入站消息。
出站消息通过RawNode::tx通道发送。节点用启动时给定的唯一节点 ID 标识;消息被包装在raft::Envelope中(见 src/raft/message.rs),指定发送者与接收者,内含raft::Message枚举——它编码了完整的 Raft 消息协议。Raft 不要求可靠消息投递(消息可随时丢失或乱序),不过 toyDB 使用 TCP 提供了更强的投递保证。
这是一个完全同步且确定性的模型:给定初始状态和相同的调用序列,结果必然相同,非常适合测试与理解。服务端部分会看到 toyDB 如何在独立线程上驱动节点、提供网络传输并按固定间隔 tick。
领导者选举与任期
稳定状态下 Raft 只有一个领导者向跟随者复制写入;但要达到稳定状态,必须先选举出领导者——这正是一致性协议中最微妙的部分。
Raft 把时间划分为任期(term):任期是从 1 开始单调递增的编号,一个任期内只可能有一个领导者(选举失败则没有),任期永不回退;复制的命令归属于其提出时所在的任期。
以新建一个 3 节点空集群为例走一遍选举:
- 节点通过
Node::new()初始化,给定空的raft::Log与raft::State,任期 0,角色为Follower且没有领导者。 - 随后每 100 ms 调用
tick(),直到到达election_timeout。new()会把election_timeout设为随机值(范围ELECTION_TIMEOUT_RANGE为 10–20 个 tick,即 1–2 秒)。若所有节点超时相同,它们很可能同时竞选导致选举平局——Raft 用随机化选举超时避免这种情况。 - 某个节点到达
election_timeout后转为Candidate,把任期递增为 1,给自己投票,并向所有对等节点发送Message::Campaign请求选票。任期不能回退、每任期每节点只能投一票(跨重启也如此),所以二者都通过Log::set_term_vote()持久化。 - 另两个节点(仍为
Follower)收到Message::Campaign后,先把任期升到 1(因为这是更新任期),检查自己尚未在任期 1 投过票后同意投票,持久化选票并回复Message::CampaignResponse { vote: true }。它们还会检查候选人的日志长度不小于自己的日志——空日志下显然成立——这是保证新领导者拥有全部已提交条目的必要条件(Raft 论文 5.4.1 节)。 - 候选人收到
Message::CampaignResponse后记录各节点选票;一旦获得法定人数(3 节点中 2 票,含自己),就成为任期 1 的领导者。 - 新领导者向所有对等节点发送
Message::Heartbeat宣告领导地位,同时向日志追加一条空条目并复制(原因见论文 5.4.2 节)。 - 其他节点收到心跳后,成为新领导者在任期 1 下的跟随者。此后领导者每 4 个 tick(
HEARTBEAT_INTERVAL)发送一次周期性Message::Heartbeat宣示领导权;跟随者记录最后一次收到领导者消息(含心跳)的时间,若一个选举超时内没有收到(如领导者崩溃或网络分区),就发起新一轮选举。
完整过程见测试脚本 src/raft/testscripts/node/election,同一目录下还有选举平局(election_tie)、争夺选举(election_contested)等更多场景。
客户端请求与转发
领导者选出后即可提交读写请求:把Message::ClientRequest以本地节点 ID 和唯一请求 ID(toyDB 用 UUIDv4)step 进节点,等待带相同 ID 的出站响应消息。请求与响应本身是任意二进制数据,由状态机解释(文档中用一个虚构协议举例:Request::Write("key=value")→Response::Write("ok"),Request::Read("key")→Response::Read("value"))。
读写请求的根本区别在于:写请求经 Raft 复制并在所有节点上执行;读请求只在本领导者上执行、不追加日志。理论上也可在跟随者上执行读以做负载均衡,但那样读会变成最终一致、破坏线性一致性,因此 toyDB 只在领导者上执行读。
如果请求提交给了跟随者,它会被转发给领导者,响应再转发回客户端(用发送者/接收者节点 ID 区分——本地客户端总是使用本地节点 ID)。为保持简单,请求提交给候选人时会被Error::Abort取消,跟随者转为候选人或发现新领导者时同理——本来可以保存请求并重定向到新领导者,但 toyDB 选择简单处理、请客户端重试。
写入复制与应用
领导者收到写请求后,把命令作为日志条目提出(propose)以复制给跟随者,并在writes中记录这个在途写入及其日志条目索引,以便条目提交并应用后把命令结果返回给客户端。
提出命令时,领导者把它追加到自己的日志,并向每个跟随者发送Message::Append复制到它们的日志。稳定状态下Message::Append只包含刚追加的单条日志条目。但跟随者可能落后(如崩溃恢复后),或日志与领导者产生分歧(如网络分区后陈旧领导者的未成功提案)。为此领导者以raft::Progress跟踪每个跟随者的复制进度。
Message::Append包含新条目紧前一条的base_index与base_term:若跟随者日志也含有该 index 与 term 的条目,则它的日志与领导者日志直到该条目处保证一致(论文 5.3 节),可安全追加新条目,并返回Message::AppendResponse确认,报告其日志与领导者匹配到match_index。领导者收到响应后更新对该跟随者的match_index视图。
一旦法定人数(3 节点中 2 个,含领导者)的日志都有了该条目,领导者即可提交该条目并应用到状态机,同时从writes中查找在途写请求,把命令结果作为Message::ClientResponse发回客户端。领导者还会通过下一次心跳把新的提交索引传播给跟随者,让它们也能把挂起的日志条目应用到各自状态机——这并非严格必需(读只在领导者执行,且节点在成为领导者前必须应用完挂起条目),但可避免跟随者应用进度落后太多。
相关测试脚本包括 src/raft/testscripts/node/append 与 src/raft/testscripts/node/heartbeat_commits_follower,以及大量其他场景。
读处理:线性一致读
如前所述,线性一致(强一致)读必须在领导者上执行。但这还不够:在网络分区等情况下,一个节点可能自以为仍是领导者,而事实上别处已在更新任期内选出了新领导者并执行了写入。
为处理这种情况,领导者每次读都必须确认自己仍是领导者:向跟随者发送包含读序号(read sequence number)的Message::Read,只有当法定人数确认它仍是领导者时才能执行读。这引入额外的一次网络往返,显然不高效——真实系统常改用领导者租约(leader lease,见 Raft 论文的姊妹篇《博士论文》6.4.1 节)——但对 toyDB 来说完全够用。
流程为:领导者收到读请求 → 递增读序号 → 把挂起读请求存入reads→ 向所有跟随者发送Message::Read;跟随者若消息来自当前领导者则回复Message::ReadResponse(过期任期的消息被忽略);领导者收到ReadResponse后记入对等节点的Progress,一旦法定人数确认了该序号,就执行读。
至此,我们就拥有一个由 Raft 管理的状态机:写入被复制、读保持线性一致。
SQL 引擎:从 SQL 文本到存储的流水线
SQL 引擎是对外的核心数据库接口,底层用键值存储存放数据、用 MVCC 提供事务、用 Raft 提供复制。整个引擎由多个组件构成一条流水线:
Client → Session → Lexer → Parser → Planner → Optimizer → Executor → Storage
实现位于 src/sql 模块,各子模块与架构文档一一对应:
- 数据模型(sql-data.md):数据类型、Schema(表/列/索引定义)、表达式。
- 存储(sql-storage.md):把表与行映射到键/值结构、Schema 目录(catalog)、行存储与事务。
- Raft 复制(sql-raft.md):SQL 写入如何经 Raft 复制。
- 解析(sql-parser.md):词法分析器(lexer)、抽象语法树(AST)、递归下降解析器(parser),对应 src/sql/parser。
- 规划(sql-planner.md):执行计划、作用域与名称解析、规划器,对应 src/sql/planner。
- 优化(sql-optimizer.md):常量折叠、谓词下推、索引查找、哈希连接、短路求值等优化规则,对应 src/sql/planner/optimizer.rs。
- 执行(sql-execution.md):计划执行器与会话管理,对应 src/sql/execution。
SQL 引擎整体由 src/sql/testscripts 下的测试脚本测试:这些脚本通常以原始 SQL 字符串为输入,在内存存储引擎上执行,并输出结果以及查询计划、存储操作、二进制键/值数据等中间状态。测试脚本按主题分为表达式(expressions)、优化器(optimizers)、查询(queries)、Schema(schema)、事务(transactions)与写入(writes)等类别。
服务端:把所有组件串起来
toydb::Server(见 src/server.rs)把前面各个组件绑在一起:它包装一个内部的 Raft 节点raft::Node(后者管理 SQL 状态机),负责在 Raft 节点、Raft 对等节点与 SQL 客户端之间路由网络流量。
协议方面,服务端使用前面讲过的Bincode 编码、走 TCP 连接,不需要额外的分帧,因为 Bincode 根据解码目标类型就能知道每个消息需要读多少字节。服务端刻意不用 async Rust(如 Tokio),改用普通 OS 线程——异步 Rust 会显著增加代码复杂度、掩盖核心概念,而效率收益对 toyDB 毫无意义;线程间用 Crossbeam 通道传递消息。
主循环Server::serve()监听入站 TCP 连接:Raft 对等节点与 SQL 客户端使用各自端口(配置文件中的listen_raft/listen_sql),并为每条连接派生线程处理。各部分分工如下:
- Raft 路由(
Server::raft_route()):服务端的心脏。它定期通过raft::Node::tick()驱动 Raft 节点、通过raft::Node::step()把来自 Raft 对等节点的入站消息 step 进节点,并把出站消息发给对等节点;同时接收来自sql::engine::RaftSQL 引擎的入站 Raft 客户端请求,step 进节点,并把节点产出的响应送回相应客户端。 - 对等节点收发:启动时为每个 Raft 对等节点派生一个
raft_send_peer()线程,持续尝试通过 TCP 连接对等节点,把通道中的出站raft::Envelope(raft::Message)用 Bincode 写入连接;raft_accept()持续监听入站 Raft TCP 连接,接受后派生raft_receive_peer()线程,读取 Bincode 编码的消息并经通道送往raft_route()。 - SQL 服务:使用
toydb::Request/toydb::Response枚举作为客户端协议(同样 Bincode 编码走 TCP)。主要请求类型是Request::Execute,它对一个sql::execution::Session执行 SQL 语句并返回sql::execution::StatementResult。服务端建立sql::engine::RaftSQL 引擎(用 Crossbeam 通道把raft::Request发给raft_route()再送入本地 Raft 节点),派生sql_accept()线程监听入站 SQL 客户端连接;接受连接后为该客户端建立新会话sql::execution::Session,派生sql_session()线程持续读取Request消息、执行并返回Response。
toydb二进制与配置文件
toydb可执行程序(src/bin/toydb.rs)是toydb::Server的薄封装:先用 clap 解析命令行并读取toydb.yaml配置文件,再初始化 Raft 日志存储与 SQL 状态机,最后启动toydb::Server。
仓库自带的示例配置 config/toydb.yaml 展示了全部配置项:
# 节点 ID(集群内必须唯一),以及 peer ID 到 Raft 地址的映射(单节点为空)。 id: 1 peers: {} # 监听 SQL 与 Raft 连接的地址。 listen_sql: localhost:9601 listen_raft: localhost:9701 # 日志级别。合法值为 DEBUG、INFO、WARN、ERROR。 log_level: INFO # 节点数据目录。Raft 日志存储在文件 "raft",SQL 数据库存储在 "sql"。 data_dir: data # 用于 Raft 日志与 SQL 数据库的存储引擎。 # # * bitcask(默认):追加写日志型存储。 # * memory:使用 Rust 标准库 BTreeMap 的内存存储。 storage_raft: bitcask storage_sql: bitcask # 是否对写入执行 fsync。关闭可获得好得多的写性能,但主机崩溃时可能丢失 # 数据并违反 Raft 保证。仅影响 Raft 日志写入(SQL 状态机从不 fsync, # 因为它可以由 Raft 日志重建)。 fsync: true # 触发 Bitcask 日志压缩的最低垃圾比例与字节数(在节点启动时执行)。 compact_threshold: 0.2 compact_min_bytes: 1000000其中值得注意的设计点(配置注释即源码级说明):
fsync只影响 Raft 日志写入,SQL 状态机从不 fsync——因为它可以从 Raft 日志完全重建,这是"状态机可以回退、靠重放日志追平"这一 Raft 设计在工程上的直接体现;compact_threshold与compact_min_bytes共同决定 BitCask 启动压缩的触发条件:垃圾占比达到阈值且垃圾字节数达到下限才触发压缩。
仓库还提供了可直接运行的 5 节点集群示例:每个节点目录(cluster/toydb1、cluster/toydb2 等)各含一个toydb.yaml,例如节点 1:
id: 1 data_dir: toydb1/data listen_sql: localhost:9601 listen_raft: localhost:9701 peers: '2': localhost:9702 '3': localhost:9703 '4': localhost:9704 '5': localhost:9705cluster/run.sh 会以 release 优化构建toydb二进制,然后后台启动节点 1–5(SQL 端口 9601–9605),并将每个节点的输出加上节点 ID 前缀以便区分:
cargo build --release --bin toydb for ID in 1 2 3 4 5; do (cargo run -q --release -- -c toydb$ID/toydb.yaml 2>&1 | sed -e "s/\(.*\)/toydb$ID \1/g") & done要连接节点 1(端口 9601)执行 SQL,运行:
cargo run --release --bin toysql客户端:客户端库与toysqlREPL
客户端位于 src/client.rs,使用与服务端相同的 Bincode 协议,发送toydb::Request、接收toydb::Response。
客户端库
主客户端库toydb::Client用于与 toyDB 服务端通信:初始化时通过 TCP 连接服务端,从而建立一个 SQL 会话;随后可以发送 Bincode 编码的toydb::Request并接收toydb::Response。其中Client::execute可在当前会话中执行任意 SQL 语句。
toysql二进制
toydb::Client是编程 API,toysql(src/bin/toysql.rs)则提供典型的REPL(读取-求值-打印循环)交互界面:用户逐行输入 SQL 语句并查看结果。它同样是一个小巧的 clap 命令,参数是待连接的 toyDB 服务端地址;启动后先用toydb::Client连接服务端,再用 Rustyline 库启动交互 shell。shell 主体就是一个循环:提示用户输入 SQL 语句 → 通过toydb::Client::execute在服务端执行 → 把响应格式化并打印输出。
至此,一个功能完整的 SQL 数据库系统就组装完毕:客户端输入 SQL,经会话、词法分析、解析、规划、优化、执行流水线落到 MVCC 键值存储,写入再经 Raft 复制到整个集群,读则由领导者确认法定人数后返回强一致结果。
总结:一条自底向上的学习路径
toyDB 架构的巧妙之处在于每一层都只解决一个问题、并作为上一层的"地基":
- BitCask / Memory 存储引擎(src/storage)提供有序键值与区间扫描;
- Bincode / Keycode 编码(src/encoding)解决数据类型与保序二进制之间的转换;
- MVCC(src/storage/mvcc.rs)在键值之上叠加快照隔离事务;
- Raft(src/raft)把状态机复制到整个集群并提供线性一致读;
- SQL 引擎(src/sql)在复制状态机上实现完整 SQL 流水线;
- Server 与 Client(src/server.rs、src/client.rs)把这一切通过 Bincode-over-TCP 连成可运行的分布式数据库。
对于想理解分布式数据库原理的读者,推荐按上述顺序阅读 docs/architecture 架构文档的各个子页面,并在每一节后对照相应源码与src/*/testscripts下的测试脚本运行验证;对于想快速上手体验的读者,直接运行 cluster/run.sh 启动 5 节点集群,再用toysql连接即可。
【免费下载链接】toydb
Distributed SQL database in Rust, written as an educational project
相关推荐
toyDB 深入解析:用 Rust 从零构建分布式 SQL 数据库的教学实践
toyDB 深入解析:用 Rust 从零构建分布式 SQL 数据库的教学实践 toyDB 是一个用 Rust 从零手写、以教学为初衷的分布式 SQL 数据库,它
如何在5分钟内上手Dockerized?新手必备快速入门指南
如何在5分钟内上手Dockerized?新手必备快速入门指南 Dockerized是一款让你无需安装即可运行流行命令行工具的终极解决方案。它通过Docker容器
RL4CO与TorchRL集成:构建高效强化学习组合优化系统
RL4CO与TorchRL集成:构建高效强化学习组合优化系统 强化学习(RL)在组合优化(CO)领域的应用正迅速改变传统解决方案的效率边界。 RL4CO 作为基
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考