【免费下载链接】toydb
Distributed SQL database in Rust, written as an educational project
toyDB 是一套用 Rust 编写的分布式 SQL 数据库(教育项目),其核心运行实体是toydb::Server——一个把 Raft 共识节点、Raft 对等节点和 SQL 客户端三方网络流量统一路由起来的枢纽。本文以 docs/architecture/server.md 为主线,结合 src/server.rs 源码与 config/toydb.yaml、cluster/run.sh 等配置,系统讲解 toyDB 服务器如何选择网络协议、如何组织多线程消息路由、如何为 SQL 客户端提供会话服务,以及toydb二进制如何把配置、存储与 Raft 状态机组装成一台可运行的节点。读完本文,你将完整掌握 toyDB 节点内部"一条 SQL 从客户端到落库"的完整链路,以及 5 节点集群从零启动的配置与运行方法。
如上图所示,Server 位于节点内部各组件之上:底层是存储引擎(Storage Engine)与 Raft 共识引擎(Raft Consensus Engine),上层是 SQL 引擎(SQL Engine),而 Server 负责管理对外网络通信,既面向 SQL 客户端,也面向其他 Raft 节点。
1. 服务器整体结构:一个包裹 Raft 节点的路由器
toydb::Server定义在 src/server.rs,其职责可以用源码注释里的三句话概括:
- 通过 TCP 监听 SQL 客户端的入站连接,把请求交给本地 Raft 节点处理;
- 通过 TCP 监听其他 toyDB 节点的入站 Raft 连接,把消息交给本地 Raft 节点;
- 主动连接其他 toyDB 节点,把本地 Raft 节点的出站消息发送出去。
Server 内部持有三个核心成员:
pub struct Server { /// 内部 Raft 节点。 node: raft::Node, /// 来自 Raft 节点的出站消息。 node_rx: Receiver<raft::Envelope>, /// Raft 对等节点 ID 与地址映射。 peers: HashMap<raft::NodeID, String>, }从源码结构看,Server 本身不实现任何共识逻辑,而是"包裹"一个raft::Node:Raft 节点负责状态机推进、选举与日志复制,Server 只负责把网络上的消息搬进搬出。Server::new()会创建一条无界 Crossbeam 通道node_tx/node_rx,把发送端交给raft::Node::new(),使 Raft 节点产生的所有出站消息都能流回 Server(src/server.rs)。
为什么不用 async Rust?
toyDB 的服务器刻意没有使用 async/await 或 Tokio,而是采用普通操作系统线程。作者在文档中给出的理由是:async Rust 会显著增加代码复杂度,反而掩盖了教学项目想展示的核心概念;而 async 带来的效率提升对 toyDB 而言完全无关紧要(toyDB 明确自我定位为"not performant、not efficient")。
在线程之间,消息统一通过 Crossbeam channels 中也能印证:crossbeam = { version = "0.8", features = ["crossbeam-channel"] }。整个服务器的主循环Server::serve()正是围绕多个 channel 的crossbeam::select!轮询实现的。
2. 网络协议:Bincode 编码 + 原生 TCP
服务器与外界(Raft 对等节点、SQL 客户端)的通信协议是Bincode 序列化 + TCP 裸连接。Bincode 是 toyDB 在编码章节讨论过的二进制序列化格式,在 src/encoding/mod.rs 中通过encoding::Valuetrait 提供encode_into()、maybe_decode_from()等读写方法。
关键设计点:协议不需要任何额外的 framing(帧边界)机制。因为 Bincode 在反序列化时能够根据目标类型精确知道每种消息需要读取多少字节,所以直接连续地把消息写进 TCP 流即可,接收方每次maybe_decode_from()都能准确地截取一条完整消息。这让 wire protocol 极其简单——一条消息就是一个 Bincode 字节流。
所有 Raft 网络消息都被封装在raft::Envelope中(src/raft/message.rs),它带有from(发送者节点 ID)、to(接收者节点 ID)、term(发送者当前任期)和message(具体消息体,如Heartbeat、Append、ClientRequest等)。Envelope同样实现了encoding::Value。
3. 端口规划与监听线程
Server::serve()(src/server.rs)是服务器的入口主循环,它创建两个 TCP 监听器:
- Raft 端口:用于 Raft 对等节点之间的消息;
- SQL 端口:用于 SQL 客户端的连接。
原文档提到的端口为Raft 9705 / SQL 9605,这是 5 节点集群中末位节点的端口。结合 cluster/run.sh 与 cluster/toydb1/toydb.yaml 可以看到,完整 5 节点集群实际使用SQL 端口 9601-9605、Raft 端口 9701-9705(节点 1 的 SQL 端口 9601、Raft 端口 9701,依此类推);单节点默认配置 config/toydb.yaml 则监听localhost:9601(SQL)与localhost:9701(Raft)。这些地址均可通过listen_sql/listen_raft配置项自由调整。
serve()内部使用std::thread::scope派生四类线程:
| 线程 | 职责 | 源码位置 |
|---|---|---|
raft_accept | 接受入站 Raft 对等连接 | src/server.rs |
raft_send_peer | 为每个 Raft 对等节点维持出站连接并发送消息 | src/server.rs |
raft_route | 服务器核心路由循环,驱动 Raft 节点 | src/server.rs |
sql_accept | 接受入站 SQL 客户端连接 | src/server.rs |
4. Raft 路由:服务器的心脏
Server::raft_route()是整个服务器最重要的线程。它持有唯一的raft::Node所有权,在一个loop里用crossbeam::select!同时监听四条事件源:
- 定时器(ticker):以
raft::TICK_INTERVAL(定义在 src/raft/mod.rs,为 100ms)为周期调用raft::Node::tick(),驱动 Raft 的选举超时、心跳等定时逻辑; - 对等节点消息(peers_rx):把远端 Raft 节点发来的
Envelope通过raft::Node::step()步进到本地节点; - 本地节点出站消息(node_rx):读取 Raft 节点产生的出站
Envelope,按to字段路由到对应 peer 的发送通道; - 本地 SQL 客户端请求(request_rx):接收来自
sql::engine::Raft引擎的raft::Request,包装成ClientRequest消息步进到 Raft 节点。
请求-响应的匹配机制
raft::Node与客户端之间是异步消息模型,如何把响应匹配回发起请求的客户端?raft_route维护了一张response_txs: HashMap<RequestID, Sender<Result<raft::Response>>>:
- 收到
request_rx上的客户端请求时,生成一个全局唯一的 UUIDv4 作为请求 ID(Uuid::new_v4(),见 src/raft/message.rs),把请求连同响应通道一起步进到 Raft 节点,并登记到映射表; - 当从
node_rx读到一条收件人是本节点、且消息体是ClientResponse的消息时,按 ID 取出对应的响应通道,把结果发回等待的客户端。
其余出站消息(心跳、追加日志等)则按目标节点 ID 从peers_tx找到对应的有界通道转发出去。若通道已满(peer 慢或不可达),消息会被丢弃并记录 error 日志——Raft 协议本身能容忍消息丢失,靠重试与心跳机制自愈。
值得一提的细节:raft_route中所有node.step()/node.tick()出错都会直接 panic,因为 Raft 节点无法从不成功的状态转换中恢复。
5. Raft 对等连接:发送与接收线程
服务器启动时,会为peers映射中的每一个Raft 对等节点派生一个raft_send_peer()线程。该线程与 peer 之间的通道是有界的,容量为RAFT_PEER_CHANNEL_CAPACITY = 1000(src/server.rs),超出后丢弃消息。
raft_send_peer的行为是"无限重连":
- 尝试
TcpStream::connect(addr)连接 peer,失败则记录错误并休眠RAFT_PEER_RETRY_INTERVAL = 1秒后重试; - 连接成功后,阻塞地从通道接收
Envelope,用encode_into()写入 TCP 并用flush()刷出; - 一旦写入失败(连接断开),跳出内层循环,回到外层继续重连。
与之相对,raft_accept()线程在 Raft 监听端口上循环accept(),每接受一条连接就派生一个raft_receive_peer()线程,从BufReader中不断raft::Envelope::maybe_decode_from()读取 Bincode 消息,并通过raft_step_tx通道转发给raft_route()步进到本地 Raft 节点。
至此,集群中所有节点的发送方向与接收方向各自建立 TCP 连接(一条消息与其响应分别走两条独立的出站 TCP 连接,见 src/raft/message.rs 的注释),Raft 集群完全互联,节点之间可以互相通信。
6. SQL 服务:客户端的会话处理
SQL 侧采用toydb::Request/toydb::Response枚举作为客户端协议(同样 Bincode 编码于 TCP 之上),定义在 src/server.rs:
pub enum Request { Execute(String), // 执行一条 SQL 语句 GetTable(String), // 获取指定表的 schema ListTables, // 列出所有表 Status, // 返回服务器状态 } pub enum Response { Execute(StatementResult), Row(Option<Row>), GetTable(Table), ListTables(Vec<String>), Status(Status), }6.1 引擎接入
服务器在serve()中创建sql::engine::Raft引擎(src/server.rs),其构造参数正是通往raft_route的raft_request_tx发送端。sql::engine::Raft的实现位于 src/sql/engine/raft.rs:它本身只是一个 Raft 客户端——每次请求都通过tx.send((request, response_tx))把raft::Request交给本地 Raft 节点,再阻塞等待响应通道。
读请求(Read)与写请求(Write)被区分对待:
- 写命令(如
Write::Insert、Write::Commit、Write::CreateTable)走raft::Request::Write,复制到所有节点并顺序应用到各节点本地的 SQL 状态机(State<E>),保证确定性; - 读命令(如
Read::Get、Read::Scan、Read::BeginReadOnly)走raft::Request::Read,只由 Leader 在本机执行,不参与日志复制,避免不必要的复制往返。
读请求不进入复制日志这一点,在 src/sql/engine/raft.rs 中体现得尤其清晰:只读事务BeginReadOnly直接以Read提交,而写事务Begin才走Write——因为只读事务不需要分配新的 MVCC 版本,也就无需任何写入。
6.2 会话线程
sql_accept()循环接受 SQL 客户端连接,每接受一条连接就为该客户端创建一个独立的sql::execution::Session<sql::engine::Raft>(通过sql_engine.session()获得),并派生sql_session()线程服务它。
sql_session()(src/server.rs)的处理循环非常直接:
- 用
BufReader从 TCP 连接读取Request::maybe_decode_from(); - 按请求类型分发执行:
Request::Execute(query)→session.execute(&query),返回StatementResult;Request::GetTable/Request::ListTables→ 在只读事务(with_txn(true, ...))中读取表结构;Request::Status→ 组装包含server、raft、mvcc三段状态的Status结构体;
- 把
Response用encode_into()写回BufWriter并 flush。
每个客户端连接对应一个独立 Session,因此每个客户端拥有独立的事务上下文;所有 SQL 执行最终都汇聚到本地 Raft 节点,从而保证整个集群的一致性视图。
7.toydb二进制:从配置文件到运行中的节点
toydb可执行文件是toydb::Server的薄包装,位于src/bin/toydb.rs。它是一个基于clap派生宏的小型命令行程序,启动流程分三步:
- 解析配置:从
toydb.yaml文件读取服务器配置(节点 ID、peers、监听地址、存储引擎、fsync 开关等); - 初始化存储与状态机:创建 Raft 日志存储(
raft::Log)、Raft 持久状态(raft::State),以及由sql::engine::Raft::new_state()构建的 SQL 状态机; - 启动服务器:调用
toydb::Server::new(...)后执行serve(listen_raft, listen_sql),服务器随即开始监听并提供服务。
7.1 配置项全解
config/toydb.yaml 是完整的单节点配置模板,各参数含义如下:
| 配置项 | 默认值/示例 | 说明 |
|---|---|---|
id | 1 | 节点 ID,集群内必须唯一 |
peers | {} | peer ID 与 Raft 地址的映射,单节点为空映射 |
listen_sql | localhost:9601 | SQL 客户端监听地址 |
listen_raft | localhost:9701 | Raft 对等节点监听地址 |
log_level | INFO | 日志级别,可选DEBUG、INFO、WARN、ERROR |
data_dir | data | 节点数据目录;Raft 日志存于raft文件,SQL 数据库存于sql文件 |
storage_raft | bitcask | Raft 日志存储引擎:bitcask(默认,追加式日志结构存储)或memory(基于标准库 BTreeMap 的内存存储) |
storage_sql | bitcask | SQL 数据库存储引擎,取值同上 |
fsync | true | 是否对写入执行 fsync。关闭可大幅提升写性能,但主机崩溃时可能丢数据并破坏 Raft 保证;仅影响 Raft 日志写入(SQL 状态机从不 fsync,因为它可从 Raft 日志重建) |
compact_threshold | 0.2 | 触发 Bitcask 日志压缩的最小垃圾比例 |
compact_min_bytes | 1000000 | 触发压缩的最小字节数 |
其中fsync的设计非常体现教学项目的用心:SQL 状态机可以被 Raft 日志完整重建,因此无需 fsync;而 Raft 日志的持久性直接关系数据安全,所以受fsync开关控制。
7.2 5 节点集群配置与运行
cluster/run.sh 提供了一条命令启动 5 节点集群的脚本,它会:
cargo build --release --bin toydb以 release 模式编译;- 为节点 1-5 分别执行
cargo run -q --release -- -c toydb$ID/toydb.yaml,后台启动并给输出加节点前缀; - 注册 trap,在脚本退出(如 Ctrl-C)时终止所有 toyDB 进程。
节点 1 的配置 cluster/toydb1/toydb.yaml 展示了 peer 映射的写法:
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:9705注意 peer 键带引号(YAML 中需将 ID 视为字符串),值是对应节点的 Raft 监听地址。集群启动后,即可通过cargo run --release --bin toysql连接节点 1 的 9601 端口执行 SQL。
8. 一条 SQL 请求的完整旅程
综合以上各节,可以梳理出一条 SQL 语句在 toyDB 节点内部的完整调用链:
- 客户端通过 TCP 连接节点(如 9601 端口),发送 Bincode 编码的
Request::Execute("INSERT ..."); sql_accept接受连接,为其创建Session<sql::engine::Raft>,sql_session线程读取请求并调用session.execute();- SQL 执行器把语句解析、规划后,通过
sql::engine::Raft把操作编码为raft::Request::Write(或Read),连同一次性响应通道发送到raft_request_tx; raft_route收到请求,生成 UUIDv4 请求 ID,包装成Envelope{ message: ClientRequest }步进到本地raft::Node,并登记响应通道;- Raft 节点作为 Leader 将写命令复制到各 follower(经
raft_send_peer线程、Bincode over TCP),达到多数派后提交,并应用到每个节点本地的 SQL 状态机State<E>(执行Local引擎的对应写操作); - 应用结果以
ClientResponse消息从 Raft 节点流出,raft_route识别出目标为本节点的响应消息,按 ID 取出响应通道发回sql::engine::Raft; sql_session把结果组装成Response::Execute,Bincode 编码写回客户端连接。
读请求的路径更短:Read不进入复制日志,Leader 在本机用State::read直接执行(为保证线性一致性,会先通过Read{ seq }消息向多数派确认领导权)。
9. 小结
toydb::Server是连接客户端、Raft 节点与对等节点的网络中枢,其设计哲学贯穿 toyDB 的整体定位:用最直白的机制讲清楚分布式数据库的核心概念。Bincode + TCP 免去了复杂的协议框架,普通线程 + Crossbeam channel 取代了 async 运行时,raft_route单循环集中呈现了消息路由的全部形态。从 docs/architecture/server.md 出发,配合 src/server.rs、src/sql/engine/raft.rs、src/raft/message.rs 以及 config/toydb.yaml、cluster/run.sh,你可以继续深入客户端协议(见 docs/architecture/client.md)、Raft 节点实现与 SQL 执行引擎,逐步拼出 toyDB 的完整图景。
【免费下载链接】toydb
Distributed SQL database in Rust, written as an educational project
相关推荐
LunaTranslator 网络服务与 API 接口完全指南:页面路由、HTTP/WebSocket 服务与源码级实现解析
LunaTranslator 网络服务与 API 接口完全指南:页面路由、HTTP/WebSocket 服务与源码级实现解析 LunaTranslator(视觉
桌面应用OCR人工智能paraphrase-albert-small-v2高级技巧:句子嵌入归一化与PyTorch性能优化方法
paraphrase albert small v2高级技巧:句子嵌入归一化与PyTorch性能优化方法 paraphrase albert small v2是
LunaTranslator 内置网络服务完全指南:HTTP API 与 WebSocket 的页面路由、接口协议与源码实现
LunaTranslator 内置网络服务完全指南:HTTP API 与 WebSocket 的页面路由、接口协议与源码实现 LunaTranslator 除了
桌面应用OCR人工智能
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考