Rust 实现 LSM-Tree KV 存储引擎:从 MemTable 到 Compaction 全栈深度实战

在现代数据库和存储引擎的版图里,LSM-Tree(Log-Structured Merge-Tree)无疑是写入密集型场景的基石架构。从 Bigtable、LevelDB、RocksDB 到 TiKV、ScyllaDB、CockroachDB,几乎所有高性能键值存储系统都在 LSM-Tree 的骨架上构建。本文将从"为什么需要 LSM-Tree"出发,深入其核心数据结构与算法,然后用 Rust 从零实现一个具备 WAL、SkipList MemTable、多级 Compaction、Bloom Filter 加速的 LSM-Tree KV 引擎,最后给出工业级实现的性能调优要点。

1. 为什么选择 LSM-Tree:写入放大 vs 读取放大的博弈

传统 B+Tree 存储引擎(如 MySQL InnoDB)的原地更新特性意味着随机写入——当一行数据被修改时,磁盘磁头需要定位到对应页并就地覆写。这在磁盘时代尚可接受,但在追求极致写入吞吐的场景下(尤其是 NVMe SSD 并发 IO 能力极强的今天),随机写入成为瓶颈。

LSM-Tree 的核心思想极其优雅:将所有随机写入转化为顺序追加写入,将"修改"推迟到后台 Compaction 阶段批量处理。这带来了天然的优势:

  • 写入极快:所有写入先进内存(MemTable),只需顺序追加 WAL 日志
  • 空间利用率高:数据不可变合并,无 B+Tree 的页分裂问题
  • IO 友好:SSTable 天然支持顺序扫描和批量 Compaction

代价是读取可能需要查找多层级——这正是 LSM-Tree 设计中最需要精心优化的部分。

2. LSM-Tree 核心数据结构:三层架构

LSM-Tree 的数据流遵循严格的层级结构:

写入路径:
 客户端 PUT → WAL (顺序追加) → MemTable (内存有序结构) → Immutable MemTable → SSTable L0 → SSTable L1 → ... → SSTable Ln

读取路径:
  客户端 GET → MemTable → Immutable MemTable → L0 SSTables → L1 SSTables → ... → Ln SSTables

MemTable 是写入的第一站,通常采用 SkipList 或 B-Tree 实现,兼顾 O(log n) 的读写性能和高效的并发安全。SkipList 的优势在于实现简单、锁粒度细、支持无锁读取。

SSTable(Sorted String Table) 是 LSM-Tree 的基石存储单元,每个 SSTable 内含按 key 排序的数据块、索引块、Bloom Filter 和元数据块。SSTable 是不可变的——一旦写入磁盘就永远不会修改,这也是 LSM-Tree 能够实现顺序写入的关键保证。

3. WAL:持久性的最后防线

写入 MemTable 的数据尚未刷盘,如果进程崩溃,这些写入就会丢失。Write-Ahead Log 解决了这个问题:在修改内存数据结构之前,先将操作追加写入持久化的日志。

Rust 实现 WAL 的核心设计:

use std::fs::{File, OpenOptions};
use std::io::{BufWriter, Write, Read, Seek, SeekFrom};
use std::path::Path;

/// WAL 条目格式: [4字节 key_len] [4字节 value_len] [1字节 op_type] [key_bytes] [value_bytes] [4字节 CRC32]
pub struct Wal {
    writer: BufWriter<File>,
    path: PathBuf,
}

impl Wal {
    pub fn new(path: PathBuf) -> io::Result<Self> {
        let file = OpenOptions::new()
            .create(true)
            .append(true)
            .open(&path)?;
        Ok(Self {
            writer: BufWriter::with_capacity(64 * 1024, file), // 64KB buffer
            path,
        })
    }

