storage

package
v0.1.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 25, 2026 License: Apache-2.0, MIT Imports: 22 Imported by: 0

Documentation

Overview

本文件是 SSTable 的类型与共用定义:磁盘布局常量、块索引、以及 SSTable 本身。 读路径见 sstable_read.go,写路径与 compaction 见 sstable_write.go,元信息管理见 sstable_meta.go。

本文件是 SSTable 的元信息管理:启动时扫描目录重建元信息,以及增删查与不可变快照发布。

本文件是 SSTable 的读路径:常驻句柄、块索引与布隆缓存、点查与全量读。

本文件是 SSTable 的写路径:flush 落盘、写尾(块索引 + 布隆 + Footer)与 compaction 合并。

Index

Constants

View Source
const (
	WALOpPut    uint8 = 1
	WALOpDelete uint8 = 2
)

WAL 操作码

View Source
const DefaultNamespaceSep = ':'

DefaultNamespaceSep 是默认的命名空间分隔符。 数据仓库场景下 key 常带"仓库"前缀,如 "log:2026-05-31"、"order:10086"。

View Source
const (
	SSTableBlockSize = 64
)

Variables

View Source
var (
	// ErrKeyNotFound 表示 key 不存在,或其最新版本是一个删除墓碑。
	// 它是正常的查询结果而非故障,调用方通常应据此返回「无此键」而不是「内部错误」。
	ErrKeyNotFound = errors.New("storage: key not found")

	// ErrMemTableUnavailable 表示 memtable 尚未初始化或已关闭,此时读写均无法进行。
	// 与 ErrKeyNotFound 不同,它指示的是引擎状态异常。
	ErrMemTableUnavailable = errors.New("storage: memtable unavailable")

	// ErrNoEntries 表示待写入 SSTable 的条目集为空,本次落盘无需进行。
	ErrNoEntries = errors.New("storage: no entries to write")
)

存储层的哨兵错误。调用方一律用 errors.Is 判别,不得比较错误文本。

哨兵的作用在于让「key 不存在」与「读盘失败」等真实故障可被判别:若二者都是匿名的 errors.New,调用方看到的只是「某个 error」,只能一视同仁当作失败处理。

错误文本遵循 Go 的约定:小写开头、无尾随标点,并以包名前缀标明来源 (对照 Pebble 的 base.ErrNotFound、Badger 的 ErrKeyNotFound)。

Functions

func ResetBlockCacheStats

func ResetBlockCacheStats()

ResetBlockCacheStats 清零计数(供 benchmark 隔离测量)。

func ResetCompactionStats

func ResetCompactionStats()

ResetCompactionStats 清零计数(供 benchmark 隔离测量)。

Types

type BlockCacheStats

type BlockCacheStats struct {
	Hits      int64
	Misses    int64
	Evictions int64
}

BlockCacheStats 是某一刻的数据块缓存计数。

func ReadBlockCacheStats

func ReadBlockCacheStats() BlockCacheStats

ReadBlockCacheStats 读取当前数据块缓存计数快照。

type BlockIndexEntry

type BlockIndexEntry struct {
	LastKey     []byte
	BlockOffset int64
}

type BloomFilter

type BloomFilter struct {
	// contains filtered or unexported fields
}

BloomFilter 标准布隆过滤器。 位数组大小 m 与哈希函数个数 k 由预期元素数 n 和目标误判率 p 按最优公式计算得出,避免拍脑袋取值导致空间浪费或误判率失控。

func DecodeBloomFilter

func DecodeBloomFilter(data []byte) (*BloomFilter, error)

DecodeBloomFilter 从 Encode 的字节还原。

func NewBloomFilter

func NewBloomFilter(n int, p float64) *BloomFilter

NewBloomFilter 按 n 个元素、目标误判率 p 计算最优 m、k:

m = ceil(-n * ln(p) / (ln2)^2)   位数
k = round((m/n) * ln2)           哈希函数个数

n<=0 视为 1;p 不在 (0,1) 区间时回退到 0.01。

func (*BloomFilter) Add

func (b *BloomFilter) Add(key []byte)

Add 插入一个 key。

