基于 Go 标准库的单机嵌入式 LSM-Tree KV 存储引擎,每个 SSTable 文件自带布隆过滤器,零第三方依赖。
# 直接运行完整演示(写入 100 万条 → 修改/删除 → 冷启动恢复 → 点查/范围查验证) go run main.go
注意:演示程序启动时会清空 ./lsm_bloom_database/ 目录,请勿把生产数据放在该路径下。
当前代码在单一 main.go(package main)中,无法直接作为 Go 库被其他项目 import。若需集成,请将引擎部分抽出到独立包(如 store/store.go),main.go 仅保留演示逻辑。以下调用方式对两种形态均适用(以独立包 store 为例):
// MaxMemBytes:MemTable 达到该字节数后自动刷盘(生产建议 64MB+)
store, _ := store.NewLSMStore("./my-db", 64*1024*1024)
defer store.Close()
// 写入(value 需支持 gob 序列化)
store.Put("user:1001", User{Name: "keven", Age: 30})
// 点查(out 类型需与写入时一致)
var u User
store.Get("user:1001", &u)
// 删除(写入墓碑,Compaction 后物理清除)
store.Delete("user:1001")
// 范围查询 [startKey, endKey] 闭区间,结果按 key 升序
users, _ := store.QueryRange("user:1", "user:9999")
| 方法 | 说明 |
|------|------|
| `NewLSMStore(dbDir string, maxMemBytes int64)` | 初始化引擎,自动加载已有 SSTable 并回放 WAL |
| `Put(key string, value interface{})` | 写入/修改(value 经 gob 序列化) |
| `Delete(key string)` | 写入墓碑标记 |
| `Get(key string, out interface{}) (bool, error)` | 点查;key 被删除时返回错误 `"key has been deleted"` |
| `QueryRange(startKey, endKey string) ([]Order, error)` | 闭区间范围查询,自动过滤已删除记录,结果按 key 升序 |
| `Close()` | 停止后台 Compaction,强制刷盘,关闭文件 |
<dbDir>/ ├── active.wal # 活跃 WAL(13 字节头 + key + value,带 CRC32) ├── 000001.sst # 已排序记录 + 尾部布隆过滤器 ├── 000002.sst └── ...
单条记录头(13 字节,大端序):
| 偏移 | 长度 | 字段 |
|------|------|------|
| 0 | 4 | CRC32(覆盖 header[4:] + key + value) |
| 4 | 4 | key 长度 |
| 8 | 4 | value 长度 |
| 12 | 1 | 类型(0=写入,1=删除墓碑) |
源代码 :
package main<>> import ( "bufio" "bytes" "encoding/binary" "encoding/gob" "errors" "fmt" "hash/crc32" "hash/fnv" "io" "os" "sort" "sync" "time" ) // 📥 业务订单结构体 type Order struct { OrderID string UserID int64 Amount float64 IsPaid bool CreatedAt time.Time } type DataEntry struct { Key string Type byte // 0: Put, 1: Delete Value []byte } // ------------------ 🧠 核心:精简版高性能布隆过滤器实现 ------------------ type Filter struct { bitmap []byte kHashes uint8 // 哈希函数的个数,默认选 3 或 4 最佳 } // NewFilter 根据元素数量自动计算并创建极致紧凑的布隆过滤器 (通常每个 Key 仅占 10-12 bits) func NewFilter(numItems int) *Filter { if numItems == 0 { numItems = 1 } // 工业标准:每个 key 分配 10 个 bit,假阳性率约为 1% numBits := numItems * 10 numBytes := (numBits + 7) / 8 return &Filter{ bitmap: make([]byte, numBytes), kHashes: 3, // 使用 3 个不同的扰动哈希函数 } } // baseHashes 使用标准库 FNV 极速生成双哈希底数,用于模拟多哈希 func (f *Filter) baseHashes(key string) (uint32, uint32) { h := fnv.New32a() _, _ = h.Write([]byte(key)) h1 := h.Sum32() // 经典算法:通过低位翻转和异或快速产生互不相关的第二个哈希值 h2 := (h1 >> 17) | (h1 << 15) ^ 0x5bd1e995 return h1, h2 } // Add 将 Key 映射到布隆过滤器的位图中 func (f *Filter) Add(key string) { h1, h2 := f.baseHashes(key) numBits := uint32(len(f.bitmap) * 8) for i := uint32(0); i < uint32(f.kHashes); i++ { // Kirsch-Mitzenmacher 技术:用两个哈希值组合出无数个哈希位置 combinedHash := h1 + i*h2 bitPos := combinedHash % numBits byteIdx := bitPos / 8 bitIdx := bitPos % 8 f.bitmap[byteIdx] |= (1 << bitIdx) } } // Contains 判断 Key 是否存在(如果返回 false,代表 100% 绝对不存在) func (f *Filter) Contains(key string) bool { if len(f.bitmap) == 0 { return false } h1, h2 := f.baseHashes(key) numBits := uint32(len(f.bitmap) * 8) for i := uint32(0); i < uint32(f.kHashes); i++ { combinedHash := h1 + i*h2 bitPos := combinedHash % numBits byteIdx := bitPos / 8 bitIdx := bitPos % 8 if (f.bitmap[byteIdx] & (1 << bitIdx)) == 0 { return false // 只要有一位为 0,绝对不存在 } } return true // 可能存在 } // ------------------ 🏪 LSM 存储引擎与过滤器深度融合 ------------------ type SSTMetaData struct { Path string MinKey string MaxKey string FileSize int64 Filter *Filter // 核心优化:稀疏索引中常驻布隆过滤器,无惧 TB 级内存开销 } type LSMStore struct { mu sync.RWMutex walFile *os.File walWriter *bufio.Writer walPath string dbDir string memTable []DataEntry memBytes int64 maxMemMax int64 sstables []SSTMetaData sstSeq int64 closeChan chan struct{} wg sync.WaitGroup } func NewLSMStore(dbDir string, maxMemBytes int64) (*LSMStore, error) { _ = os.MkdirAll(dbDir, 0755) s := &LSMStore{ dbDir: dbDir, walPath: dbDir + "/active.wal", maxMemMax: maxMemBytes, closeChan: make(chan struct{}), } if err := s.loadSSTables(); err != nil { return nil, err } if err := s.recoverWAL(); err != nil { return nil, err } walFile, err := os.OpenFile(s.walPath, os.O_CREATE|os.O_RDWR|os.O_APPEND, 0666) if err != nil { return nil, err } s.walFile = walFile s.walWriter = bufio.NewWriterSize(walFile, 256*1024) s.wg.Add(1) go s.compactionLoop() return s, nil } func (s *LSMStore) Put(key string, value interface{}) error { return s.writeEntry(key, value, false) } func (s *LSMStore) Delete(key string) error { return s.writeEntry(key, nil, true) } func (s *LSMStore) writeEntry(key string, val interface{}, isDelete bool) error { s.mu.Lock() defer s.mu.Unlock() var vBuf []byte if !isDelete { var buf bytes.Buffer _ = gob.NewEncoder(&buf).Encode(val) vBuf = buf.Bytes() } kBytes := []byte(key) kLen, vLen := uint32(len(kBytes)), uint32(len(vBuf)) header := make([]byte, 13) binary.BigEndian.PutUint32(header[4:8], kLen) binary.BigEndian.PutUint32(header[8:12], vLen) if isDelete { header[12] = 1 } else { header[12] = 0 } crc := crc32.ChecksumIEEE(append(header[4:], append(kBytes, vBuf...)...)) binary.BigEndian.PutUint32(header[0:4], crc) _, _ = s.walWriter.Write(header) _, _ = s.walWriter.Write(kBytes) _, _ = s.walWriter.Write(vBuf) entry := DataEntry{Key: key, Type: header[12], Value: vBuf} s.memTable = append(s.memTable, entry) s.memBytes += int64(13 + kLen + vLen) if s.memBytes >= s.maxMemMax { s.flushMemTable() } return nil } // flushMemTable 将内存固化为具备布隆过滤器的 SSTable 物理文件 func (s *LSMStore) flushMemTable() { if len(s.memTable) == 0 { return } _ = s.walWriter.Flush() sort.Slice(s.memTable, func(i, j int) bool { return s.memTable[i].Key < s.memTable[j].Key }) s.sstSeq++ sstPath := fmt.Sprintf("%s/%06d.sst", s.dbDir, s.sstSeq) file, _ := os.OpenFile(sstPath, os.O_CREATE|os.O_RDWR, 0666) writer := bufio.NewWriterSize(file, 512*1024) // 1. 动态为该批数据初始化一个紧凑的布隆过滤器 filter := NewFilter(len(s.memTable)) for _, entry := range s.memTable { filter.Add(entry.Key) // 将 key 注入过滤器中 kBytes := []byte(entry.Key) kLen, vLen := uint32(len(kBytes)), uint32(len(entry.Value)) header := make([]byte, 13) binary.BigEndian.PutUint32(header[4:8], kLen) binary.BigEndian.PutUint32(header[8:12], vLen) header[12] = entry.Type crc := crc32.ChecksumIEEE(append(header[4:], append(kBytes, entry.Value...)...)) binary.BigEndian.PutUint32(header[0:4], crc) _, _ = writer.Write(header) _, _ = writer.Write(kBytes) _, _ = writer.Write(entry.Value) } // 2. 将计算好的布隆过滤器元数据紧凑地追加在文件的最后面 (Footer 角色) filterOffset, _ := file.Seek(0, io.SeekCurrent) filterOffset += int64(writer.Buffered()) // 纠正缓冲区带来的偏移量差异 // 写入过滤器基本参数 _ = binary.Write(writer, binary.BigEndian, uint32(len(filter.bitmap))) _, _ = writer.Write(filter.bitmap) // 在整个文件的最后 8 字节,写入布隆过滤器的起始偏移量,方便冷启动定位 _ = binary.Write(writer, binary.BigEndian, int64(filterOffset)) _ = writer.Flush() fi, _ := file.Stat() _ = file.Close() s.sstables = append(s.sstables, SSTMetaData{ Path: sstPath, MinKey: s.memTable[0].Key, MaxKey: s.memTable[len(s.memTable)-1].Key, FileSize: fi.Size(), Filter: filter, }) _ = s.walFile.Close() _ = os.Remove(s.walPath) walFile, _ := os.OpenFile(s.walPath, os.O_CREATE|os.O_RDWR|os.O_APPEND, 0666) s.walFile = walFile s.walWriter = bufio.NewWriterSize(walFile, 256*1024) s.memTable = make([]DataEntry, 0) s.memBytes = 0 } // Get 完美融合布隆过滤器与稀疏索引 func (s *LSMStore) Get(key string, out interface{}) (bool, error) { s.mu.RLock() // 1. 检索内存 for i := len(s.memTable) - 1; i >= 0; i-- { if s.memTable[i].Key == key { if s.memTable[i].Type == 1 { s.mu.RUnlock() return false, errors.New("key has been deleted") } err := gob.NewDecoder(bytes.NewReader(s.memTable[i].Value)).Decode(out) s.mu.RUnlock() return true, err } } sstFiles := make([]SSTMetaData, len(s.sstables)) copy(sstFiles, s.sstables) s.mu.RUnlock() // 2. 检索磁盘有序文件 for i := len(sstFiles) - 1; i >= 0; i-- { sst := sstFiles[i] // 🔴 智能拦截关卡 1:检查稀疏边界 if key < sst.MinKey || key > sst.MaxKey { continue } // 🔴 智能拦截关卡 2:布隆过滤器闪电判定(TB 级高并发检索的核心) if sst.Filter != nil && !sst.Filter.Contains(key) { // 如果过滤器说没有,100% 绝对没有,直接彻底免除本次昂贵的磁盘 I/O 寻道! continue } // 通过双重拦截,极高概率确定在此文件中,进入磁盘 found, deleted, vBuf, _ := s.searchInSSTFile(sst.Path, key) if found { if deleted { return false, errors.New("key has been deleted") } err := gob.NewDecoder(bytes.NewReader(vBuf)).Decode(out) return true, err } } return false, errors.New("key not found") } // loadSSTables 冷启动补充:在不读取主体订单数据的情况下,优雅抽取出文件尾部的布隆过滤器 func (s *LSMStore) loadSSTables() error { files, err := os.ReadDir(s.dbDir) if err != nil { return nil } var sstPaths []string for _, f := range files { if !f.IsDir() && len(f.Name()) > 4 && f.Name()[len(f.Name())-4:] == ".sst" { sstPaths = append(sstPaths, s.dbDir+"/"+f.Name()) } } sort.Strings(sstPaths) for _, path := range sstPaths { file, err := os.OpenFile(path, os.O_RDWR, 0666) if err != nil { return err } fi, _ := file.Stat() if fi.Size() < 8 { file.Close() continue } // A. 首先读取文件最后 8 字节,获取布隆过滤器的物理 Offset _, _ = file.Seek(-8, io.SeekEnd) var filterOffset int64 _ = binary.Read(file, binary.BigEndian, &filterOffset) // B. 移动到过滤器所在位置,加载过滤器元数据 _, _ = file.Seek(filterOffset, io.SeekStart) reader := bufio.NewReader(file) var bitmapLen uint32 _ = binary.Read(reader, binary.BigEndian, &bitmapLen) bitmap := make([]byte, bitmapLen) _, _ = io.ReadFull(reader, bitmap) filter := &Filter{bitmap: bitmap, kHashes: 3} // C. 快速提取头部 MinKey 和 MaxKey,流程与前版本一致 _, _ = file.Seek(0, io.SeekStart) headReader := bufio.NewReader(file) header := make([]byte, 13) var minKey, maxKey string isFirst := true // 流式读取直至撞到布隆过滤器的边界,确保 100% 安全重建 var curOffset int64 for curOffset < filterOffset { _, err := io.ReadFull(headReader, header) if err != nil { break } kLen := binary.BigEndian.Uint32(header[4:8]) vLen := binary.BigEndian.Uint32(header[8:12]) payload := make([]byte, kLen+vLen) _, _ = io.ReadFull(headReader, payload) curKey := string(payload[:kLen]) if isFirst { minKey = curKey isFirst = false } maxKey = curKey curOffset += int64(13 + kLen + vLen) } file.Close() if !isFirst { s.sstables = append(s.sstables, SSTMetaData{ Path: path, MinKey: minKey, MaxKey: maxKey, FileSize: fi.Size(), Filter: filter, // 完美还原 }) var seq int64 _, _ = fmt.Sscanf(path, s.dbDir+"/%06d.sst", &seq) if seq > s.sstSeq { s.sstSeq = seq } } } return nil } // [此处的 searchInSSTFile, compactionLoop, recoverWAL, Close 保持逻辑一致,Compaction 写入新文件时同样按上述 Footer 格式附带布隆过滤器即可] func (s *LSMStore) searchInSSTFile(path string, target string) (bool, bool, []byte, error) { file, _ := os.Open(path) defer file.Close() // 获取文件的布隆过滤器起始线,防止扫描越界 _, _ = file.Seek(-8, io.SeekEnd) var filterOffset int64 _ = binary.Read(file, binary.BigEndian, &filterOffset) _, _ = file.Seek(0, io.SeekStart) reader := bufio.NewReader(file) header := make([]byte, 13) var curOffset int64 for curOffset < filterOffset { if _, err := io.ReadFull(reader, header); err != nil { break } kLen := binary.BigEndian.Uint32(header[4:8]) // === 补全开始:承接上文的 searchInSSTFile 解析循环 === vLen := binary.BigEndian.Uint32(header[8:12]) t := header[12] // 提取操作类型 (0: Put, 1: Delete) payload := make([]byte, kLen+vLen) if _, err := io.ReadFull(reader, payload); err != nil { break } curKey := string(payload[:kLen]) if curKey == target { // 找到了目标 Key,判断它是否是被删除的墓碑标记 return true, t == 1, payload[kLen:], nil } // 因为 SSTable 文件内部的 Key 是绝对升序排列的, // 一旦当前遇到的 curKey 已经大于目标 target,说明后面不可能再有了,提前终止,减少磁盘 I/O if curKey > target { break } curOffset += int64(13 + kLen + vLen) } return false, false, nil, nil } // compactionLoop 核心空间回收:多路有序合并,粉碎死数据与删除标记,控制 TB 级体积 func (s *LSMStore) compactionLoop() { defer s.wg.Done() ticker := time.NewTicker(5 * time.Second) defer ticker.Stop() for { select { case <-ticker.C: s.mu.Lock() // 当磁盘上积累了 3 个以上的有序小 SST 文件时,触发合并 if len(s.sstables) < 3 { s.mu.Unlock() continue } fmt.Println("\n🧹 [Compaction] 正在执行多路归并,粉碎历史删除/修改死数据,并重构新布隆过滤器...") allEntries := make(map[string]DataEntry) // 顺序合并所有老文件。因为 sstables 切片中老文件在前、新文件在后, // 新文件的相同 Key 会在 map 赋值时天然覆盖老文件的旧数据(实现修改生效) for _, sst := range s.sstables { file, _ := os.Open(sst.Path) _, _ = file.Seek(-8, io.SeekEnd) var filterOffset int64 _ = binary.Read(file, binary.BigEndian, &filterOffset) _, _ = file.Seek(0, io.SeekStart) reader := bufio.NewReader(file) header := make([]byte, 13) var curOffset int64 for curOffset < filterOffset { _, _ = io.ReadFull(reader, header) kLen := binary.BigEndian.Uint32(header[4:8]) vLen := binary.BigEndian.Uint32(header[8:12]) t := header[12] payload := make([]byte, kLen+vLen) _, _ = io.ReadFull(reader, payload) kStr := string(payload[:kLen]) allEntries[kStr] = DataEntry{Key: kStr, Type: t, Value: payload[kLen:]} curOffset += int64(13 + kLen + vLen) } file.Close() _ = os.Remove(sst.Path) // 物理删除被合并的老文件,释放磁盘空间 } // 将合并去重后的所有 Key 进行严格的字典序升序排序 var sortedKeys []string for k := range allEntries { sortedKeys = append(sortedKeys, k) } sort.Strings(sortedKeys) s.sstSeq++ mergedPath := fmt.Sprintf("%s/%06d.sst", s.dbDir, s.sstSeq) file, _ := os.OpenFile(mergedPath, os.O_CREATE|os.O_RDWR, 0666) writer := bufio.NewWriter(file) // 过滤掉带有删除标记 (Tombstone) 的死数据,将其在物理磁盘上彻底抹去! var activeEntries []DataEntry for _, k := range sortedKeys { if allEntries[k].Type == 1 { continue // 核心:在这里粉碎被用户删除的数据 } activeEntries = append(activeEntries, allEntries[k]) } // 如果合并后还有存活数据,重写成一个新的超级 SSTable 并附带最新的紧凑布隆过滤器 if len(activeEntries) > 0 { newFilter := NewFilter(len(activeEntries)) for _, entry := range activeEntries { newFilter.Add(entry.Key) kBytes := []byte(entry.Key) kLen, vLen := uint32(len(kBytes)), uint32(len(entry.Value)) header := make([]byte, 13) binary.BigEndian.PutUint32(header[4:8], kLen) binary.BigEndian.PutUint32(header[8:12], vLen) header[12] = entry.Type crc := crc32.ChecksumIEEE(append(header[4:], append(kBytes, entry.Value...)...)) binary.BigEndian.PutUint32(header[0:4], crc) _, _ = writer.Write(header) _, _ = writer.Write(kBytes) _, _ = writer.Write(entry.Value) } // 追加写入全新的布隆过滤器到文件末尾 fOffset, _ := file.Seek(0, io.SeekCurrent) fOffset += int64(writer.Buffered()) _ = binary.Write(writer, binary.BigEndian, uint32(len(newFilter.bitmap))) _, _ = writer.Write(newFilter.bitmap) _ = binary.Write(writer, binary.BigEndian, int64(fOffset)) _ = writer.Flush() fi, _ := file.Stat() file.Close() // 原子替换内存里的稀疏快照索引,内存开销瞬间缩减回 1 个文件 s.sstables = []SSTMetaData{{ Path: mergedPath, MinKey: activeEntries[0].Key, MaxKey: activeEntries[len(activeEntries)-1].Key, FileSize: fi.Size(), Filter: newFilter, }} fmt.Printf("✅ [布隆 Compaction] 归并重构完成。生成超压缩主文件: %s\n", mergedPath) } else { file.Close() _ = os.Remove(mergedPath) s.sstables = nil } s.mu.Unlock() case <-s.closeChan: return } } } // recoverWAL 崩溃恢复:将未刷盘的 active.wal 捞回并固化 func (s *LSMStore) recoverWAL() error { if _, err := os.Stat(s.walPath); os.IsNotExist(err) { return nil } walFile, _ := os.OpenFile(s.walPath, os.O_RDONLY, 0666) defer walFile.Close() reader := bufio.NewReader(walFile) header := make([]byte, 13) for { if _, err := io.ReadFull(reader, header); err != nil { break } expectCrc := binary.BigEndian.Uint32(header[0:4]) kLen := binary.BigEndian.Uint32(header[4:8]) vLen := binary.BigEndian.Uint32(header[8:12]) t := header[12] payload := make([]byte, kLen+vLen) if _, err := io.ReadFull(reader, payload); err != nil { break // 截断安全 } if crc32.ChecksumIEEE(append(header[4:], payload...)) != expectCrc { break // 校验失败中断 } s.memTable = append(s.memTable, DataEntry{ Key: string(payload[:kLen]), Type: t, Value: payload[kLen:], }) s.memBytes += int64(13 + kLen + vLen) } if len(s.memTable) > 0 { s.flushMemTable() } return nil } // QueryRange 范围查询(TB 级数据库核心:利用稀疏索引锁定边界进行多路流式扫描) func (s *LSMStore) QueryRange(startKey, endKey string) ([]Order, error) { s.mu.RLock() defer s.mu.RUnlock() // 聚合临时结果哈希表以去重(最新文件或内存的数据会天然覆盖老文件的数据) mergedMap := make(map[string]DataEntry) // 1. 扫描内存有序表 (MemTable) 中的最新存活数据 for _, entry := range s.memTable { if entry.Key >= startKey && entry.Key <= endKey { mergedMap[entry.Key] = entry } } // 2. 扫描命中的磁盘不变量有序文件 (SSTable) for _, sst := range s.sstables { // 智能加速判定:如果当前 SST 物理文件的范围与查询的 start/end 范围完全没有交集,直接过滤掉! if sst.MaxKey < startKey || sst.MinKey > endKey { continue } file, err := os.Open(sst.Path) if err != nil { continue } // 获取该 SST 文件的布隆过滤器起始线,防止扫描越界 _, _ = file.Seek(-8, io.SeekEnd) var filterOffset int64 _ = binary.Read(file, binary.BigEndian, &filterOffset) _, _ = file.Seek(0, io.SeekStart) reader := bufio.NewReader(file) header := make([]byte, 13) // 4B CRC + 4B KLen + 4B VLen + 1B Type var curOffset int64 // 顺序扫描当前物理文件直到遇到布隆过滤器 Footer 边界 for curOffset < filterOffset { if _, err := io.ReadFull(reader, header); err != nil { break } kLen := binary.BigEndian.Uint32(header[4:8]) vLen := binary.BigEndian.Uint32(header[8:12]) t := header[12] // ✨ 核心修复:精准提取第 12 位索引上的 1 字节 byte 值 payload := make([]byte, kLen+vLen) if _, err = io.ReadFull(reader, payload); err != nil { break } curKey := string(payload[:kLen]) // 过滤出符合范围条件的数据块 if curKey >= startKey && curKey <= endKey { mergedMap[curKey] = DataEntry{ Key: curKey, Type: t, // ✨ 此时 t 的类型为 byte,完美契合结构体字面量 Value: payload[kLen:], } } // 因为文件是字典升序的,如果当前 key 已经越过了最大边界,可以直接跳出本文件的检索 if curKey > endKey { break } curOffset += int64(13 + kLen + vLen) } file.Close() } // 3. 过滤掉被标记删除的墓碑数据 (Type == 1),并反序列化还原成最终的订单切片 var result []Order for _, entry := range mergedMap { if entry.Type == 1 { continue // 被逻辑删除的订单不予返回 } var o Order _ = gob.NewDecoder(bytes.NewReader(entry.Value)).Decode(&o) result = append(result, o) } // 为返回结果进行按 Key 字典序升序重排序(保证范围查询结果的确定有序性) sort.Slice(result, func(i, j int) bool { return result[i].OrderID < result[j].OrderID }) return result, nil } // Close 安全关闭存储 func (s *LSMStore) Close() error { close(s.closeChan) s.wg.Wait() s.mu.Lock() defer s.mu.Unlock() s.flushMemTable() // 强制执行最后一次内存数据固化 if s.walFile != nil { _ = s.walFile.Close() } return nil }