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-Tree | LMDB B+Tree |
|---|---|---|
| 随机写 (ops/s) | 450,000 | 30,000 |
| 随机读 (ops/s) | 50,000 | 180,000 |
| 顺序扫描 (MB/s) | 3,200 | 2,400 |
| 写入放大系数 | 8-20x | 1-2x |
| 空间放大系数 | 1.1-1.5x | 1.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 交流。

发表评论 取消回复