func (*BloomFilter) Encode

func (b *BloomFilter) Encode() []byte

Encode 序列化为字节:[m(8B)][k(8B)][bits...],大端。

func (*BloomFilter) MayContain

func (b *BloomFilter) MayContain(key []byte) bool

MayContain 判断 key 是否可能存在。 返回 false 一定不存在;返回 true 可能存在(存在误判率 p)。

type CompactionStats

type CompactionStats struct {
	FlushBytes      int64 // flush(memtable → L0 SSTable)写出的字节
	CompactionBytes int64 // compaction 写出的字节
}

CompactionStats 是某一刻的写放大统计。

func ReadCompactionStats

func ReadCompactionStats() CompactionStats

ReadCompactionStats 读取当前写放大计数快照。

type Engine

type Engine struct {
	// contains filtered or unexported fields
}

Engine 是 LSM 存储引擎:它同时持有内存中的表与磁盘上的 SSTable 集合,并驱动二者 之间的流转。

注意区分:这里的「内存表」是 SkipList(见 skiplist.go),Engine 是持有它们并连同 SSTable 一起管理的整条链路。

内存侧采用双表:

  • active 接收所有 Put/Delete;
  • dirty 是正在 flush 的不可变快照,flush 完成后置 nil。

双表使 flush 能在锁外进行:交换后 dirty 不再被写入,Get 仍可回退查询它, 从而避免 flush 与写入相互阻塞。

读顺序 active → dirty → SSTable;写入受字节级信用背压约束,超预算即阻塞等待 flush 归还信用。后台两个协程分别负责 flush 与 compaction。

func NewEngine

func NewEngine(opts Options) *Engine

SkipNode 跳表节点 NewEngine 按 opts 构造存储引擎。参数在此固定,此后不再读全局配置。 需要沿用进程配置时传 DefaultOptions()。

func (*Engine) Close

func (m *Engine) Close() error

delete 从跳表中删除节点,返回是否成功删除

func (*Engine) CompactSSTable

func (m *Engine) CompactSSTable(startLevel int)

func (*Engine) Delete

func (m *Engine) Delete(key []byte) error

insert 在跳表中插入键值对(无锁版本,由调用者保证线程安全) insert 插入或覆盖 key,并返回本次操作使 byteSize 变化的增量(覆盖写可能为负)。 调用方用该增量做背压信用对账。

func (*Engine) Flush

func (m *Engine) Flush()

Flush 将 dirty 表数据刷入 SSTable 流程:

  1. 持锁交换 active → dirty(active 变为 dirty 的不可变快照)
  2. 创建新的空 active 表用于接受后续写入
  3. 释放锁,在锁外将 dirty 数据写入 SSTable
  4. flush 完成后将 dirty 置 nil

func (*Engine) FlushToSSTable

func (m *Engine) FlushToSSTable(entries []LogEntry) error

FlushToSSTable 将 entries 写入临时跳表并立即 Flush 到 SSTable 不经过 active 表,不影响正常读写,专用于快照重放等场景

func (*Engine) FlushWorker

func (m *Engine) FlushWorker()

func (*Engine) Get

func (m *Engine) Get(key []byte) ([]byte, error)

Get 获取指定 key 的值 查找顺序:active → dirty → SSTable

func (*Engine) InflightBytes

func (m *Engine) InflightBytes() int64

InflightBytes 返回当前未 flush(active + 正在刷的 dirty)占用的字节信用,供观测/压测使用。

func (*Engine) ListenCompactCh

func (m *Engine) ListenCompactCh()

func (*Engine) Put

func (m *Engine) Put(key []byte, value []byte) error

search 在跳表中查找指定 key,返回值和是否找到

func (*Engine) ReclaimUpTo

func (m *Engine) ReclaimUpTo(bound []byte) int

ReclaimUpTo 丢弃 MaxKey 严格小于 bound 的 SSTable 整个文件,返回回收的文件数。

用于「投递前置缓冲」的保留期回收:投递按 key 升序推进游标,故 bound 之前的数据已全部 被投递读过,整份文件不再需要留在本地。

