2.9 简易 Redis 服务器实战(上):架构设计与核心数据结构实现
引言:从理论到综合实战
在过去的几章中,我们已经学习了 Rust 并发编程的几大核心支柱:
- 多线程 (std::thread)
- 共享内存与锁 (Arc, Mutex, RwLock, dashmap)
- 消息传递 (mpsc, crossbeam-channel)
- 异步编程 (Future, async/await, tokio)
现在,是时候将这些知识融会贯通,应用到一个真实、有趣且具有挑战性的项目中了。我们将从零开始,分两部分构建一个功能简化的、支持并发的、异步的 Redis 服务器。
Redis 是一个非常流行的高性能内存键值数据库。它的协议简单,核心功能清晰,是学习网络编程和并发设计的绝佳案例。通过这个项目,你将亲手实践如何设计一个高性能网络服务,如何管理服务状态,以及如何处理客户端命令。
在本章(上篇),我们将专注于项目的设计阶段:
在下一章(下篇),我们将实现服务器的主循环、并发处理客户端连接以及命令的执行逻辑。
1. 功能范围定义
完整的 Redis 拥有数百个命令,我们不可能全部实现。作为一个学习项目,我们将目标锁定在最核心、最有代表性的几类命令上。
我们要实现的命令
- 连接管理:
- PING: 测试服务器是否在线。
- ECHO message: 回显消息。
- 字符串 (String) 操作:
- SET key value: 设置键值。
- GET key: 获取键值。
- INCR key: 将键的值加一。
- DECR key: 将键的值减一。
- 哈希 (Hash) 操作:
- HSET key field value: 设置哈希中的字段值。
- HGET key field: 获取哈希中的字段值。
- HGETALL key: 获取哈希中的所有字段和值。
- 列表 (List) 操作:
- LPUSH key value [value …]: 将一个或多个值插入到列表头部。
- RPUSH key value [value …]: 将一个或多个值插入到列表尾部。
- LPOP key: 移除并获取列表的第一个元素。
- RPOP key: 移除并获取列表的最后一个元素。
这些命令涵盖了 Redis 最常用的数据类型,足以让我们构建一个有实际意义的服务器。
2. 架构设计
在构建网络服务时,一个核心的决策是选择 I/O 模型和并发模型。
I/O 与并发模型选择
每连接一线程(Thread-per-connection):
- 模型:主线程 accept 连接,然后为每个连接 spawn 一个新的操作系统线程。
- 优点:实现简单直观,代码是同步阻塞式的,易于编写和推理。
- 缺点:无法扩展到大量连接(C10K 问题)。线程是昂贵资源,创建上万个线程会耗尽系统内存和调度能力。
异步 I/O 与事件循环(Asynchronous I/O with Event Loop):
- 模型:使用 tokio 这样的异步运行时。所有 I/O 操作都是非阻塞的。少量的工作线程通过事件循环来管理成千上万的并发连接(任务)。
- 优点:极高的并发能力,资源占用少,性能优异。
- 缺点:代码涉及 async/await,心智负担比同步代码稍高。
考虑到我们的目标是构建一个高性能的服务器,并且我们已经学习了 tokio,异步 I/O 模型是必然的选择。
整体架构
我们的服务器将由以下几个关键部分组成:
主监听循环 (Listener Loop):
- 在主 async fn main 中。
- 使用 tokio::net::TcpListener 绑定端口并异步地 accept 新的 TCP 连接。
- 每当有新连接进来,就 tokio::spawn 一个新的异步任务来处理这个客户端。
连接处理任务 (Connection Handler Task):
- 每个客户端连接都由一个独立的异步任务处理。
- 在一个循环中,从 TCP 流中读取数据,解析成 Redis 命令。
- 执行命令,并将结果写回 TCP 流。
- 处理连接关闭和错误。
共享状态/数据存储 (Shared State / Data Store):
- 所有客户端任务都需要访问和修改同一个数据集(我们的键值存储)。
- 这个数据存储必须是线程安全的。
- 它将由一个独立的模块管理,并通过 Arc 在所有任务间共享。
graph TD
subgraph Tokio 运行时
direction LR
subgraph Worker线程池
W1[Worker 1]
W2[Worker 2]
W…
Wn[Worker N]
end
subgraph 主任务
A[main()] –> B(TcpListener.bind);
B –> C{loop: listener.accept().await};
end
subgraph "共享数据存储 (Arc<Db>)"
DB[(HashMap/DashMap)]
end
C — 新连接 –> D(tokio::spawn);
D — 创建 –> T1[连接1处理任务];
D — 创建 –> T2[连接2处理任务];
D — 创建 –> T…
D — 创建 –> Tn[连接N处理任务];
T1 — 访问 –> DB;
T2 — 访问 –> DB;
Tn — 访问 –> DB;
T1 — 调度到 –> W1;
T2 — 调度到 –> W2;
Tn — 调度到 –> W1;
end
3. 核心数据结构设计
这是我们项目中最重要的部分之一。如何设计一个既能支持多种 Redis 数据类型,又是线程安全的数据存储?
数据类型枚举
首先,Redis 的一个 key 可以对应不同类型的值(String, List, Hash)。我们可以用一个 enum 来表示。
// src/data.rs
use std::collections::{HashMap, VecDeque};
// Redis 值可以是不同的类型
#[derive(Debug, Clone)]
pub enum RedisValue {
String(Vec<u8>), // 使用 Vec<u8> 而不是 String 来存储原始字节
List(VecDeque<Vec<u8>>),
Hash(HashMap<Vec<u8>, Vec<u8>>),
}
- 我们使用 Vec<u8> 而不是 String,因为 Redis 的键和值都是二进制安全的。
- 列表使用 VecDeque(双端队列),因为它在头部和尾部插入/删除都是 O(1) 的,完美匹配 LPUSH/RPUSH/LPOP/RPOP 的需求。
线程安全的数据存储
现在,我们需要一个地方来存储从 key (通常是 String 或 Vec<u8>) 到 RedisValue 的映射。这个地方必须是线程安全的。我们有几个选择:
我们选择 DashMap。
// src/db.rs
use crate::data::RedisValue;
use dashmap::DashMap;
use std::sync::Arc;
// 定义数据库的核心结构
#[derive(Debug, Clone)]
pub struct Db {
// 使用 Arc 包裹 DashMap,以便在多个地方共享它
// DashMap 的 Key 是 String,Value 是我们定义的 RedisValue
entries: Arc<DashMap<String, RedisValue>>,
}
impl Db {
pub fn new() -> Self {
Db {
entries: Arc::new(DashMap::new()),
}
}
// 这里可以添加与数据库交互的辅助方法
// 例如: get_string, set_string 等
// 这样可以将 DashMap 的具体实现细节封装起来
}
impl Default for Db {
fn default() -> Self {
Self::new()
}
}
我们将 DashMap 包装在一个 Db 结构体中。这是一种良好的设计实践,称为“Newtype Pattern”。它允许我们:
- 封装实现细节:将来如果我们想把 DashMap 换成其他实现,只需要修改 Db 内部,而不需要改动服务器的其他代码。
- 添加特定于领域的方法:我们可以在 Db 上实现 get_string, lpush 等方法,提供更高级、更安全的接口。
4. 命令解析:RESP 协议
Redis 使用一种名为 REdis Serialization Protocol (RESP) 的协议。这是一种人类可读、实现简单的文本协议。RESP 可以序列化不同的数据类型,如简单字符串、错误、整数、批量字符串和数组。
RESP 数据类型
| 简单字符串 | + | +OK\\r\\n | 用于传输非二进制安全的简单字符串,如状态回复。 |
| 错误 | – | -ERR unknown command\\r\\n | 用于传输错误信息。 |
| 整数 | : | :1000\\r\\n | 64 位有符号整数。 |
| 批量字符串 | $ | $6\\r\\nfoobar\\r\\n | 用于传输二进制安全的字符串,最长 512MB。 |
| $0\\r\\n\\r\\n | 空字符串。 | ||
| $-1\\r\\n | Null 值(不存在)。 | ||
| 数组 | * | *2\\r\\n$3\\r\\nfoo\\r\\n$3\\r\\nbar\\r\\n | 用于传输其他 RESP 类型的集合。 |
| *-1\\r\\n | Null 数组。 |
客户端发送的命令总是以 RESP 数组 的形式。例如,命令 SET mykey myvalue 会被编码为:
*3\\r\\n$3\\r\\nSET\\r\\n$5\\r\\nmykey\\r\\n$7\\r\\nmyvalue\\r\\n
这表示一个包含 3 个元素的数组,每个元素都是一个批量字符串。
解析器设计
我们需要编写一个解析器,它能从 TcpStream 读取字节流,并将其转换为我们能理解的命令结构。
首先,定义命令的 enum。
// src/command.rs
// 代表解析后的命令
#[derive(Debug)]
pub enum Command {
Ping,
Echo(Vec<u8>),
Set(String, Vec<u8>),
Get(String),
Incr(String),
Decr(String),
HSet(String, Vec<u8>, Vec<u8>),
HGet(String, Vec<u8>),
HGetAll(String),
LPush(String, Vec<Vec<u8>>),
RPush(String, Vec<Vec<u8>>),
LPop(String),
RPop(String),
Unknown(String),
}
接下来,是解析器的核心逻辑。从 TCP 流中解析 RESP 是一个有状态的过程,因为一个 TCP 包可能不包含完整的命令。我们需要一个缓冲和解析循环。
为了简化,我们可以先假设我们一次性收到了完整的命令数据。
// src/resp.rs
use crate::command::Command;
// 一个简单的 RESP 帧的枚举
#[derive(Debug, Clone)]
pub enum Frame {
Simple(String),
Error(String),
Integer(i64),
Bulk(Vec<u8>),
Array(Vec<Frame>),
Null,
}
// 这是一个简化的解析函数,假设 buf 包含一个完整的 RESP 帧
// 生产级的解析器要复杂得多,需要处理不完整的流
pub fn parse_frame(buf: &[u8]) -> Result<(Frame, usize), &'static str> {
match buf.get(0) {
Some(b'+') => parse_simple_string(buf),
Some(b'*') => parse_array(buf),
Some(b'$') => parse_bulk_string(buf),
_ => Err("Unsupported frame type"),
}
}
// … (各种 parse_* 函数的实现) …
// 解析命令数组
pub fn parse_command(frame: Frame) -> Result<Command, &'static str> {
let array = match frame {
Frame::Array(arr) => arr,
_ => return Err("Client command must be an array"),
};
if array.is_empty() {
return Err("Empty command");
}
// 将 Frame::Bulk 转换为 String 或 Vec<u8>
let command_name_frame = array.get(0).unwrap();
let command_name = match command_name_frame {
Frame::Bulk(data) => String::from_utf8(data.clone()).map_err(|_| "Command name is not valid UTF-8")?,
_ => return Err("Command name must be a bulk string"),
};
let command_name = command_name.to_uppercase();
match command_name.as_str() {
"PING" => Ok(Command::Ping),
"ECHO" => {
// … 实现 ECHO 的参数解析 …
Ok(Command::Unknown("ECHO not implemented".into()))
},
"SET" => {
// … 实现 SET 的参数解析 …
Ok(Command::Unknown("SET not implemented".into()))
},
// … 其他命令 …
_ => Ok(Command::Unknown(command_name)),
}
}
注意:一个完整的、健壮的 RESP 解析器需要处理很多边界情况,例如不完整的帧、缓冲区管理等。在 tokio 生态中,通常使用 tokio_util::codec 来实现这种有状态的流解析。为了聚焦于项目整体,我们这里的解析器会进行简化。
项目结构整合
现在,我们的项目目录看起来像这样:
src/
├── main.rs # 服务器入口,监听循环
├── db.rs # 数据库核心,使用 DashMap
├── data.rs # 定义 RedisValue 枚举
├── command.rs # 定义 Command 枚举
├── resp.rs # RESP 协议解析逻辑 (待详细实现)
└── connection.rs # (将在下篇实现) 处理单个客户端连接的逻辑
main.rs 的骨架
// src/main.rs
mod command;
mod data;
mod db;
mod resp;
// mod connection; // 稍后添加
use db::Db;
use tokio::net::TcpListener;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let listener = TcpListener::bind("127.0.0.1:6379").await?;
println!("Mini-Redis is listening on port 6379");
// 创建一个全局的、线程安全的数据库实例
let db = Db::new();
loop {
let (socket, _) = listener.accept().await?;
// 克隆 Db 的句柄 (内部是 Arc,所以是轻量操作)
let db_clone = db.clone();
println!("Accepted new connection");
tokio::spawn(async move {
// 在这里调用连接处理逻辑
// process_connection(socket, db_clone).await;
});
}
}
上篇总结
在本章,我们为简易 Redis 服务器项目打下了坚实的地基。
我们已经做好了所有的前期准备。现在,我们有了一个清晰的蓝图,知道需要哪些组件以及它们如何协同工作。
在下一章中,我们将卷起袖子,完成最核心的编码工作:实现完整的 RESP 解析、处理客户端连接的循环、执行具体的命令逻辑,最终让我们的服务器真正“活”起来。
网硕互联帮助中心





评论前必须登录!
注册