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/ 集成测试。