之所以按整文件丢弃、而不是逐 key 写墓碑:墓碑会让写入量翻倍,且自身还要再经一轮 compaction 才消失;而文件级丢弃是 O(1),无写放大。

保守之处(宁可少回收,不可误删):

  • 仅在 MaxKey 可信时回收。没有可读 footer 的文件(老格式或尾部残缺)MaxKey 未知,一律跳过。
  • 用严格小于:恰好含 bound 的文件保留,因为 bound 本身尚未被投递消费。
  • bound 为空表示尚无已提交游标,不回收任何文件。

与 compaction 互斥(共用 fileMu),避免删掉 compaction 正在读的源文件。

func (*Engine) ScanRange

func (m *Engine) ScanRange(start, end []byte, fn func(key, value []byte) bool)

ScanRange 在 [start,end] 闭区间升序遍历最新可见键值,跳过墓碑,对每条命中调用 fn; fn 返回 false 可提前停止。start/end 为空分别表示下界/上界不限。

覆盖 active + dirty + 全部 SSTable。此前它只遍历两张内存表,已 flush 的数据对扫描完全 不可见——SCAN 命令因此只返回热数据,而按游标取数的下游投递会跳过所有已落盘的记录。

实现为多路归并:内存表在锁内按范围拷出快照,SSTable 以文件迭代器参与,源序按新旧排列 (SSTable 由旧到新,其后 dirty,最后 active),故同一 key 只保留最新版本;最新版本是 墓碑时整条跳过。

拷贝内存表而非全程持锁,是因为归并要读磁盘:持锁会让写入停等 I/O,无锁遍历跳表又会与 并发写相争。拷贝量由内存表大小天然有界。

fn 内若需在调用返回后继续持有 key/value,应自行拷贝。

func (*Engine) Size

func (m *Engine) Size() int

newSkipNode 创建新的跳表节点

func (*Engine) SnapshotLive

func (m *Engine) SnapshotLive() []LogEntry

SnapshotLive 返回 active+dirty 合并后的全部键值快照(active 覆盖 dirty), 与 ScanRange 不同:保留墓碑(value==nil)。用于 WAL checkpoint 重写——未 flush 的 热数据(含删除墓碑)必须完整保留,否则重放时被删的 key 会从 SSTable 复活。 拷贝底层字节,返回后可安全持有。

func (*Engine) StartFlush

func (m *Engine) StartFlush()

type LogEntry

type LogEntry struct {
	Key   []byte
	Value []byte
}

type Options

type Options struct {
	// Dir 是 SSTable 文件目录。
	Dir string

	// MaxMemTableSize 是 active 表的条目数阈值,超过即触发 flush。
	MaxMemTableSize int

	// MaxCompactionSize 是单层文件数阈值,达到即触发该层 compaction。
	MaxCompactionSize int

	// MaxInflightBytes 是未 flush 数据(active + 正在 flush 的 dirty)的字节预算,
	// 超出即阻塞写入等待 flush 归还信用。<=0 关闭背压。
	MaxInflightBytes int64

	// BlockCacheBytes 是 SSTable 数据块缓存的字节预算。<=0 关闭缓存。
	BlockCacheBytes int64

	// SkipListMaxLevel 与 SkipListP 是跳表的最大层高与升层概率。
	SkipListMaxLevel int
	SkipListP        float64
}

Options 是存储引擎的全部可调参数,构造时一次性传入。

之所以显式传参而非在包内读全局 config.G:全局配置让同一进程内无法并存两套不同配置的 引擎(多节点集成测试正需要),也让测试之间经由全局变量互相影响——本包此前正因如此 出现过偶发失败。它还迫使生产代码写防御性代码:构造时把配置「快照」一份,以避开与测试 中并发改配置形成的数据竞争。参数一旦由调用方传入,这些问题都不复存在。

func DefaultOptions

func DefaultOptions() Options

DefaultOptions 从全局配置取一份参数。

这是全局配置进入存储层的唯一入口:调用方在构造时读一次,此后引擎只认自己那份参数, 不再受 config.G 后续变动影响。

func OptionsFromConfig

func OptionsFromConfig(cfg *config.GlobalConfig) Options

