Documentation
¶
Overview ¶
本文件是 SSTable 的类型与共用定义:磁盘布局常量、块索引、以及 SSTable 本身。 读路径见 sstable_read.go,写路径与 compaction 见 sstable_write.go,元信息管理见 sstable_meta.go。
本文件是 SSTable 的元信息管理:启动时扫描目录重建元信息,以及增删查与不可变快照发布。
本文件是 SSTable 的读路径:常驻句柄、块索引与布隆缓存、点查与全量读。
本文件是 SSTable 的写路径:flush 落盘、写尾(块索引 + 布隆 + Footer)与 compaction 合并。
Index ¶
- Constants
- Variables
- func ResetBlockCacheStats()
- func ResetCompactionStats()
- type BlockCacheStats
- type BlockIndexEntry
- type BloomFilter
- type CompactionStats
- type Engine
- func (m *Engine) Close() error
- func (m *Engine) CompactSSTable(startLevel int)
- func (m *Engine) Delete(key []byte) error
- func (m *Engine) Flush()
- func (m *Engine) FlushToSSTable(entries []LogEntry) error
- func (m *Engine) FlushWorker()
- func (m *Engine) Get(key []byte) ([]byte, error)
- func (m *Engine) InflightBytes() int64
- func (m *Engine) ListenCompactCh()
- func (m *Engine) Put(key []byte, value []byte) error
- func (m *Engine) ReclaimUpTo(bound []byte) int
- func (m *Engine) ScanRange(start, end []byte, fn func(key, value []byte) bool)
- func (m *Engine) Size() int
- func (m *Engine) SnapshotLive() []LogEntry
- func (m *Engine) StartFlush()
- type LogEntry
- type Options
- type PartitionedBloom
- type SSTable
- func (ss *SSTable) AddMeta(meta *SSTableMeta)
- func (ss *SSTable) DeleteSSTable(meta *SSTableMeta)
- func (ss *SSTable) LevelFiles(level int) []*SSTableMeta
- func (ss *SSTable) LoadSSTableMetaList()
- func (ss *SSTable) MergeSSTable(files []*SSTableMeta, targetLevel int) *SSTableMeta
- func (ss *SSTable) Metas() []*SSTableMeta
- func (ss *SSTable) ReadAllFromSSTable(filepath string) ([]*LogEntry, error)
- func (ss *SSTable) ReadFromSSTable(filepath string, key []byte) ([]byte, bool)
- func (ss *SSTable) RemoveMeta(target *SSTableMeta)
- func (ss *SSTable) WriteToSSTable(entries []LogEntry) error
- type SSTableMeta
- type SkipList
- type SkipNode
- type WAL
- type WALRecord
Constants ¶
const ( WALOpPut uint8 = 1 WALOpDelete uint8 = 2 )
WAL 操作码
const DefaultNamespaceSep = ':'
DefaultNamespaceSep 是默认的命名空间分隔符。 数据仓库场景下 key 常带"仓库"前缀,如 "log:2026-05-31"、"order:10086"。
const (
SSTableBlockSize = 64
)
Variables ¶
var ( // ErrKeyNotFound 表示 key 不存在,或其最新版本是一个删除墓碑。 // 它是正常的查询结果而非故障,调用方通常应据此返回「无此键」而不是「内部错误」。 ErrKeyNotFound = errors.New("storage: key not found") // 与 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 ¶
BlockCacheStats 是某一刻的数据块缓存计数。
func ReadBlockCacheStats ¶
func ReadBlockCacheStats() BlockCacheStats
ReadBlockCacheStats 读取当前数据块缓存计数快照。
type BlockIndexEntry ¶
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) 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 ¶
SkipNode 跳表节点 NewEngine 按 opts 构造存储引擎。参数在此固定,此后不再读全局配置。 需要沿用进程配置时传 DefaultOptions()。
func (*Engine) CompactSSTable ¶
func (*Engine) Delete ¶
insert 在跳表中插入键值对(无锁版本,由调用者保证线程安全) insert 插入或覆盖 key,并返回本次操作使 byteSize 变化的增量(覆盖写可能为负)。 调用方用该增量做背压信用对账。
func (*Engine) Flush ¶
func (m *Engine) Flush()
Flush 将 dirty 表数据刷入 SSTable 流程:
- 持锁交换 active → dirty(active 变为 dirty 的不可变快照)
- 创建新的空 active 表用于接受后续写入
- 释放锁,在锁外将 dirty 数据写入 SSTable
- flush 完成后将 dirty 置 nil
func (*Engine) FlushToSSTable ¶
FlushToSSTable 将 entries 写入临时跳表并立即 Flush 到 SSTable 不经过 active 表,不影响正常读写,专用于快照重放等场景
func (*Engine) FlushWorker ¶
func (m *Engine) FlushWorker()
func (*Engine) InflightBytes ¶
InflightBytes 返回当前未 flush(active + 正在刷的 dirty)占用的字节信用,供观测/压测使用。
func (*Engine) ListenCompactCh ¶
func (m *Engine) ListenCompactCh()
func (*Engine) ReclaimUpTo ¶
ReclaimUpTo 丢弃 MaxKey 严格小于 bound 的 SSTable 整个文件,返回回收的文件数。
用于「投递前置缓冲」的保留期回收:投递按 key 升序推进游标,故 bound 之前的数据已全部 被投递读过,整份文件不再需要留在本地。
之所以按整文件丢弃、而不是逐 key 写墓碑:墓碑会让写入量翻倍,且自身还要再经一轮 compaction 才消失;而文件级丢弃是 O(1),无写放大。
保守之处(宁可少回收,不可误删):
- 仅在 MaxKey 可信时回收。没有可读 footer 的文件(老格式或尾部残缺)MaxKey 未知,一律跳过。
- 用严格小于:恰好含 bound 的文件保留,因为 bound 本身尚未被投递消费。
- bound 为空表示尚无已提交游标,不回收任何文件。
与 compaction 互斥(共用 fileMu),避免删掉 compaction 正在读的源文件。
func (*Engine) ScanRange ¶
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) SnapshotLive ¶
SnapshotLive 返回 active+dirty 合并后的全部键值快照(active 覆盖 dirty), 与 ScanRange 不同:保留墓碑(value==nil)。用于 WAL checkpoint 重写——未 flush 的 热数据(含删除墓碑)必须完整保留,否则重放时被删的 key 会从 SSTable 复活。 拷贝底层字节,返回后可安全持有。
func (*Engine) StartFlush ¶
func (m *Engine) StartFlush()
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 (*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 (*SSTable) ReadFromSSTable ¶
func (*SSTable) RemoveMeta ¶
func (ss *SSTable) RemoveMeta(target *SSTableMeta)
func (*SSTable) WriteToSSTable ¶
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 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 落盘。