Skip to main content

raft_rust/storage/
bitcask.rs

1// 内存 keydir:有序映射键 → 值位置
2use std::collections::BTreeMap;
3// keydir 范围扫描的迭代器类型
4use std::collections::btree_map::Range;
5// 底层日志文件句柄
6use std::fs::File;
7// 缓冲读写与定位
8use std::io::{BufReader, BufWriter, Read as _, Seek as _, SeekFrom, Write as _};
9// 范围边界,用于 scan
10use std::ops::{Bound, RangeBounds};
11// 日志文件路径
12use std::path::PathBuf;
13// 构建 keydir 时用 std IO Result,与库 Result 区分
14use std::result::Result as StdResult;
15
16// 跨进程排他锁,防止双开同一数据文件
17use fs4::fs_std::FileExt;
18// 打开/压缩/截断时的运维日志
19use log::{error, info};
20
21// Engine trait 与对外 Status
22use super::{Engine, Status};
23// 库错误类型
24use crate::error::{Error, Result};
25
26/// BitCask 的极简变体。BitCask 是日志结构键值引擎(如 Riak 所用)。
27/// 与其他实现生成的 BitCask 库不兼容。参见:
28/// <https://riak.com/assets/bitcask-intro.pdf>
29///
30/// BitCask 将键值对追加写入日志文件,并在内存中维护键到文件偏移的映射。
31/// 所有存活键须能装入内存。删除会向日志写入墓碑值。为清理旧垃圾
32///(已删除或被覆盖的键),可通过只写入存活数据的新日志来压缩,丢弃被替换的值与墓碑。
33///
34/// 本实现比标准 BitCask 简单得多:
35///
36/// * 不写多个固定大小日志文件,而是使用单个任意大小的追加日志。
37///   这会增加压缩量(每次压缩须重写整文件),也可能超过文件系统文件大小限制。
38///   不过本库预期数据库较小。
39///
40/// * 压缩期间会锁住数据库的读写。可接受:仅在节点启动时压缩,且文件预期较小。
41///
42/// * 不使用 hint 文件,打开时扫描日志本身构建 keydir。Hint 只省略值,
43///   而本库值预期较小,hint 几乎与压缩后的日志一样大。
44///
45/// * 日志条目不含时间戳或校验和。
46///
47/// 编码后的日志条目结构:
48///
49/// 1. 键长度,大端 u32 [4 字节]。
50/// 2. 值长度,大端 i32;墓碑为 -1 [4 字节]。
51/// 3. 键原始字节 [<= 2 GB]。
52/// 4. 值原始字节 [<= 2 GB]。
53pub struct BitCask {
54    /// 当前追加写日志文件。
55    log: Log,
56    /// 将键映射到 [`BitCask::log`] 中值的偏移与长度。
57    keydir: KeyDir,
58// 结束当前作用域
59}
60
61/// 将键映射到日志文件中值的位置。
62// 使用 BTreeMap 以支持有序范围扫描(Raft 日志按键序迭代)
63type KeyDir = BTreeMap<Vec<u8>, ValueLocation>;
64
65/// 日志文件中值的位置。
66// Copy 便于从 keydir 取出后按位置随机读
67#[derive(Clone, Copy)]
68// 定义数据结构
69struct ValueLocation {
70    /// 值在日志文件中的字节偏移。
71    offset: u64,
72    /// 值的字节长度。
73    length: usize,
74// 结束当前作用域
75}
76
77// 为类型实现方法
78impl ValueLocation {
79    // 值在文件中的结束偏移(不含),用于 EOF 校验
80    fn end(&self) -> u64 {
81        // offset + length
82        self.offset + self.length as u64
83    // 结束当前作用域
84    }
85// 结束当前作用域
86}
87
88// 为类型实现方法
89impl BitCask {
90    /// 在给定路径打开或创建 BitCask 数据库。
91    pub fn new(path: PathBuf) -> Result<Self> {
92        // 打开/创建追加日志
93        let mut log = Log::new(path.clone())?;
94        // 扫描日志重建内存索引
95        let keydir = log.build_keydir()?;
96        // 记录存活键数量,便于启动观测
97        info!("Opened {} with {} live keys", path.display(), keydir.len());
98        // 组装引擎实例
99        Ok(Self { log, keydir })
100    // 结束当前作用域
101    }
102
103    /// 打开 BitCask 数据库;若打开时垃圾占比与字节数超过阈值则自动压缩。
104    pub fn new_maybe_compact(
105        // 列表或字段续项
106        path: PathBuf,
107        // 垃圾占比阈值(0.0–1.0)
108        garbage_min_fraction: f64,
109        // 垃圾绝对字节阈值
110        garbage_min_bytes: u64,
111    // 业务逻辑步骤
112    ) -> Result<Self> {
113        // 先正常打开
114        let mut engine = Self::new(path)?;
115
116        // 读取盘占用与存活统计
117        let status = engine.status()?;
118        // 磁盘总大小(含垃圾)
119        let total_size = status.disk_size;
120        // 垃圾字节数
121        let garbage_size = status.garbage_disk_size();
122        // 垃圾占比
123        let garbage_fraction = garbage_size as f64 / total_size as f64;
124        // 同时满足:有垃圾、字节数与占比都超阈值才压缩
125        if garbage_size > 0
126            // 业务逻辑步骤
127            && garbage_size >= garbage_min_bytes
128            // 业务逻辑步骤
129            && garbage_fraction >= garbage_min_fraction
130        // 进入代码块
131        {
132            // 压缩前记录预期可回收空间
133            info!(
134                // 列表或字段续项
135                "Compacting {} to remove {:.0}% garbage ({:.1} MB out of {:.1} MB)",
136                // 列表或字段续项
137                engine.log.path.display(),
138                // 列表或字段续项
139                garbage_fraction * 100.0,
140                // 列表或字段续项
141                garbage_size as f64 / 1024.0 / 1024.0,
142                // 业务逻辑步骤
143                total_size as f64 / 1024.0 / 1024.0
144            // 结束调用参数列表
145            );
146            // 重写仅含存活键的新日志
147            engine.compact()?;
148            // 压缩后记录目标大小
149            info!(
150                // 列表或字段续项
151                "Compacted {} to size {:.1} MB",
152                // 列表或字段续项
153                engine.log.path.display(),
154                // 业务逻辑步骤
155                (total_size - garbage_size) as f64 / 1024.0 / 1024.0
156            // 结束调用参数列表
157            );
158        // 结束当前作用域
159        }
160
161        // 返回可能已压缩的引擎
162        Ok(engine)
163    // 结束当前作用域
164    }
165// 结束当前作用域
166}
167
168// 实现通用 Engine 接口,供 Raft Log 使用
169impl Engine for BitCask {
170    // 关联扫描迭代器
171    type ScanIterator<'a> = ScanIterator<'a>;
172
173    // 删除:写墓碑并从 keydir 移除
174    fn delete(&mut self, key: &[u8]) -> Result<()> {
175        // 追加墓碑条目(value = None)
176        self.log.write_entry(key, None)?;
177        // 内存索引删除,后续 get 视为不存在
178        self.keydir.remove(key);
179        // 成功返回
180        Ok(())
181    // 结束当前作用域
182    }
183
184    // 将日志刷到稳定存储
185    fn flush(&mut self) -> Result<()> {
186        // 测试中不做 fsync 以加速。在此禁用而非在测试里设
187        // raft::Log::fsync = false,以便断言即使 flush 是空操作也会刷盘。
188        #[cfg(not(test))]
189        // 生产路径:fsync 整个文件
190        self.log.file.sync_all()?;
191        // 成功返回
192        Ok(())
193    // 结束当前作用域
194    }
195
196    // 按 keydir 定位后从日志随机读值
197    fn get(&mut self, key: &[u8]) -> Result<Option<Vec<u8>>> {
198        // 无索引项则键不存在
199        let Some(location) = self.keydir.get(key) else {
200            // 提前返回
201            return Ok(None);
202        // 结束当前作用域
203        };
204        // 按偏移读取值并包装为 Some
205        self.log.read_value(*location).map(Some)
206    // 结束当前作用域
207    }
208
209    // 在 keydir 上做有序范围扫描,值懒加载
210    fn scan(&mut self, range: impl RangeBounds<Vec<u8>>) -> Self::ScanIterator<'_> {
211        // 持有 keydir 范围迭代器 + 日志可变借用
212        ScanIterator { inner: self.keydir.range(range), log: &mut self.log }
213    // 结束当前作用域
214    }
215
216    // dyn Engine 路径:装箱静态 scan 结果
217    fn scan_dyn(
218        // 借用参数
219        &mut self,
220        // 列表或字段续项
221        range: (Bound<Vec<u8>>, Bound<Vec<u8>>),
222    // 业务逻辑步骤
223    ) -> Box<dyn super::ScanIterator + '_> {
224
225        // 复用 scan 实现
226        Box::new(self.scan(range))
227    // 结束当前作用域
228    }
229
230    // 设置/覆盖键:追加新值并更新 keydir
231    fn set(&mut self, key: &[u8], value: Vec<u8>) -> Result<()> {
232        // 写日志并拿到新值位置
233        let value_location = self.log.write_entry(key, Some(&*value))?;
234        // 覆盖内存索引(旧位置成为垃圾,待压缩回收)
235        self.keydir.insert(key.to_vec(), value_location);
236        // 成功返回
237        Ok(())
238    // 结束当前作用域
239    }
240
241    // 汇总逻辑大小与磁盘占用,供 Status 与压缩阈值
242    fn status(&mut self) -> Result<Status> {
243        // 存活键数
244        let keys = self.keydir.len() as u64;
245        // 逻辑大小:各存活键长 + 值长之和
246        let size =
247            // 链式变换结果
248            self.keydir.iter().map(|(key, value_loc)| (key.len() + value_loc.length) as u64).sum();
249        // 文件实际字节数(含垃圾)
250        let disk_size = self.log.file.metadata()?.len();
251        // 存活盘占用估算:逻辑大小 + 每键 8 字节长度前缀
252        let live_disk_size = size + 8 * keys; // 计入长度前缀
253        // 引擎名固定为 bitcask
254        Ok(Status { name: "bitcask".to_string(), keys, size, disk_size, live_disk_size })
255    // 结束当前作用域
256    }
257// 结束当前作用域
258}
259
260// 为类型实现方法
261impl BitCask {
262    /// 压缩当前日志:写出仅含存活键的新日志,并替换当前文件。
263    pub fn compact(&mut self) -> Result<()> {
264        // 创建临时日志文件;若已存在则截断。
265        // 同目录 `.new` 扩展名,便于 rename 替换
266        let new_path = self.log.path.with_extension("new");
267        // 打开临时日志
268        let mut new_log = Log::new(new_path)?;
269        // 确保为空文件
270        new_log.file.set_len(0)?;
271
272        // 将全部存活条目写入新日志,并生成新 KeyDir。
273        let mut new_keydir = KeyDir::new();
274        // 遍历当前存活键
275        for (key, value_loc) in &self.keydir {
276            // 从旧日志读出当前值
277            let value = self.log.read_value(*value_loc)?;
278            // 写入新日志并记录新位置
279            let value_loc = new_log.write_entry(key, Some(&value))?;
280            // 填入新索引
281            new_keydir.insert(key.clone(), value_loc);
282        // 结束当前作用域
283        }
284
285        // 用新日志替换当前日志。
286        // 原子 rename 到原路径
287        std::fs::rename(&new_log.path, &self.log.path)?;
288        // 更新临时 Log 的 path 字段为正式路径
289        new_log.path = self.log.path.clone();
290
291        // 切换到新日志与新 keydir
292        self.log = new_log;
293        // 业务逻辑步骤
294        self.keydir = new_keydir;
295        // 成功返回
296        Ok(())
297    // 结束当前作用域
298    }
299// 结束当前作用域
300}
301
302/// 数据库关闭时尝试刷盘。
303impl Drop for BitCask {
304    // 定义函数
305    fn drop(&mut self) {
306        // 尽力 fsync;失败只记日志,析构不能返回错误
307        if let Err(error) = self.flush() {
308            // 业务逻辑步骤
309            error!("failed to flush file: {}", error)
310        // 结束当前作用域
311        }
312    // 结束当前作用域
313    }
314// 结束当前作用域
315}
316
317/// BitCask 范围扫描迭代器。
318pub struct ScanIterator<'a> {
319    /// keydir 的范围迭代器。
320    inner: Range<'a, Vec<u8>, ValueLocation>,
321    /// 用于按位置读取值的日志引用。
322    log: &'a mut Log,
323// 结束当前作用域
324}
325
326// 为类型实现方法
327impl ScanIterator<'_> {
328    /// 真正从数据文件按位置读取值。
329    fn map(&mut self, item: (&Vec<u8>, &ValueLocation)) -> <Self as Iterator>::Item {
330        // 解构键与位置
331        let (key, value_loc) = item;
332        // 克隆键并从日志读值
333        Ok((key.clone(), self.log.read_value(*value_loc)?))
334    // 结束当前作用域
335    }
336// 结束当前作用域
337}
338
339// 正向迭代:按键升序产出
340impl Iterator for ScanIterator<'_> {
341    // 每项为可能失败的 (key, value)
342    type Item = Result<(Vec<u8>, Vec<u8>)>;
343
344    /// 迭代器逐步从数据文件读取数据。
345    fn next(&mut self) -> Option<Self::Item> {
346        // 取下一个 keydir 项并加载值
347        self.inner.next().map(|item| self.map(item))
348    // 结束当前作用域
349    }
350// 结束当前作用域
351}
352
353// 反向迭代:支持从尾部扫描 Raft 日志等场景
354impl DoubleEndedIterator for ScanIterator<'_> {
355    // 定义函数
356    fn next_back(&mut self) -> Option<Self::Item> {
357        // 从范围高端取项并加载值
358        self.inner.next_back().map(|item| self.map(item))
359    // 结束当前作用域
360    }
361// 结束当前作用域
362}
363
364/// BitCask 追加写日志文件,包含按如下格式编码的键值条目序列:
365///
366/// 1. 键长度,大端 u32 [4 字节]。
367/// 2. 值长度,大端 i32;墓碑为 -1 [4 字节]。
368/// 3. 键原始字节 [<= 2 GB]。
369/// 4. 值原始字节 [<= 2 GB]。
370struct Log {
371    /// 已打开的日志文件。
372    file: File,
373    /// 日志文件路径。
374    path: PathBuf,
375// 结束当前作用域
376}
377
378// 为类型实现方法
379impl Log {
380    /// 打开日志文件;不存在则创建。在关闭前对文件加排他锁;
381    /// 若锁已被持有则报错。
382    fn new(path: PathBuf) -> Result<Self> {
383        // 确保父目录存在
384        if let Some(dir) = path.parent() {
385            // 业务逻辑步骤
386            std::fs::create_dir_all(dir)?
387        // 结束当前作用域
388        }
389        // 读写创建、不截断,以便恢复已有日志
390        let file = std::fs::OpenOptions::new()
391            // 方法链调用
392            .read(true)
393            // 方法链调用
394            .write(true)
395            // 方法链调用
396            .create(true)
397            // 方法链调用
398            .truncate(false)
399            // 方法链调用
400            .open(&path)?;
401        // 排他锁:同一数据目录只允许一个进程打开
402        if !file.try_lock_exclusive()? {
403            // 提前返回
404            return Err(Error::IO(format!("file {path:?} is already is use")));
405        // 结束当前作用域
406        }
407        // 返回打开的日志
408        Ok(Self { file, path })
409    // 结束当前作用域
410    }
411
412    /// 扫描日志文件构建 keydir。若遇到不完整条目,视为不完整写入,截断文件余下部分。
413    fn build_keydir(&mut self) -> Result<KeyDir> {
414        // 复用 4 字节缓冲读长度字段
415        let mut len_buf = [0u8; 4];
416        // 累积最终内存索引
417        let mut keydir = KeyDir::new();
418        // 文件总长度,作为扫描上界
419        let file_len = self.file.metadata()?.len();
420        // 从文件偏移 0 开始,用 BufReader 逐条读取记录。
421        let mut r = BufReader::new(&mut self.file);
422        // 当前条目起始偏移
423        let mut offset = r.seek(SeekFrom::Start(0))?;
424
425        // 顺序扫描直到文件尾
426        while offset < file_len {
427            // 读取下一条目,返回键与值位置;墓碑返回 None。
428            // 闭包内用 std IO 错误,便于识别 UnexpectedEof
429            let result = || -> StdResult<(Vec<u8>, Option<ValueLocation>), std::io::Error> {
430                // 读取键长度:4 字节 u32。
431                r.read_exact(&mut len_buf)?;
432                // 大端键长度
433                let key_len = u32::from_be_bytes(len_buf);
434
435                // 读取值长度:4 字节 i32,墓碑为 -1。
436                r.read_exact(&mut len_buf)?;
437                // 负长度 → 墓碑;否则记录值在文件中的位置
438                let value_loc = match i32::from_be_bytes(len_buf) {
439                    // 值为 -1(..0)表示删除标记,索引会从 keydir 移除该键。
440                    ..0 => None, // 墓碑
441                    // 值从 header(8) + key 之后开始
442                    len => Some(ValueLocation {
443                        // 列表或字段续项
444                        offset: offset + 8 + key_len as u64,
445                        // 列表或字段续项
446                        length: len as usize,
447                    // 列表或字段续项
448                    }),
449                // 结束当前作用域
450                };
451
452                // 读取键。
453                let mut key: Vec<u8> = vec![0; key_len as usize];
454                // 业务逻辑步骤
455                r.read_exact(&mut key)?;
456
457                // 跳过值。
458                if let Some(value_loc) = value_loc {
459                    // 值越过 EOF 视为截断写入
460                    if value_loc.end() > file_len {
461                        // 提前返回
462                        return Err(std::io::Error::new(
463                            // 列表或字段续项
464                            std::io::ErrorKind::UnexpectedEof,
465                            // 列表或字段续项
466                            "value extends beyond end of file",
467                        // 语法续行
468                        ));
469                    // 结束当前作用域
470                    }
471                    // 索引只需知道值位置,用 seek 跳过值内容以加速扫描。
472                    r.seek_relative(value_loc.length as i64)?;
473                // 结束当前作用域
474                }
475
476                // 更新文件偏移。
477                // 下一条起点 = 本条起点 + 头 8 + 键 + 值(墓碑无值)
478                offset += 8 + key_len as u64 + value_loc.map_or(0, |v| v.length) as u64;
479
480                // 返回本条解析结果
481                Ok((key, value_loc))
482            // 业务逻辑步骤
483            }();
484
485            // 用该条目更新 keydir。
486            match result {
487                // 正常 put:覆盖该键最新位置
488                Ok((key, Some(value_loc))) => keydir.insert(key, value_loc),
489                // 墓碑:删除索引
490                Ok((key, None)) => keydir.remove(&key),
491                // 若在文件末尾发现不完整条目,视为不完整写入并截断文件。
492                // 场景:写入中断电可能导致文件末尾只写了一半数据。
493                // 处理:假定末尾不完整数据无效,set_len(offset) 截到最后一条完整记录。
494                Err(err) if err.kind() == std::io::ErrorKind::UnexpectedEof => {
495                    // 记录截断位置
496                    error!("Found incomplete entry at offset {offset}, truncating file");
497                    // 截断到最后完整条目末尾
498                    self.file.set_len(offset)?; // 截断文件
499                    // 结束扫描
500                    break;
501                // 结束当前作用域
502                }
503                // 其它 IO 错误向上传播
504                Err(err) => return Err(err.into()),
505            // 结束当前作用域
506            };
507        // 结束当前作用域
508        }
509
510        // 返回重建的索引
511        Ok(keydir)
512    // 结束当前作用域
513    }
514
515    /// 从日志文件给定位置读取值。
516    fn read_value(&mut self, location: ValueLocation) -> Result<Vec<u8>> {
517        // 预先分配要读取的空间。
518        let mut value: Vec<u8> = vec![0; location.length];
519        // SeekFrom::Start(n) 表示从文件开头向后偏移 n 字节。
520        // seek 将文件读取指针移到指定位置。
521        self.file.seek(SeekFrom::Start(location.offset))?;
522        // read_exact 从当前指针读取,直到填满缓冲区(location.length 字节)。
523        self.file.read_exact(&mut value)?;
524        // 返回读取的值字节
525        Ok(value)
526    // 结束当前作用域
527    }
528
529    /// 向日志追加键值条目;值为 None 表示墓碑。返回条目中值在日志的位置,
530    /// 供 [`KeyDir`] 使用。
531    fn write_entry(&mut self, key: &[u8], value: Option<&[u8]>) -> Result<ValueLocation> {
532        // 本条总字节:8 字节头 + 键 + 值(墓碑无值)
533        let length = 8 + key.len() + value.map_or(0, |v| v.len());
534        // 追加写:定位到文件末尾,即为本条起始 offset
535        let offset = self.file.seek(SeekFrom::End(0))?;
536        // 按条目大小分配写缓冲,减少系统调用
537        let mut w = BufWriter::with_capacity(length, &mut self.file);
538
539        // 键长度:4 字节 u32。
540        w.write_all(&(key.len() as u32).to_be_bytes())?;
541
542        // 值长度:4 字节 i32,墓碑为 -1。
543        w.write_all(&value.map_or(-1, |v| v.len() as i32).to_be_bytes())?;
544
545        // 实际键与值。
546        w.write_all(key)?;
547        // 墓碑时 unwrap_or_default 写空切片
548        w.write_all(value.unwrap_or_default())?;
549        // 冲刷 BufWriter 到文件
550        w.flush()?;
551
552        // 将条目位置转换为值位置。
553        // 值从 header+key 之后开始;墓碑 length=0
554        Ok(ValueLocation {
555            // 列表或字段续项
556            offset: offset + 8 + key.len() as u64,
557            // 列表或字段续项
558            length: value.map_or(0, |v| v.len()),
559        // 语法续行
560        })
561    // 结束当前作用域
562    }
563// 结束当前作用域
564}
565
566
567// 单元测试(goldenscript 套件)已剥离;见 tests/ 集成测试。