    /// 追加一条 WAL 记录,使用 buffered write 聚合小写入
    pub fn append(&mut self, key: &[u8], value: Option<&[u8]>) -> io::Result<()> {
        let op: u8 = if value.is_some() { 1u8 } else { 0u8 }; // 1=PUT, 0=DELETE
        let key_len = key.len() as u32;
        let val_len = value.map_or(0, |v| v.len() as u32);

        // 序列化头部
        self.writer.write_all(&key_len.to_be_bytes())?;
        self.writer.write_all(&val_len.to_be_bytes())?;
        self.writer.write_all(&[op])?;
        self.writer.write_all(key)?;
        if let Some(v) = value {
            self.writer.write_all(v)?;
        }
        // CRC32 校验,用于崩溃恢复时检测部分写入
        let crc = crc32fast::hash(&self.writer.buffer()[..]);
        self.writer.write_all(&crc.to_be_bytes())?;
        Ok(())
    }

    /// 批量 fsync,在关键写入点调用
    pub fn flush_and_sync(&mut self) -> io::Result<()> {
        self.writer.flush()?;
        self.writer.get_ref().sync_all()?;
        Ok(())
    }

    /// 崩溃恢复:读取所有 WAL 条目
    pub fn recover(path: &Path) -> io::Result<Vec<WalEntry>> {
        let mut file = File::open(path)?;
        let mut entries = Vec::new();
        let mut buf = Vec::new();
        file.read_to_end(&mut buf)?;

        let mut offset = 0;
        while offset < buf.len() - 4 {
            let key_len = u32::from_be_bytes(buf[offset..offset+4].try_into().unwrap()) as usize;
            offset += 4;
            let val_len = u32::from_be_bytes(buf[offset..offset+4].try_into().unwrap()) as usize;
            offset += 4;
            let op = buf[offset];
            offset += 1;
            let key = buf[offset..offset+key_len].to_vec();
            offset += key_len;
            let value = if val_len > 0 {
                Some(buf[offset..offset+val_len].to_vec())
            } else { None };
            offset += val_len;
            offset += 4; // skip CRC
            entries.push(WalEntry { op, key, value });
        }
        Ok(entries)
    }
}

上述实现中,64KB 的 BufWriter buffer 将大量小写入聚合为大块顺序 IO,sync_all 的调用频率需要在"数据安全"和"写入吞吐"之间取得平衡。LevelDB 采用了 group commit 策略,将多个线程的 fsync 合并为一个,大幅提升了并发写入性能。

4. SkipList MemTable:高并发的内存有序结构

MemTable 需要支持高效的插入(O(log n))和范围扫描(迭代器),SkipList 是最佳选择之一——相比红黑树或 B-Tree,SkipList 的多层链表结构天然适合并发操作,且实现简洁。

use std::sync::atomic::{AtomicUsize, Ordering};
use std::ptr;
use rand::Rng;

const MAX_HEIGHT: usize = 12;
const BRANCHING_FACTOR: u32 = 4;

struct Node {
    key: Vec<u8>,
    value: Option<Vec<u8>>,
    next: Vec<AtomicUsize>, // 每层的原子指针
}

pub struct SkipList {
    head: *mut Node,
    height: AtomicUsize,
    size: AtomicUsize,
}

impl SkipList {
    pub fn new() -> Self {
        let head = Box::into_raw(Box::new(Node {
            key: Vec::new(),
            value: None,
            next: (0..MAX_HEIGHT).map(|_| AtomicUsize::new(0)).collect(),
        }));
        Self {
            head,
            height: AtomicUsize::new(1),
            size: AtomicUsize::new(0),
        }
    }

    /// 返回 key 的插入位置和每层的前驱节点
    pub fn insert(&self, key: Vec<u8>, value: Option<Vec<u8>>) {
        let mut prev = [ptr::null_mut(); MAX_HEIGHT];
        let mut curr = self.head;

        // 从最高层向下查找插入位置
        for level in (0..self.height.load(Ordering::Relaxed)).rev() {
            unsafe {
                loop {
                    let next_ptr = (*curr).next[level].load(Ordering::Acquire);
                    let next = if next_ptr == 0 { ptr::null_mut() } else { next_ptr as *mut Node };
                    if !next.is_null() && (*next).key < key {
                        curr = next;
                    } else {
                        break;
                    }
                }
                prev[level] = curr;
            }
        }

        // 随机决定新节点高度
        let node_height = self.random_height();
        let new_node = Box::into_raw(Box::new(Node {
            key,
            value,
            next: (0..node_height).map(|_| AtomicUsize::new(0)).collect(),
        }));

        // 逐层插入,使用 CAS 支持并发
        for level in 0..node_height {
            unsafe {
                loop {
                    (*new_node).next[level].store(
                        (*prev[level]).next[level].load(Ordering::Acquire),
                        Ordering::Relaxed,
                    );
                    let expected = (*prev[level]).next[level].load(Ordering::Relaxed);
                    if (*prev[level]).next[level]
                        .compare_exchange(expected, new_node as usize, Ordering::AcqRel, Ordering::Acquire)
                        .is_ok()
                    {
                        break;
                    }
                    // CAS 失败,重新定位前驱(简化处理:这里直接重试)
                }
            }
        }
        self.size.fetch_add(1, Ordering::Relaxed);
    }