OptionsFromConfig 从一份显式配置构造 Options——与 DefaultOptions 的唯一区别 是数据来源(显式 cfg 而不是包级 config.G)。kairosflux 引擎 (service/kairosflux.go)为每个实例使用独立数据目录构造存储时走这里, 与 DefaultOptions 保持同一份字段映射,不另抄一遍。

type PartitionedBloom

type PartitionedBloom struct {
	// contains filtered or unexported fields
}

PartitionedBloom 按 key 的命名空间前缀("仓库")分区,每个仓库一个 独立的 BloomFilter,例如 "log"、"order" 各自一份,互不污染。

分区内部采用「前缀删除」:把仓库前缀去掉后只对区分性后缀做哈希。 数据仓库里同一仓库的 key 高度相似(共享长前缀),删除冗余前缀既省去 无意义的哈希输入,又让命名空间本身成为第一道过滤——查 "order:x" 时 若 order 仓库根本不存在,可立即返回 false,不会被 log 仓库的 key 干扰。

func BuildPartitionedBloom

func BuildPartitionedBloom(keys [][]byte, sep byte, p float64) *PartitionedBloom

BuildPartitionedBloom 根据全部 key 一次性构建:先按仓库预统计元素数, 再为每个仓库按其各自的 n 计算最优 m、k,最后前缀删除后插入。 适合 SSTable 这类 key 集合已知且构建后不再变更的场景,sizing 比 Add 逐个懒创建(共用一个 n)更精确。

func DecodePartitionedBloom

func DecodePartitionedBloom(data []byte, nPerPartition int, p float64) (*PartitionedBloom, error)

DecodePartitionedBloom 从 Encode 的字节还原。n/p 仅用于后续新增分区, 已有分区从字节恢复其原始 m、k。

func NewPartitionedBloom

func NewPartitionedBloom(sep byte, nPerPartition int, p float64) *PartitionedBloom

NewPartitionedBloom 创建分区布隆过滤器。 nPerPartition / p 用于按需创建各分区的子过滤器(最优 m、k 计算)。

func (*PartitionedBloom) Add

func (pb *PartitionedBloom) Add(key []byte)

Add 按命名空间路由到对应分区,前缀删除后插入后缀。

func (*PartitionedBloom) Encode

func (pb *PartitionedBloom) Encode() []byte

Encode 序列化:[sep(1B)][partCount(4B)] 然后每个分区 [nsLen(4B)][ns][bloomLen(4B)][bloomBytes],大端。

func (*PartitionedBloom) MayContain

func (pb *PartitionedBloom) MayContain(key []byte) bool

MayContain 命名空间不存在直接返回 false;否则在对应分区查后缀。

func (*PartitionedBloom) Namespaces

func (pb *PartitionedBloom) Namespaces() []string

Namespaces 返回当前已建立的所有命名空间(仓库)。

type SSTable

type SSTable struct {
	// contains filtered or unexported fields
}

blockExtent 返回第 i 块在文件中的 [start,end) 字节范围。

func NewSSTable

func NewSSTable(opts Options) *SSTable

NewSSTable 按 opts 构造 SSTable 集合的管理者。

func (*SSTable) AddMeta

func (ss *SSTable) AddMeta(meta *SSTableMeta)

func (*SSTable) DeleteSSTable

func (ss *SSTable) DeleteSSTable(meta *SSTableMeta)

func (*SSTable) LevelFiles

func (ss *SSTable) LevelFiles(level int) []*SSTableMeta

LevelFiles 获取指定层级的文件列表

func (*SSTable) LoadSSTableMetaList

func (ss *SSTable) LoadSSTableMetaList()

openFile 取该路径的常驻只读句柄,未缓存则打开并缓存。返回 nil 表示打开失败。

func (*SSTable) MergeSSTable

func (ss *SSTable) MergeSSTable(files []*SSTableMeta, targetLevel int) *SSTableMeta

cacheBloom 将过滤器写入缓存(应在 file.Sync() 成功后调用)。

func (*SSTable) Metas

func (ss *SSTable) Metas() []*SSTableMeta

