基于 Go 标准库的单机嵌入式 LSM-Tree KV 存储引擎

这是一个 TB 级 数据存储 研究,基于golang 和 lsm 技术,查询采用布隆过滤器

TB-Store

基于 Go 标准库的单机嵌入式 LSM-Tree KV 存储引擎,每个 SSTable 文件自带布隆过滤器,零第三方依赖。

核心特性

  • LSM-Tree 架构:MemTable → WAL → SSTable 分层写入,写入全部追加式,对磁盘友好
  • 布隆过滤器加速点查:每个 SSTable 尾部内嵌位图过滤器(双哈希 / Kirsch-Mitzenmacher),key 不存在时 100% 拦截,省去无效磁盘 I/O
  • 墓碑删除:删除操作写入 Type=1 标记,读取时拦截,Compaction 时物理清除
  • 崩溃恢复:WAL 记录带 CRC32 校验,重启时自动回放并固化
  • 后台 Compaction:SSTable 累积到 3 个时每 5 秒触发多路归并,合并数据、丢弃墓碑、重建过滤器
  • 范围查询:利用每个文件的 MinKey/MaxKey 稀疏索引定位边界,多文件流式扫描

快速开始

# 直接运行完整演示(写入 100 万条 → 修改/删除 → 冷启动恢复 → 点查/范围查验证)
go run main.go
注意:演示程序启动时会清空  ./lsm_bloom_database/  目录,请勿把生产数据放在该路径下。


使用方式

当前代码在单一 main.gopackage 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")

API

| 方法 | 说明 |

|------|------|

| `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
}


© GVGNN 2013-2026