    pub fn get(&self, key: &[u8]) -> Option<Vec<u8>> {
        let mut curr = self.head;
        unsafe {
            for level in (0..self.height.load(Ordering::Relaxed)).rev() {
                loop {
                    let next_ptr = (*curr).next[level].load(Ordering::Acquire);
                    let next = if next_ptr == 0 { ptr::null_mut() } else { next_ptr as *mut Node };
                    if next.is_null() { break; }
                    match (*next).key.as_slice().cmp(key) {
                        std::cmp::Ordering::Less => curr = next,
                        std::cmp::Ordering::Equal => return (*next).value.clone(),
                        std::cmp::Ordering::Greater => break,
                    }
                }
            }
        }
        None
    }

    fn random_height(&self) -> usize {
        let mut height = 1;
        let mut rng = rand::thread_rng();
        while rng.gen_ratio(1, BRANCHING_FACTOR) && height < MAX_HEIGHT {
            height += 1;
        }
        height
    }
}

上述 SkipList 的关键优化点在于:逐层 CAS 无锁插入 + 指数分布的高度选择。每层有 1/4 的概率继续增长,使得高层节点数约为低层的 1/4,整体空间开销约为 O(n × 4/3)。

5. SSTable 序列化格式:磁盘上的有序数据块

SSTable 的内存布局决定了读取效率。一个标准的 SSTable 文件结构如下:p>

[Data Block 1] [Data Block 2] ... [Index Block] [Bloom Filter Block] [Footer]

使用 Rust 实现 SSTable 的序列化与反序列化:

use std::collections::BTreeMap;

pub const BLOCK_SIZE: usize = 4 * 1024; // 4KB  blockSize

/// SSTable 构建器
pub struct SsTableBuilder {
    data_blocks: Vec<Vec<u8>>       index: Vec<(Vec<u8>, u32)>, // (last_key, block_offset)
    bloom: BloomFilter,
    current_block: Vec<u8>,
}

impl SsTableBuilder {
    pub fn new() -> Self {
        Self {
            data_blocks: Vec::new(),
            index: Vec::new(),
            bloom: BloomFilter::new(10), // 每个 key 约 10 bits
            current_block: Vec::with_capacity(BLOCK_SIZE),
        }
    }

    pub fn add_entry(&mut self, key: &[u8], value: Option<&[u8]>) {
        if self.current_block.len() >= BLOCK_SIZE && !self.index.is_empty() {
            self.flush_block();
        }
        self.bloom.insert(key);

        // 编码: [key_len(u32)] [key] [value_len(u32)] [value?]
        self.current_block.extend_from_slice(&(key.len() as u32).to_be_bytes());
        self.current_block.extend_from_slice(key);
        if let Some(v) = value {
            self.current_block.extend_from_slice(&(v.len() as u32).to_be_bytes());
            self.current_block.extend_from_slice(v);
        } else {
            self.current_block.extend_from_slice(&0u32.to_be_bytes());
        }
    }

    fn flush_block(&mut self) {
        let offset = self.data_blocks.iter().map(|b| b.len()).sum::<usize>() as u32;
        // 记录该 block 的最后一个 key 作为索引
        // (实际实现需从 current_block 末尾解析出最后一个 key)
        self.index.push((self.last_key_in_current_block.clone(), offset));
        self.data_blocks.push(std::mem::take(&mut self.current_block));
        self.current_block = Vec::with_capacity(BLOCK_SIZE);
    }