WriteToSSTable 将有序 entries 写入 SSTable 文件(含块索引)

func (*SSTable) ReadAllFromSSTable

func (ss *SSTable) ReadAllFromSSTable(filepath string) ([]*LogEntry, error)

func (*SSTable) ReadFromSSTable

func (ss *SSTable) ReadFromSSTable(filepath string, key []byte) ([]byte, bool)

func (*SSTable) RemoveMeta

func (ss *SSTable) RemoveMeta(target *SSTableMeta)

func (*SSTable) WriteToSSTable

func (ss *SSTable) WriteToSSTable(entries []LogEntry) error

type SSTableMeta

type SSTableMeta struct {
	Level    int
	Filepath string
	MinKey   []byte
	MaxKey   []byte
	Size     int64

	MaxKeyKnown bool
}

SSTableMeta 是一个 SSTable 文件的内存元信息。

MaxKeyKnown 表示 MaxKey 是否可信:仅当它取自文件尾部的块索引(新格式)或写入时直接 填入才为 true。不可信时读路径不施加上界过滤——猜一个 MaxKey 比不过滤危险得多:猜低了 会把命中 key 整段跳过,表现为数据「消失」,而不过滤只是多扫一个文件。

type SkipList

type SkipList struct {
	// contains filtered or unexported fields
}

SkipList 是一张有序跳表,Engine 用它承载内存中的活跃表与待 flush 的不可变表。 它自身不加锁——并发由持有者(Engine)统一同步。

type SkipNode

type SkipNode struct {
	Next  []*SkipNode
	Key   []byte
	Value []byte
}

SkipNode 是跳表节点。Next 的长度即该节点的层高。

type WAL

type WAL struct {
	// contains filtered or unexported fields
}

WAL 存储层预写日志:standalone 模式下,写先 append + fsync 到此处再进 memtable, 提供单机崩溃恢复。记录格式 [op u8][klen u32][vlen u32][key][value](BigEndian)。 与 Raft/raft_wal.go 一致:重放读到残缺尾部记录时直接停止(撕裂的尾写按 EOF 处理), 不使用 CRC。重放是幂等盲写(Put/Delete),未截断的 WAL 反复重放也安全。

写入走 group commit:所有 Append 把请求投递到 reqCh,由唯一的 flushLoop 攒批—— 把当前排队的并发写一次性写入后只 fsync 一次,再唤醒整批等待者。这既把 N 次 fsync 摊销为 1 次(并发越高摊销越充分),又使 flushLoop 成为文件的唯一写者,从而不存在 并发 Write+Sync。持久化契约:Append 返回即代表该记录已 fsync 落盘。

func NewWAL

func NewWAL(path string) (*WAL, error)

NewWAL 打开(或创建)WAL 文件,以追加模式准备写入,并启动 group commit flushLoop。

func (*WAL) Append

func (w *WAL) Append(op uint8, key, value []byte) error

Append 追加一条记录,返回时该记录已随所在批次一起 fsync 落盘。 删除墓碑用 op=WALOpDelete、value 传 nil。

func (*WAL) Close

func (w *WAL) Close() error

Close 停止 flushLoop(其会先排空已投递的请求)并关闭底层文件。

func (*WAL) Replay

func (w *WAL) Replay(fn func(op uint8, key, value []byte) error) error

Replay 从头读取全部记录,对每条调用 fn。读到残缺记录(撕裂尾写)即停止重放, 返回 nil;底层 IO 错误或 fn 返回错误则向上抛出。

func (*WAL) Rewrite

func (w *WAL) Rewrite(records []WALRecord) error

Rewrite 把 WAL 内容整体替换为 records(截断后重写并 fsync),用于 checkpoint 自清洁。 调用方必须保证调用期间没有并发 Append(否则会丢写):本项目由 KVServer.cpMu 独占锁 让所有写静默后再调用。重写经 flushLoop 执行以维持"文件唯一写者"不变式。

type WALRecord

type WALRecord struct {
	Op    uint8
	Key   []byte
	Value []byte
}

WALRecord 一条完整 WAL 记录,供 Rewrite 批量重建 WAL 内容时使用。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL