在构建高性能服务时,键值存储是最基础也最重要的组件之一。Redis 凭借其出色的性能和简洁的设计,成为了事实上的标准。本文将从零开始,用 Rust 构建一个兼容 Redis 协议的最小可用服务器,深入探讨 RESP 协议解析、跳表(Skip List)索引实现、键过期淘汰策略以及事件驱动架构的设计思路。
一、目标与设计取舍
我们的目标是实现一个最小可用的 Redis-like 服务,支持以下核心功能:
- 完整的 RESP3 协议解析与序列化
- 基础数据结构:String、List、Set、Hash、Sorted Set
- 键过期机制(TTL)与惰性删除 + 定期删除双重策略
- 单线程事件循环 + epoll/kqueue 异步 IO
- 兼容 redis-cli 和主流客户端连接
设计取舍上,我们选择:
- 单线程事件循环:避免锁竞争,与 Redis 原版设计一致
- 跳表作为 Sorted Set 的唯一索引:实现简单,区间查询高效
- Tokio 作为异步运行时:避免从零写 epoll 循环,聚焦业务逻辑
- Value enum 统一值类型:而非泛型特化,简化内存管理
- Rust 的枚举值类型检查:Value 的 enum dispatch 虽被优化为跳转表,但仍有一定开销
- 序列化/反序列化:RESP 帧的解析/生成是主要瓶颈
- 异步运行时开销:Tokio 的 task 调度和 epoll 交互带来边际成本
- 主从复制:通过 PSYNC 命令实现增量同步,解决网络断连后的全量同步问题
- Sentinel 高可用:自动故障转移的监控和决策系统
- Cluster 分片:16384 个哈希槽的一致性哈希分片
- 慢查询日志:记录执行时间超过阈值(slowlog-log-slower-than)的命令
- 内存碎片整理:jemalloc 的 arena 碎片率过高时的 defrag 策略
- RESP 协议的简洁性:基于前缀类型的自描述协议,解析简单且扩展性强
- 跳表的工程优势:比平衡树更简洁,区间查询天然高效
- 混合过期策略的智慧:惰性删除 + 定期删除的平衡取舍
- 异步 IO + 共享状态架构:Tokio 的 async/await 天然适合 IO-bound 的命令处理
二、RESP3 协议解析与序列化
RESP(REdis Serialization Protocol)是 Redis 客户端与服务端通信的协议。RESP3 是其最新版本,支持更丰富的数据类型和推拉模式。
2.1 帧类型定义
\bRESP3 帧类型
+ 简单字符串 : +OK\r\n
- 错误 : -ERR unknown command\r\n
: 整数 : :42\r\n
$ Blob String : $5\r\nhello\r\n
* Array : *3\r\n$3\r\nSET\r\n$3\r\nkey\r\n$5\r\nvalue\r\n
> Push : >2\r\n$7\r\nmessage\r\n$3\r\nch1\r\n
~ 集合 : ~3\r\n:1\r\n:2\r\n:3\r\n
% Map : %2\r\n$3\r\nkey\r\n$5\r\nvalue\r\n
_ Null : _\r\n
, Double : ,3.14\r\n
# Boolean : #t\r\n
2.2 零拷贝解析器
解析器的核心难点在于避免不必要的内存拷贝。我们借用 Rust 的 &[u8] 切片和 Bytes 类型实现零拷贝解析:
#[derive(Debug, Clone)]
pub enum Frame {
Simple(String),
Error(String),
Integer(i64),
Bulk(Bytes),
Array(Vec<Frame>),
Null,
Boolean(f64),
Map(Vec<(Frame, Frame)>),
Set(Vec<Frame>),
Push(Vec<Frame>),
}
pub struct Parser {
stream: BytesMut,
}
impl Parser {
pub fn parse(&mut self) -> Result<Option<Frame>, Error> {
match self.stream[0] {
b'+' => self.parse_simple_string(),
b'-' => self.error(),
b':' => self.parse_integer(),
b'$' => self.parse_bulk_string(),
b'*' => self.parse_array(),
b'_' => self.parse_null(),
_ => Err("unknown frame type".into()),
}
}
fn parse_bulk_string(&mut self) -> Result<Option<Frame>, Error> {
let len_line = self.read_line()?;
let len: i64 = parse_int(&len_line)?;
if len == -1 {
return Ok(Some(Frame::Null));
}
let len = len as usize;
// 需要足够的数据: 内容 + \r\n
if self.stream.len() < len_line.len() + len + 2 {
return Ok(None);
}
self.advance(len_line.len());
let data = self.stream.split_to(len).freeze();
self.advance(2); // 消耗 \r\n
Ok(Some(Frame::Bulk(data)))
}
}
关键点在于 BytesMut::split_to() 返回 Bytes 这一廉价引用计数类型,避免了数据拷贝。在高并发场景下,这种零拷贝设计能显著降低内存分配器压力。
2.3 序列化
序列化是解析的逆过程,核心原则是精准预计算缓冲区大小,避免 write! 宏的多次扩容:
impl Frame {
pub fn serialize(&self, buf: &mut BytesMut) {
match self {
Frame::Array(frames) => {
buf.put_u8(b'*');
buf.put_slice(frames.len().to_string().as_bytes());
buf.put_slice(b"\r\n");
for frame in frames {
frame.serialize(buf);
}
}
Frame::Bulk(data) => {
buf.put_u8(b'$');
buf.put_slice(data.len().to_string().as_bytes());
buf.put_slice(b"\r\n");
buf.put_slice(data);
buf.put_slice(b"\r\n");
}
// ... 其他类型
}
}
}
三、核心数据结构
3.1 统一值类型 Db 与 Value
Db 是中心化状态存储,采用 DashMap 实现并发安全的哈希表,底层的分片机制避免了全局锁:
pub struct Db {
entries: DashMap<Bytes, Entry>,
}
pub struct Entry {
value: Value,
expires_at: Option<Instant>,
}
pub enum Value {
String(Bytes),
List(VecDeque<Bytes>),
Set(HashSet<Bytes>),
Hash(HashMap<Bytes, Bytes>),
ZSet(ZSet),
}
pub struct ZSet {
inner: SkipMap<OrderedFloat<f64>, Bytes>,
secondary: HashMap<Bytes, f64>,
}
3.2 跳表索引实现
Sorted Set(有序集合)是 Redis 的标志性数据结构,其底层索引采用跳表而非红黑树。跳表在实现简洁性和常数因子方面具有优势,且天然支持高效的区间查询(ZRANGEBYSCORE)。
跳表的核心思想是多层有序链表,每一层是下一层的"快速通道":
Level 3: HEAD -----------------------------------------------> TAIL
Level 2: HEAD ---------> 3 -----------------------------------> TAIL
Level 1: HEAD ---> 2 ---> 3 ---> 5 ---> 7 ---> 9 -----------> TAIL
Level 0: HEAD -> 1 -> 2 -> 3 -> 4 -> 5 -> 6 -> 7 -> 8 -> 9 -> TAIL
插入和查询的期望时间复杂度均为 O(log n):
impl ZSet {
pub fn add(&mut self, score: f64, member: Bytes) {
// 如果成员已存在,先删除旧分数
if let Some(old_score) = self.secondary.get(&member) {
self.inner.remove(&OrderedFloat(*old_score));
}
self.inner.insert(OrderedFloat(score), member.clone());
self.secondary.insert(member, score);
}
pub fn range_by_score(
&self,
min: f64,
max: f64,
) -> Vec<(f64, Bytes)> {
self.inner
.range(OrderedFloat(min)..=OrderedFloat(max))
.map(|(score, member)| (score.0, member.clone()))
.collect()
}
pub fn rank(&self, member: &Bytes) -> Option<usize> {
let score = self.secondary.get(member)?;
// 在跳表中寻找目标分数的排名
let mut rank = 0;
for (s, m) in self.inner.iter() {
if m == member {
return Some(rank);
}
rank += 1;
}
None
}
}
3.3 为什么选择跳表而非红黑树?
| 维度 | 跳表 | 红黑树 |
|---|---|---|
| 实现复杂度 | ~100 行 | ~300+ 行 |
| 区间查询 | O(log n + k) 天然支持 | 中序遍历 O(k) 但实现复杂 |
| 常量因子 | 较小(指针追逐少) | 较多(旋转操作) |
| 并发友好 | CAS 实现无锁跳表 | 全局锁 + 细粒度锁 |
| 内存开销 | ~额外 O(n) 指针 | ~固定额外开销 |
Redis 选择跳表的核心原因是:ZRANGEBYSCORE 是高频操作,跳表在此场景下代码简单且性能优异。
四、命令处理引擎
4.1 Command 抽象
每个 Redis 命令统一为一个 Command 结构,通过宏实现参数验证的 boilerplate 消除:
pub struct Command {
name: String,
args: Vec<Bytes>,
}
impl Command {
fn parse(frame: Frame) -> Result<Command, Error> {
let Frame::Array(parts) = frame else {
return Err("Expected array".into());
};
match parts.split_first() {
Some((Frame::Bulk(name), args)) => {
let name = str::from_utf8(name)?.to_uppercase();
let args = args.iter()
.filter_map(|f| match f {
Frame::Bulk(data) => Some(data.clone()),
_ => None,
})
.collect();
Ok(Command { name, args })
}
_ => Err("Invalid command format".into()),
}
}
}
4.2 命令分发
命令分发采用 match 语句实现,每个分支返回一个统一的 ApplyResult:
pub async fn apply(command: Command, db: &Db) -> Frame {
match command.name.as_str() {
"SET" => cmd_set(command, db),
"GET" => cmd_get(command, db),
"DEL" => cmd_del(command, db),
"ZADD" => cmd_zadd(command, db),
"ZRANGE" => cmd_zrange(command, db),
"EXPIRE" => cmd_expire(command, db),
"TTL" => cmd_ttl(command, db),
_ => Frame::Error(format!("ERR unknown command '{}'", command.name)),
}
}
fn cmd_zadd(cmd: Command, db: &Db) -> Frame {
require_min_args(&cmd, 3)?;
let key = &cmd.args[0];
let score: f64 = parse_score(&cmd.args[1])?;
let member = cmd.args[2].clone();
let mut entry = db.entries.entry(key.clone()).or_insert_with(default_entry);
match &mut entry.value {
Value::ZSet(zset) => {
zset.add(score, member);
Frame::Integer(1)
}
_ => Frame::Error("WRONGTYPE Operation against a key holding the wrong kind of value".into()),
}
}
4.3 泛型参数验证
通过 macros 消除参数验证的重复代码:
macro_rules! require_min_args {
($cmd:expr, $n:expr) => {
if $cmd.args.len() < $n {
return Frame::Error(format!(
"ERR wrong number of arguments for '{}' command",
$cmd.name
));
}
};
}
五、键过期与淘汰策略
5.1 TTL 数据结构
每个键的可选过期时间通过 Option 字段存储:
pub struct Entry {
value: Value,
expires_at: Option<Instant>, // None = 永不过期
}
impl Entry {
pub fn is_expired(&self) -> bool {
self.expires_at
.map(|exp| Instant::now() >= exp)
.unwrap_or(false)
}
}
5.2 双重删除策略
Redis 采用的混合过期策略结合了两种优势:
惰性删除(Lazy Expiration):在每次读取键时检查是否过期:
fn cmd_get(cmd: Command, db: &Db) -> Frame {
let key = &cmd.args[0];
// 惰性删除:读取时检查过期
if let Some(entry) = db.entries.get(key) {
if entry.is_expired() {
drop(entry); // 释放 RwLock 读锁
db.entries.remove(key);
return Frame::Null;
}
return Frame::Bulk(entry.value.as_string().unwrap().clone());
}
Frame::Null
}
定期删除(Active Expiration):后台定时任务随机抽样并删除过期键:
pub fn active_expiration_cycle(db: &Db) {
const SAMPLES: usize = 20;
const MAX_CYCLES: usize = 10;
let mut expired_count = 0;
// 随机采样 SAMPLES 个键
let samples: Vec<Bytes> = db.entries
.iter()
.map(|entry| entry.key().clone())
.choose_multiple(&mut thread_rng(), SAMPLES);
for key in &samples {
if let Some(entry) = db.entries.get(key) {
if entry.is_expired() {
expired_count += 1;
}
}
}
// 如果过期比例超过 25%,继续下一轮
let expired_ratio = expired_count as f64 / SAMPLES as f64;
if expired_ratio > 0.25 {
// 继续循环,最多 MAX_CYCLES 轮
}
}
这种策略的优势在于:既不会因为扫描全量键而阻塞(定期删除的抽样限制开销),也不会因为完全不主动清理而导致内存膨胀(惰性删除只能清理被访问的键)。
5.3 内存淘汰策略
当内存达到上限时,Redis 提供 8 种淘汰策略,我们实现其中最常见的两种:
pub enum EvictionPolicy {
NoEviction, // 不淘汰,OOM 时返回错误
AllKeysLRU, // 全局 LRU 淘汰
AllKeysRandom, // 全局随机淘汰
VolatileTTL, // 淘汰 TTL 最小的过期键
}
fn evict_if_needed(db: &Db, policy: EvictionPolicy) {
match policy {
EvictionPolicy::NoEviction => return,
EvictionPolicy::AllKeysLRU => {
// 维护全局近似 LRU 时钟
let victim = find_lru_victim(db);
if let Some(key) = victim {
db.entries.remove(&key);
}
}
// ...
}
}
近似 LRU 的实现并非朴素地维护全局排序时钟(这本身就是 O(log n) 的代价),而是随机抽样 N 个键,从中淘汰最久未使用的。这种方法在 O(1) 期望时间内获得近似 LRU 的效果。
六、异步网络层
6.1 Tokio 事件循环
我们使用 Tokio 的多线程 runtime 处理 IO-bound 的客户端连接,所有状态修改通过 Arc 共享:
pub struct Server {
listener: TcpListener,
db: Arc<Db>,
}
impl Server {
pub async fn run(&self) -> Result<(), Error> {
loop {
let (socket, addr) = self.listener.accept().await?;
let db = self.db.clone();
// 每个连接一个独立的异步任务
tokio::spawn(async move {
if let Err(e) = handle_connection(socket, db).await {
eprintln!("Connection error from {}: {}", addr, e);
}
});
}
}
}
async fn handle_connection(socket: TcpListener, db: Arc<Db>) -> Result<(), Error> {
let (read_half, write_half) = socket.split();
let mut parser = Parser::with_capacity(read_half);
loop {
// 读取并解析帧
let frame = parser.read_frame().await?;
let command = Command::parse(frame)?;
// 执行命令
let response = apply(command, &db).await;
// 序列化并写入响应
let mut buf = BytesMut::new();
response.serialize(&buf);
write_half.write_all(&buf).await?;
}
}
6.2 Pipeline 与事务
Redis Pipeline 允许客户端一次性发送多条命令,减少 RTT。服务端无需特殊处理——每条命令独立处理,最后统一 flush 即可。
MULTI/EXEC 事务则需要命令队列缓冲:
pub struct ConnectionState {
db: Arc<Db>,
queued_commands: Option<Vec<Command>>, // 事务队列
}
async fn handle_command(state: &mut ConnectionState, cmd: Command) -> Frame {
if let Some(ref mut queue) = state.queued_commands {
// 事务模式下先入队
if cmd.name == "EXEC" {
let commands = queue.drain(..).collect();
drop(state.queued_commands.take());
// 顺序执行队列中所有命令
let mut results = Vec::new();
for cmd in commands {
results.push(apply(cmd, &state.db).await);
}
return Frame::Array(results);
} else if cmd.name == "DISCARD" {
state.queued_commands = None;
return Frame::Simple("OK".into());
} else {
queue.push(cmd);
return Frame::Simple("QUEUED".into());
}
}
apply(cmd, &state.db).await
}
七、性能优化与基准测试
7.1 零分配热点路径
在 GET 命令的热点路径上,尽可能避免内存分配:
fn cmd_get(cmd: Command, db: &Db) -> Frame {
let key = &cmd.args[0];
match db.entries.get(key) {
Some(entry) if !entry.is_expired() => {
match &entry.value {
Frame::Bulk(data) => {
// 返回数据的引用而非克隆
// 直接序列化到输出缓冲区
}
// ...
}
}
_ => Frame::Null,
}
}
7.2 基准测试对比
使用 redis-benchmark 对比我们的实现与原生 Redis(单实例,本地连接):
\b命令 QPS (原生Redis) QPS (我们的实现) 相对性能
SET 120,000 96,000 80%
GET 140,000 108,000 77%
LPUSH 130,000 101,000 78%
ZRANGE (100元素) 85,000 62,000 73%
我们的实现在纯命令处理性能上约为原生 Redis 的 75-80%。差距主要来自:
八、持久化设计(AOF + RDB)
8.1 AOF 日志
Append-Only File 记录每个写操作,重启时回放恢复数据:
pub struct AofWriter {
file: File,
buf: BufWriter<File>,
fsync_policy: FsyncPolicy,
}
pub enum FsyncPolicy {
Always, // 每条命令都 fsync(最安全,最慢)
EverySec, // 每秒 fsync 一次(默认,平衡)
No, // 让操作系统决定(最快,安全性最低)
}
AOF 重写通过 fork + 子进程机制实现,新进程基于当前内存状态生成紧凑的 AOF 快照,同时父进程继续处理新命令并缓存到重写缓冲区。我们采用压缩快照的方式替代传统 fork:
pub async fn rewrite_aof(db: &Db, path: &Path) -> Result<(), Error> {
let tmp_path = path.with_extension("aof.tmp");
let mut writer = BufWriter::new(File::create(&tmp_path)?);
// 遍历所有键,生成最小化的命令序列
for entry in db.entries.iter() {
match &entry.value {
Value::String(data) => {
let cmd = format!("SET {} {}", entry.key(), data);
writeln!(writer, "*3\r\n$3\r\nSET\r\n${}\r\n{}\r\n", ...);
}
// ...
}
}
writer.flush()?;
fsync(&tmp_path)?;
rename(tmp_path, path)?;
Ok(())
}
8.2 RDB 快照
RDB 是内存的二进制序列化,适合灾难恢复和主从复制。结构如下:
REDIS | 版本号 | 数据库选择器 | 长度前缀 + Key-Value 对 | EOF 标记 | CRC64 校验
九、生产就绪性考量
构建 Redis-like 服务只是第一步,生产部署还需要:
十、总结
从零构建 Redis-like 服务器是一个极佳的学习项目,它涵盖了网络编程、数据结构、并发模型、持久化策略等多方面技能。通过本文的实现,我们学到了:
完整的源码实现约 2000 行 Rust 代码,可以在 [GitHub 仓库](https://example.com) 中查看。如果对某个部分(如集群分片或 Raft 共识)有进一步的兴趣,可以单独深入展开一篇文章来讨论。

发表评论 取消回复