    pub fn build(mut self) -> Vec<u8> {
        if !self.current_block.is_empty() {
            self.flush_block();
        }

        let mut file = Vec::new();
        let data_offset = 0u32;

        // 写入所有 data blocks
        for block in &self.data_blocks {
            file.extend_from_slice(block);
        }

        // 写入 index block
        let index_offset = file.len() as u32;
        for (key, block_offset) in &self.index {
            file.extend_from_slice(&(key.len() as u32).to_be_bytes());
            file.extend_from_slice(key);
            file.extend_from_slice(&block_offset.to_be_bytes());
        }

        // 写入 bloom filter
        let bloom_offset = file.len() as u32;
        file.extend_from_slice(&self.bloom.to_bytes());

        // Footer: [data_offset(u32)] [index_offset(u32)] [bloom_offset(u32)]
        file.extend_from_slice(&data_offset.to_be_bytes());
        file.extend_from_slice(&index_offset.to_be_bytes());
        file.extend_from_slice(&bloom_offset.to_be_bytes());
        file
    }
}

6. Compaction:写入放大的根源与平衡艺术

Compaction 是 LSM-Tree 最核心的后台操作——它将上层 SSTable 合并到下层,消除过期数据和删除标记。两种主流策略决定了 LSM-Tree 的性能特征:

Size-Tiered Compaction(RocksDB 的 FIFO/Universal 模式):当同一层级的 SSTable 数量达到阈值(如 4 个)时,将它们合并为一个更大的下一层 SSTable。优势是写入放大低(~1.1x),但读取放大高(需要查多个 SSTable)且空间放大严重。

Leveled Compaction(LevelDB/RocksDB 默认模式):每层有大小限制(L1=10MB, L2=100MB, ...),当某层超出限制时,选取该层的一个 SSTable 与下一层重叠的 SSTable 合并。优势是读取放大低、空间放大小(~1.1x),但写入放大较高(~10x)。

下面是用 Rust 实现的 Compaction 核心逻辑:

use std::sync::Arc;
use crossbeam_channel::{bounded, Sender, Receiver};

pub struct CompactionScheduler {
    levels: Vec<Vec<Arc<SsTable>>>,
    max_sizes: Vec<usize>,
    tx: Sender<CompactionTask>,
}

type CompactionTask = (Vec<Arc<SsTable>>, Vec<Arc<SsTable>>); // (target_level, next_level_overlapping)

impl CompactionScheduler {
    pub fn should_compact(&self, level: usize) -> bool {
        let total_size: usize = self.levels[level].iter().map(|s| s.size()).sum();
        total_size > self.max_sizes.get(level).copied().unwrap_or(usize::MAX)
    }

    /// 选择需要 Compaction 的 SSTable(简化版 Leveled 策略)
    pub fn pick_compaction(&self) -> Option<CompactionTask> {
        for level in 0..self.levels.len() - 1 {
            if self.should_compact(level) {
                // 选择该层第一个 SSTable
                let sst = self.levels[level][0].clone();
                // 查找下一层中 key 范围重叠的 SSTable
                let overlapping = self.find_overlapping(level + 1, &sst);
                return Some((vec![sst], overlapping));
            }
        }
        None
    }

    fn find_overlapping(&self, level: usize, sst: &SsTable) -> Vec<Arc<SsTable>> {
        self.levels[level]
            .iter()
            .filter(|other| {
                key_ranges_overlap(
                    sst.first_key(), sst.last_key(),
                    other.first_key(), other.last_key(),
                )
            })
            .cloned()
            .collect()
    }

    /// 多路归并 Compaction 核心
    pub static fn compact(tables: Vec<Vec<Arc<SsTable>>>) -> Vec<Arc<SsTable>> {
        // 所有输入表的两两归并:创建 MergeIterator
        let merged = MergeIterator::new(
            tables.iter().flatten().map(|t| t.iter()).collect()
        );

        let mut builder = SsTableBuilder::new();
        let mut output_tables = Vec::new();
        let target_table_size = 2 * 1024 * 1024; // 2MB per SSTable

        while let Some((key, value)) = merged.next() {
            builder.add_entry(&key, value.as_ref().map(|v| v.as_slice()));
            if builder.estimated_size() >= target_table_size {
                let bytes = builder.build();
                output_tables.push(Arc::new(SsTable::from_bytes(bytes)));
                builder = SsTableBuilder::new();
            }
        }
        if builder.estimated_size() > 0 {
            let bytes = builder.build();
            output_tables.push(Arc::new(SsTable::from_bytes(bytes)));
        }
        output_tables
    }
}

/// K路归并迭代器
pub struct MergeIterator {
    heap: BinaryHeap<(Vec<u8>, usize, usize)>, // (key, table_index, entry_index)
}

impl Iterator for MergeIterator {
    type Item = (Vec<u8>, Option<Vec<u8>>);

    fn next(&mut self) -> Option<Self::Item> {
        self.heap.pop().map(|(key, table_idx, entry_idx)| {
            // 从对应表推进下一个 entry,重新推入 heap
            // 相同 key 的处理:以最新写入的(table_idx 最小的)为准
            (key, value)
        })
    }
}

7. Bloom Filter:点查加速器

当 key 不存在时,LSM-Tree 需要遍历所有层级的 SSTable 才能确定"查无此 key"。Bloom Filter 以极小的内存开销为每个 SSTable 提供"key 可能存在"或"key 一定不存在"的判断,将大量无效 IO 消灭在查询路径之外。

pub struct BloomFilter {
    bits: Vec<u64>,
    num_bits: usize,
    num_hashes: usize,
}

impl BloomFilter {
    pub fn new(bits_per_key: usize) -> Self {
        Self {
            bits: Vec::new(),
            num_bits: 0,
            num_hashes: (bits_per_key as f64 * 0.693) as usize, // ln(2) ≈ 0.693
        }
    }

    pub fn insert(&mut self, key: &[u8]) {
        let (h1, h2) = Self::hash(key);
        for i in 0..self.num_hashes {
            let bit_idx = ((h1.wrapping_add(h2.wrapping_mul(i as u64))) % self.num_bits as u64) as usize;
            self.bits[bit_idx / 64] |= 1 << (bit_idx % 64);
        }
    }

    pub fn may_contain(&self, key: &[u8]) -> bool {
        let (h1, h2) = Self::hash(key);
        for i in 0..self.num_hashes {
            let bit_idx = ((h1.wrapping_add(h2.wrapping_mul(i as u64))) % self.num_bits as u64) as usize;
            if self.bits[bit_idx / 64] & (1 << (bit_idx % 64)) == 0 {
                return false;
            }
        }
        true
    }

    /// 使用 murmurhash3 的变体:f(k) = h1 + i × h2
    fn hash(key: &[u8]) -> (u64, u64) {
        let h1 = murmur3_x64_128(key, 0).unwrap_or(0) as u64;
        let h2 = murmur3_x64_128(key, h1 as u32).unwrap_or(0) as u64;
        (h1, h2)
    }

    pub fn to_bytes(&self) -> Vec<u8> {
        self.bits.iter().flat_map(|b| b.to_be_bytes()).collect()
    }
}

每个 key 分配 10 bits 的典型配置下,Bloom Filter 的误判率约为 1%,意味着 99% 的不存在 key 查询可以直接跳过该 SSTable,无需任何磁盘 IO。

8. 读取路径:多级查找的完整实现

GET 操作需要从最新到最旧逐级查找,命中后立即返回。这是典型的"空间换时间"——Bloom Filter 加速不存在的 key,分层结构确保找到最新的 value:

impl KvEngine {
    pub fn get(&self, key: &[u8]) -> Option<Vec<u8>> {
        // 1. 查 MemTable(最新写入,最快路径)
        if let Some(value) = self.memtable.get(key) {
            return value; // Option<Vec<u8>>:None 表示墓碑标记(已删除)
        }

        // 2. 查 Immutable MemTable(即将刷盘)
        if let Some(value) = self.immutable_memtables.front().and_then(|m| m.get(key)) {
            return value;
        }

        // 3. 从 L0 到Ln 逐级查找
        for (level, ssts) in self.levels.iter().enumerate() {
            if level == 0 {
                // L0 的数据按写入时间从新到旧排列
                for sst in ssts.iter().rev() {
                    if let Some(value) = self.get_from_sstable(sst, key) {
                        return value;
                    }
                }
            } else {
                // L1+ 的数据按 key 范围不重叠,只需查可能包含 key 的那个 SSTable
                if let Some(sst) = self.find_sstable_in_level(level, key) {
                    if sst.bloom_filter().may_contain(key) {
                        if let Some(value) = self.get_from_sstable(sst, key) {
                            return value;
                        }
                    }
                }
            }
        }
        None
    }

