在构建高性能服务时,键值存储是最基础也最重要的组件之一。Redis 凭借其出色的性能和简洁的设计,成为了事实上的标准。本文将从零开始,用 Rust 构建一个兼容 Redis 协议的最小可用服务器,深入探讨 RESP 协议解析、跳表(Skip List)索引实现、键过期淘汰策略以及事件驱动架构的设计思路。

一、目标与设计取舍

我们的目标是实现一个最小可用的 Redis-like 服务,支持以下核心功能:

  • 完整的 RESP3 协议解析与序列化
  • 基础数据结构:String、List、Set、Hash、Sorted Set
  • 键过期机制(TTL)与惰性删除 + 定期删除双重策略
  • 单线程事件循环 + epoll/kqueue 异步 IO
  • 兼容 redis-cli 和主流客户端连接

设计取舍上,我们选择:

  1. 单线程事件循环:避免锁竞争,与 Redis 原版设计一致
  2. 跳表作为 Sorted Set 的唯一索引:实现简单,区间查询高效
  3. Tokio 作为异步运行时:避免从零写 epoll 循环,聚焦业务逻辑
  4. Value enum 统一值类型:而非泛型特化,简化内存管理
  5. 二、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%。差距主要来自:

    1. Rust 的枚举值类型检查:Value 的 enum dispatch 虽被优化为跳转表,但仍有一定开销
    2. 序列化/反序列化:RESP 帧的解析/生成是主要瓶颈
    3. 异步运行时开销:Tokio 的 task 调度和 epoll 交互带来边际成本
    4. 八、持久化设计(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 服务只是第一步,生产部署还需要:

      1. 主从复制:通过 PSYNC 命令实现增量同步,解决网络断连后的全量同步问题
      2. Sentinel 高可用:自动故障转移的监控和决策系统
      3. Cluster 分片:16384 个哈希槽的一致性哈希分片
      4. 慢查询日志:记录执行时间超过阈值(slowlog-log-slower-than)的命令
      5. 内存碎片整理:jemalloc 的 arena 碎片率过高时的 defrag 策略
      6. 十、总结

        从零构建 Redis-like 服务器是一个极佳的学习项目,它涵盖了网络编程、数据结构、并发模型、持久化策略等多方面技能。通过本文的实现,我们学到了:

        • RESP 协议的简洁性:基于前缀类型的自描述协议,解析简单且扩展性强
        • 跳表的工程优势:比平衡树更简洁,区间查询天然高效
        • 混合过期策略的智慧:惰性删除 + 定期删除的平衡取舍
        • 异步 IO + 共享状态架构:Tokio 的 async/await 天然适合 IO-bound 的命令处理

        完整的源码实现约 2000 行 Rust 代码,可以在 [GitHub 仓库](https://example.com) 中查看。如果对某个部分(如集群分片或 Raft 共识)有进一步的兴趣,可以单独深入展开一篇文章来讨论。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿
网站二维码

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部
/* 跳过导航链接 (无障碍) */ position: absolute; top: -100px; left: 15px; z-index: 99999; padding: 8px 16px; background: #007bff; color: #fff; font-size: 14px; border-radius: 0 0 4px 4px; text-decoration: none; transition: top 0.2s; } top: 0; outline: 3px solid #0056b3; }