    fn get_from_sstable(&self, sst: &SsTable, key: &[u8]) -> Option<Vec<u8>> {
        // 利用 index block 定位 data block
        let block = sst.find_block(key)?;
        // 在 block 内二分查找
        block.binary_search(key)
    }
}

9. 工程实践:工业级实现的关键优化

从零实现一个能用的 LSM-Tree 引擎,到生产级可用之间还有一段距离。以下是 RocksDB、Pebble(CockroachDB 的引擎)、Titan(腾讯基于 RocksDB 优化的引擎)等工业实践中的核心优化:

9.1 Block Cache 与 Table Cache

将频繁访问的 data block 和 index block 缓存在内存中 (LRU 或 W-TinyLFU 策略),避免重复 SSTable 读取。RocksDB 默认 8MB Block Cache,生产环境通常配置到 GB 级别。

9.2 分离 Key-Value (Titan/Pebble 的 Merge Operator)

对于 value 较大的场景(如 Blob 数据),将 key 存在 LSM-Tree 中而 value 存在独立的 Blob 文件中。这大幅减小了 Compaction 的开销——Compaction 只需搬移小 key,不搬移大 value。Tinn 存储引擎即采用此策略,写入吞吐可提升 5-10 倍。

9.3 IO 优化:Direct IO 与 io_uring

RocksDB 支持 Direct IO,绕过 Page Cache,减少一次内核态到用户态的数据拷贝。在 Linux 5.1+ 中,io_uring 进一步减少了系统调用开销,Pebble 引擎对 io_uring 的利用使其在 NVMe SSD 上的性能超越了 RocksDB。

9.4 前缀压缩与前缀 Bloom Filter

当 key 有公共前缀时(如 user:1001:profile),前缀编码可以大幅减小 SSTable 的体积。按前缀构建的 Bloom Filter 进一步提升了前缀范围查询的效率。

9.5 Compaction 限速与优先级

后台 Compaction 会消耗 IO 带宽,影响前台读写。RocksDB 提供了 max_background_compactions 和 rate_limiter 来动态控制 Compaction 的 IO 配额,确保生产延迟不被 Compaction 抖动。

10. 性能基准:LSM-Tree vs B+Tree

在 Samsung PM9A3 NVMe SSD (QD1, 4KB) 上的典型表现:

操作RocksDB LSM-TreeLMDB B+Tree
随机写 (ops/s)450,00030,000
随机读 (ops/s)50,000180,000
顺序扫描 (MB/s)3,2002,400
写入放大系数8-20x1-2x
空间放大系数1.1-1.5x1.0x

结论:LSM-Tree 在写入密集型场景有数量级的优势,B+Tree 在读密集场景更优。选择引擎前务必明确 workload 特征。

11. 总结

LSM-Tree 的核心魅力在于"用写入换读取"的范式——以 Compaction 的后台代价换取前台写入的高吞吐。在 Rust 生态中,sled、rocksdb、pebble(通过 C FFI 绑定)都是可用的生产级实现。但理解其内部机制,对于调优 RocksDB 参数(max_bytes_for_level_base、target_file_size_base、bloom_filter_bits_per_key)至关重要。

如果你正在构建一个 KV 存储引擎,记住三条黄金法则:(1)测量你的写入模式——point delete 频率决定 Compaction 压力;(2)合理配置 Bloom Filter——每个 key 10 bits 几乎是无损优化;(3)分离大 value——Titan 模式让 Compaction 不再成为瓶颈。LSM-Tree 不会因为新技术而没落,它会像 B+Tree 一样,作为基础架构的"隐形冠军"持续演进。

完整代码实现已开源在 github.com,欢迎 star 和 issue 交流。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部