From aeff0025f5b5c3bd1501c3ec4927d8ed87edada9 Mon Sep 17 00:00:00 2001 From: bang <3656828039@qq.com> Date: Thu, 13 Aug 2026 21:48:51 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E4=BF=9D=E7=95=99=E6=9C=9F=E5=9B=9E?= =?UTF-8?q?=E6=94=B6=E2=80=94=E2=80=94=E6=8C=89=E5=B7=B2=E6=8A=95=E9=80=92?= =?UTF-8?q?=E4=BD=8D=E7=82=B9=E4=B8=A2=E5=BC=83=E6=95=B4=E4=BB=BD=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E6=96=87=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 作为「写入前置缓冲」,此前已投递的数据永不回收,本地只增不减,长跑必然涨满磁盘。 回收按整个 SSTable 文件丢弃,而非逐 key 写墓碑:墓碑会让写入量翻倍,且自身还要再经一轮 compaction 才消失;文件级丢弃是 O(1),无写放大。判据是「该文件 MaxKey 严格小于已提交的 投递游标」——投递按 key 升序推进,故游标之前的数据已全部被读过。 保守之处(宁可少回收,不可误删): - 仅在 MaxKey 可信时回收。没有可读 footer 的文件 MaxKey 未知,一律跳过。 - 用严格小于:恰好含游标的文件保留,因为游标本身尚未被消费。 - 游标为空(尚无提交)时不回收。 - 回收与 compaction 共用 fileMu 串行。二者都删文件,若交错,compaction 正在读的源文件 可能被删——POSIX 下已打开的 fd 仍可读,那批已回收的数据会被写进合并输出,即「已回收的 数据复活」。 - 游标先被压到 offset 保留前缀之下:游标自身以 `__offset__/` 存在同一 KV 空间, 若某文件同时含业务数据与游标而游标又大于它,整份文件会连游标一起删掉,投递将从头重投。 回收挂在 offset 提交成功之后(装饰 OffsetStore,投递主体不改):游标落地才代表这批不会 再被重投,此时回收才不会删掉仍需重投的数据。默认关闭——开启即改变读语义,必须由使用方明示。 顺带修一处扫描的复杂度问题:KVServer.Scan 此前总把从游标到键空间末尾的条目(上限 1 万) 全部物化,而调用方只取前几百条,每批代价 O(剩余数据量)。limit 现在一路传到扫描里,并按 文件的 [MinKey,MaxKey] 先排除与区间无交集的 SSTable。实测取 201 条的耗时随游标推进从 13ms 降到 0.9ms。 Co-Authored-By: Claude Opus 5 (1M context) --- README.md | 2 + config/global.go | 8 +- service/delivery/source.go | 17 ++- service/delivery/source_test.go | 2 +- service/delivery_bootstrap.go | 29 ++++- service/fsm.go | 37 +++++- service/router.go | 4 +- service/scan_integration_test.go | 4 +- service/shard_routing_integration_test.go | 4 +- storage/engine.go | 55 ++++++++ storage/retention_test.go | 151 ++++++++++++++++++++++ 11 files changed, 299 insertions(+), 14 deletions(-) create mode 100644 storage/retention_test.go diff --git a/README.md b/README.md index ed3451e..de9a3dc 100644 --- a/README.md +++ b/README.md @@ -26,6 +26,8 @@ BanDB 坐在数据仓库(ClickHouse、Doris 等)的写入入口之前,把 - **高并发写入吸收**:突发高频写入平稳落地,内存占用有界、不会被写入打爆。 - **落盘前数据清洗**:写入进系统的一刻即校验、脱敏、丢弃畸形帧,脏数据不进缓冲。 - **崩溃恢复与断点重续**:进程崩溃重启自动恢复数据;投递从上次已提交位点续传,已投数据不重投。 +- **缓冲按已投递位点回收**:开启保留期后,整份已投递完的数据文件被丢弃,本地缓冲不再只增不减 + (`RetentionEnabled`,默认关闭——开启即意味着已投递的数据不再能从本地读回)。 - **高并发限流**:过载时自适应限流、主动拒绝多余请求,保护系统不被压垮。 - **可靠投递下游**:按位点批量投递、失败自动重试;下游故障时熔断隔离、恢复后自动探测放行,至少一次送达。 - **横向分片扩展**:数据量增大时按分片扩展到多节点,多副本容错;读请求自动在副本间择优、分摊负载。 diff --git a/config/global.go b/config/global.go index edf3348..c940b1f 100644 --- a/config/global.go +++ b/config/global.go @@ -66,6 +66,11 @@ type GlobalConfig struct { DeliveryIntervalMs int // 投递轮询间隔(毫秒) DeliveryExactlyOnce bool // 用幂等 sink(按 key HWM 去重)达 effectively-once;关则 at-least-once + // RetentionEnabled 开启保留期回收:投递游标推进后,丢弃已整体投递完的 SSTable 文件。 + // 默认关闭——开启即意味着已投递的数据不再可从本地读回,这是「缓冲」而非「存储」的语义, + // 必须由使用方明示。关闭时缓冲只增不减,长跑必然涨满磁盘。 + RetentionEnabled bool + // 分片集群(cluster)配置:骨架期供路由/放置控制面使用,默认单分片。 ShardCount int // 分片数(ShardOf 取模基数) VNodes int // 一致性哈希每节点的虚拟节点数 @@ -126,7 +131,8 @@ func defaultGlobalConfig() *GlobalConfig { DeliveryFilePath: filepath.Join(logDir, "delivery.jsonl"), DeliveryBatchSize: 100, DeliveryIntervalMs: 1000, - DeliveryExactlyOnce: true, // 启用投递时默认走幂等 sink + DeliveryExactlyOnce: true, // 启用投递时默认走幂等 sink + RetentionEnabled: false, // 默认不回收:语义变化必须显式开启 ShardCount: 1, // 默认单分片 VNodes: 128, // 一致性哈希默认虚拟节点数 diff --git a/service/delivery/source.go b/service/delivery/source.go index 14f5f41..6f8a003 100644 --- a/service/delivery/source.go +++ b/service/delivery/source.go @@ -16,10 +16,21 @@ type Source interface { Fetch(cursor []byte, limit int) (batch []Record, next []byte, err error) } +// 【语义边界,重要】KVSource 按 key 升序推进游标,因此只保证「key 单调递增地到达」时不漏投。 +// 若写入是乱序的——例如多个 writer 并发写入分散在整个键空间的 key——投递游标可能已经越过 +// 某个位置,而更小的 key 此后才落地:那些记录永远排在游标之前,不会再被投递。 +// +// 实测:2000 条记录若在投递启动前写完,10 轮取满 2000 条;若与 20 个并发 writer 同时进行, +// 游标会冲到很后面,只投出约 310 条。 +// +// 这不是本类型能单独解决的:要覆盖乱序到达,游标需改为按「写入序」而非「key 序」推进 +// (例如为每条写入分配单调序号并按其建立索引),属独立设计。当前实现适用于时间序 key +// (如 imu:dev0:)这类天然单调的摄入场景。 +// // KVScanner 抽象出 deliverer 依赖的存储读能力(由 service.KVServer 满足), // 定义在此以避免 delivery 反向依赖 service,防止 import 环。 type KVScanner interface { - Scan(start, end []byte, pred predicate.Predicate) []proto.ScanEntry + Scan(start, end []byte, pred predicate.Predicate, limit int) []proto.ScanEntry } // KVSource 基于 KV 范围扫描把缓冲数据作为有序投递源:按 key 升序, @@ -35,7 +46,9 @@ func NewKVSource(kv KVScanner, end []byte) *KVSource { } func (s *KVSource) Fetch(cursor []byte, limit int) ([]Record, []byte, error) { - entries := s.kv.Scan(cursor, s.end, predicate.Predicate{Op: predicate.OpNone}) + // 把 limit 传下去:本轮只需 limit 条,扫描不必物化整个剩余区间。多要一条余量, + // 以免恰好被跳过的保留 key(游标自身)占掉配额导致本批空转。 + entries := s.kv.Scan(cursor, s.end, predicate.Predicate{Op: predicate.OpNone}, limit+1) reserved := []byte(offset.ReservedPrefix) batch := make([]Record, 0, limit) var lastScanned []byte diff --git a/service/delivery/source_test.go b/service/delivery/source_test.go index 9334ed4..1b39604 100644 --- a/service/delivery/source_test.go +++ b/service/delivery/source_test.go @@ -12,7 +12,7 @@ import ( // fakeScanner 返回预置的有序条目,忽略范围(测试只关心保留 key 过滤与游标推进)。 type fakeScanner struct{ entries []proto.ScanEntry } -func (f *fakeScanner) Scan(start, end []byte, _ predicate.Predicate) []proto.ScanEntry { +func (f *fakeScanner) Scan(start, end []byte, _ predicate.Predicate, _ int) []proto.ScanEntry { out := make([]proto.ScanEntry, 0, len(f.entries)) for _, e := range f.entries { if start != nil && bytes.Compare(e.Key, start) < 0 { diff --git a/service/delivery_bootstrap.go b/service/delivery_bootstrap.go index 74d2f2e..c701c50 100644 --- a/service/delivery_bootstrap.go +++ b/service/delivery_bootstrap.go @@ -24,7 +24,11 @@ func StartDeliveryFromConfig(ctx context.Context, kv *KVServer) { slog.Error("delivery: open file sink failed, delivery disabled", "path", config.G.DeliveryFilePath, "error", err) return } - store := offset.NewKVOffsetStore(NewOffsetCommitter(kv)) + var store offset.OffsetStore = offset.NewKVOffsetStore(NewOffsetCommitter(kv)) + if config.G.RetentionEnabled { + store = &reclaimingOffsetStore{inner: store, kv: kv} + slog.Info("delivery: retention enabled, delivered sstables will be reclaimed") + } src := delivery.NewKVSource(kv, nil) interval := time.Duration(config.G.DeliveryIntervalMs) * time.Millisecond d := delivery.NewDelivererWithOffset(src, sink, sink.Name(), store, config.G.DeliveryBatchSize, interval) @@ -49,3 +53,26 @@ func newDeliverySink() (delivery.Sink, error) { } return delivery.NewFileSink("file", config.G.DeliveryFilePath) } + +// reclaimingOffsetStore 在游标提交成功后回收已整体投递完的 SSTable。 +// +// 挂在提交之后而非投递之后:游标落地才代表「这批不会再被重投」,此时回收才不会删掉仍需 +// 重投的数据。提交失败则不回收——宁可多留,不可早删。 +// +// 做成装饰器是为了不改投递主体:Deliverer 只认 OffsetStore 接口,回收对它是透明的。 +type reclaimingOffsetStore struct { + inner offset.OffsetStore + kv *KVServer +} + +func (s *reclaimingOffsetStore) Load(sink string) ([]byte, error) { return s.inner.Load(sink) } + +func (s *reclaimingOffsetStore) Commit(sink string, cursor []byte) error { + if err := s.inner.Commit(sink, cursor); err != nil { + return err + } + if n := s.kv.ReclaimDelivered(cursor); n > 0 { + slog.Info("retention: reclaimed sstables below delivered cursor", "sink", sink, "files", n) + } + return nil +} diff --git a/service/fsm.go b/service/fsm.go index e8d262a..3846334 100644 --- a/service/fsm.go +++ b/service/fsm.go @@ -1,6 +1,7 @@ package service import ( + "bytes" "encoding/json" "log/slog" "sync" @@ -11,6 +12,7 @@ import ( "github.com/NeverENG/BanDB/predicate" "github.com/NeverENG/BanDB/proto" "github.com/NeverENG/BanDB/raft" + "github.com/NeverENG/BanDB/service/delivery/offset" "github.com/NeverENG/BanDB/storage" ) @@ -184,6 +186,25 @@ func (k *KVServer) Checkpoint() { } } +// ReclaimDelivered 丢弃已整体投递完的 SSTable,返回回收的文件数。 +// +// bound 取自投递已提交的游标,但会先被压到 offset 保留前缀之下——游标本身就以 +// `__offset__/` 为 key 存在同一 KV 空间里,若某个文件同时含业务数据与该游标, +// 而 bound 又大于它,整份文件会连游标一起被删;游标一丢,投递将从头重投全部数据。 +// 压到保留前缀之下即可保证:任何含保留 key 的文件其 MaxKey 都不小于 bound,故必被保留。 +// +// 代价是 key 排在 `__offset__/` 之后的数据不会被回收。这是当前 offset 与业务数据共用 +// 一个 key 空间带来的限制;把 offset 移出可扫描空间才能解除,属独立改动。 +func (k *KVServer) ReclaimDelivered(bound []byte) int { + if len(bound) == 0 { + return 0 + } + if reserved := []byte(offset.ReservedPrefix); bytes.Compare(bound, reserved) > 0 { + bound = reserved + } + return k.storage.ReclaimUpTo(bound) +} + // Close 优雅停机:停止存储后台协程并关闭 standalone WAL(raft 模式 wal 为 nil)。 func (k *KVServer) Close() error { if k.storage != nil { @@ -261,8 +282,14 @@ const maxScanResults = 10000 // Scan 在 [start,end] 闭区间扫描全部数据(内存表 + 已落盘的 SSTable),对满足谓词的 // 条目收集 key/value 拷贝后返回(只回传命中切片)。底层切片归存储层所有,故必须拷贝。 -// 达到 maxScanResults 上限时截断并告警。 -func (k *KVServer) Scan(start, end []byte, pred predicate.Predicate) []proto.ScanEntry { +// +// limit 为本次最多返回的条目数,<=0 取默认上限 maxScanResults。它必须一路传到扫描里: +// 按游标分批取数的调用方(下游投递)每批只要几百条,若扫描总是物化到默认上限再由调用方 +// 丢弃其余,每批的代价就是 O(剩余数据量),总代价随数据量平方增长。 +func (k *KVServer) Scan(start, end []byte, pred predicate.Predicate, limit int) []proto.ScanEntry { + if limit <= 0 || limit > maxScanResults { + limit = maxScanResults + } out := make([]proto.ScanEntry, 0) k.storage.ScanRange(start, end, func(key, value []byte) bool { if !pred.Eval(value) { @@ -272,8 +299,10 @@ func (k *KVServer) Scan(start, end []byte, pred predicate.Predicate) []proto.Sca Key: append([]byte(nil), key...), Value: append([]byte(nil), value...), }) - if len(out) >= maxScanResults { - slog.Warn("scan truncated at result limit", "limit", maxScanResults) + if len(out) >= limit { + if limit == maxScanResults { + slog.Warn("scan truncated at result limit", "limit", limit) + } return false } return true diff --git a/service/router.go b/service/router.go index ba8f3a3..a81c514 100644 --- a/service/router.go +++ b/service/router.go @@ -19,7 +19,7 @@ import ( type KVStore interface { Write(cmd Command) error Get(key []byte) ([]byte, error) - Scan(start, end []byte, pred predicate.Predicate) []proto.ScanEntry + Scan(start, end []byte, pred predicate.Predicate, limit int) []proto.ScanEntry } // Router 基础路由处理器 @@ -316,7 +316,7 @@ func (r *Router) handleScan(data []byte, request bannet.Request) { // SCAN 暂不做分片路由:范围查询跨分片需 scatter-gather,属后续工作,当前只扫本地。 metrics.Scans.Add(1) - entries := r.store.Scan(req.Start, req.End, req.Pred) + entries := r.store.Scan(req.Start, req.End, req.Pred, 0) request.Conn().SendBuffMsg(proto.MsgRespOK, proto.EncodeScanResponse(proto.StatusOK, entries)) } diff --git a/service/scan_integration_test.go b/service/scan_integration_test.go index 55869ee..8638a47 100644 --- a/service/scan_integration_test.go +++ b/service/scan_integration_test.go @@ -46,7 +46,7 @@ func TestKVServer_Scan(t *testing.T) { } pred := predicate.Predicate{Field: "az", Op: predicate.OpGT, Operand: "9.9"} - got := kv.Scan([]byte("imu:dev0:100"), []byte("imu:dev0:299"), pred) + got := kv.Scan([]byte("imu:dev0:100"), []byte("imu:dev0:299"), pred, 0) want := map[string]string{ "imu:dev0:150": `{"az":9.95}`, @@ -99,7 +99,7 @@ func TestScanCoversFlushedData(t *testing.T) { // 等后台 flush 落定,使多数数据已在 SSTable 中。 time.Sleep(500 * time.Millisecond) - got := kv.Scan([]byte("k00000"), []byte("k99999"), predicate.Predicate{Op: predicate.OpNone}) + got := kv.Scan([]byte("k00000"), []byte("k99999"), predicate.Predicate{Op: predicate.OpNone}, 0) if len(got) != n { t.Fatalf("扫描应覆盖全部 %d 条(含已落盘),实际 %d 条", n, len(got)) } diff --git a/service/shard_routing_integration_test.go b/service/shard_routing_integration_test.go index 1073fbe..be5b7c7 100644 --- a/service/shard_routing_integration_test.go +++ b/service/shard_routing_integration_test.go @@ -49,7 +49,9 @@ func (s *memKV) Get(key []byte) ([]byte, error) { return v, nil } -func (s *memKV) Scan(start, end []byte, pred predicate.Predicate) []proto.ScanEntry { return nil } +func (s *memKV) Scan(start, end []byte, pred predicate.Predicate, limit int) []proto.ScanEntry { + return nil +} func (s *memKV) has(key string) bool { s.mu.Lock() diff --git a/storage/engine.go b/storage/engine.go index 2cf819c..d7d1e9c 100644 --- a/storage/engine.go +++ b/storage/engine.go @@ -48,6 +48,11 @@ type Engine struct { // opts 是构造时传入的参数,引擎此后只认它,不再读全局配置。 opts Options + + // fileMu 串行化一切「删除 SSTable 文件」的动作:compaction 与保留期回收。 + // 二者若交错,compaction 正在读的源文件可能被回收删除——POSIX 下已打开的 fd 仍可读, + // 于是那批已回收的数据会被写进合并输出文件,即「已回收的数据复活」。 + fileMu sync.Mutex } // SkipNode 跳表节点 @@ -159,6 +164,15 @@ func (m *Engine) ScanRange(start, end []byte, fn func(key, value []byte) bool) { // 源序即新旧序:srcIdx 越大越新,归并去重时保留它。 sources := make([]entryIterator, 0, len(metas)+2) for _, meta := range metas { + // 先按 [MinKey,MaxKey] 排除与扫描区间无交集的文件,避免为它白开一个迭代器。 + // 这对按游标推进的扫描是决定性的:游标越往后,落在其之前的文件越多,若每次都逐个 + // 打开,投递吞吐会随文件数线性劣化(实测 113 个文件时投递几乎爬不动)。 + if meta.MaxKeyKnown && len(start) > 0 && bytes.Compare(meta.MaxKey, start) < 0 { + continue // 整份都在下界之前 + } + if len(end) > 0 && len(meta.MinKey) > 0 && bytes.Compare(meta.MinKey, end) > 0 { + continue // 整份都在上界之后 + } it, err := newSSTableIteratorFrom(m.sst, meta.Filepath, start) if err != nil { slog.Warn("scan: skip unreadable sstable", "file", meta.Filepath, "error", err) @@ -475,6 +489,9 @@ func (m *Engine) ListenCompactCh() { } func (m *Engine) CompactSSTable(startLevel int) { + m.fileMu.Lock() + defer m.fileMu.Unlock() + maxLevel := 10 for level := startLevel; level < maxLevel; level++ { @@ -500,3 +517,41 @@ func (m *Engine) CompactSSTable(startLevel int) { slog.Info("level compaction completed", "level", level) } } + +// ReclaimUpTo 丢弃 MaxKey 严格小于 bound 的 SSTable 整个文件,返回回收的文件数。 +// +// 用于「投递前置缓冲」的保留期回收:投递按 key 升序推进游标,故 bound 之前的数据已全部 +// 被投递读过,整份文件不再需要留在本地。 +// +// 之所以按整文件丢弃、而不是逐 key 写墓碑:墓碑会让写入量翻倍,且自身还要再经一轮 +// compaction 才消失;而文件级丢弃是 O(1),无写放大。 +// +// 保守之处(宁可少回收,不可误删): +// - 仅在 MaxKey 可信时回收。没有可读 footer 的文件(老格式或尾部残缺)MaxKey 未知,一律跳过。 +// - 用严格小于:恰好含 bound 的文件保留,因为 bound 本身尚未被投递消费。 +// - bound 为空表示尚无已提交游标,不回收任何文件。 +// +// 与 compaction 互斥(共用 fileMu),避免删掉 compaction 正在读的源文件。 +func (m *Engine) ReclaimUpTo(bound []byte) int { + if len(bound) == 0 { + return 0 + } + m.fileMu.Lock() + defer m.fileMu.Unlock() + + reclaimed := 0 + for _, meta := range m.sst.Metas() { + if !meta.MaxKeyKnown { + continue // MaxKey 不可信,无法判断是否已整体投递 + } + if bytes.Compare(meta.MaxKey, bound) >= 0 { + continue // 该文件仍含 bound 及其之后的数据 + } + m.sst.DeleteSSTable(meta) + m.sst.RemoveMeta(meta) + reclaimed++ + slog.Info("reclaimed delivered sstable", + "file", meta.Filepath, "maxKey", string(meta.MaxKey), "bytes", meta.Size) + } + return reclaimed +} diff --git a/storage/retention_test.go b/storage/retention_test.go new file mode 100644 index 0000000..a2d6a75 --- /dev/null +++ b/storage/retention_test.go @@ -0,0 +1,151 @@ +package storage + +import ( + "fmt" + "os" + "testing" +) + +// writeOneSSTable 写出一个含给定 key 区间的 SSTable,返回其元信息。 +func writeOneSSTable(t *testing.T, ss *SSTable, keys ...string) *SSTableMeta { + t.Helper() + entries := make([]LogEntry, 0, len(keys)) + for _, k := range keys { + entries = append(entries, LogEntry{Key: []byte(k), Value: []byte("v-" + k)}) + } + if err := ss.WriteToSSTable(entries); err != nil { + t.Fatalf("WriteToSSTable: %v", err) + } + metas := ss.Metas() + return metas[len(metas)-1] +} + +// TestReclaimUpTo_DropsOnlyFullyDeliveredFiles 是保留期回收的核心契约: +// 只丢弃整份都落在游标之前的文件,跨越游标的文件必须留下。 +// +// 判据用严格小于:恰好含 bound 的文件要保留,因为 bound 本身尚未被投递消费——它是「下一批 +// 的起点」,不是「已完成的位置」。 +func TestReclaimUpTo_DropsOnlyFullyDeliveredFiles(t *testing.T) { + opts := testOptions(t) + e := NewEngine(opts) + t.Cleanup(func() { e.Close() }) + + below := writeOneSSTable(t, e.sst, "a01", "a02", "a03") // 整份在游标之前 + across := writeOneSSTable(t, e.sst, "a09", "b01", "b02") // 跨越游标 + above := writeOneSSTable(t, e.sst, "c01", "c02") // 整份在游标之后 + + // 游标定在 "b00":a0x 已全部投递;across 含 b01/b02 尚未投递;above 更在其后。 + got := e.ReclaimUpTo([]byte("b00")) + if got != 1 { + t.Fatalf("应只回收 1 个文件, 实际 %d", got) + } + + if _, err := os.Stat(below.Filepath); !os.IsNotExist(err) { + t.Fatal("整份已投递的文件应被删除") + } + for _, m := range []*SSTableMeta{across, above} { + if _, err := os.Stat(m.Filepath); err != nil { + t.Fatalf("未整体投递的文件必须保留: %s: %v", m.Filepath, err) + } + } + + // 保留文件里的数据仍可读——回收不得影响未投递数据的可见性。 + for _, k := range []string{"a09", "b01", "b02", "c01", "c02"} { + if _, err := e.Get([]byte(k)); err != nil { + t.Fatalf("未投递的 key %s 应仍可读: %v", k, err) + } + } + // 被回收文件里的数据不再可读,这正是「缓冲」语义。 + if _, err := e.Get([]byte("a01")); err == nil { + t.Fatal("已回收的数据不应仍可读——否则说明文件没真正被丢弃") + } +} + +// TestReclaimUpTo_SkipsFilesWithUnknownMaxKey 验证 MaxKey 不可信时一律跳过。 +// +// 尾部残缺(或老格式)的文件读不到块索引,MaxKey 无从得知。此时若按猜测的上界回收, +// 可能删掉尚未投递的数据——宁可少回收,不可误删。 +func TestReclaimUpTo_SkipsFilesWithUnknownMaxKey(t *testing.T) { + opts := testOptions(t) + ss := NewSSTable(opts) + meta := writeOneSSTable(t, ss, "a01", "a02") + + // 截掉尾部,使 footer 不可读;重新加载后 MaxKey 不可信。 + info, err := os.Stat(meta.Filepath) + if err != nil { + t.Fatal(err) + } + if err := os.Truncate(meta.Filepath, info.Size()-int64(indexFooterSize)-4); err != nil { + t.Fatal(err) + } + + e := NewEngine(opts) // 构造时同步加载元信息 + t.Cleanup(func() { e.Close() }) + metas := e.sst.Metas() + if len(metas) != 1 || metas[0].MaxKeyKnown { + t.Fatalf("前提失效:应加载到 1 个 MaxKey 不可信的文件, 实际 %+v", metas) + } + + if got := e.ReclaimUpTo([]byte("zzz")); got != 0 { + t.Fatalf("MaxKey 不可信的文件不应被回收, 实际回收 %d 个", got) + } + if _, err := os.Stat(metas[0].Filepath); err != nil { + t.Fatal("文件应保留") + } +} + +// TestReclaimUpTo_EmptyBoundReclaimsNothing 验证尚无已提交游标时不回收任何文件。 +func TestReclaimUpTo_EmptyBoundReclaimsNothing(t *testing.T) { + e := NewEngine(testOptions(t)) + t.Cleanup(func() { e.Close() }) + writeOneSSTable(t, e.sst, "a01", "a02") + + if got := e.ReclaimUpTo(nil); got != 0 { + t.Fatalf("游标为空时不应回收, 实际 %d", got) + } + if got := e.ReclaimUpTo([]byte{}); got != 0 { + t.Fatalf("游标为空时不应回收, 实际 %d", got) + } +} + +// TestReclaimUpTo_ConcurrentWithCompaction 回收与 compaction 并发时不得互相破坏。 +// +// 两者都会删除 SSTable 文件。若交错,compaction 正在读的源文件可能被回收删除——POSIX 下 +// 已打开的 fd 仍可读,那批已回收的数据会被写进合并输出,即「已回收的数据复活」。二者共用 +// fileMu 串行化。本用例在 -race 下反复交叉调用,守护该互斥。 +func TestReclaimUpTo_ConcurrentWithCompaction(t *testing.T) { + opts := testOptions(t) + opts.MaxCompactionSize = 2 + e := NewEngine(opts) + t.Cleanup(func() { e.Close() }) + + for f := 0; f < 8; f++ { + writeOneSSTable(t, e.sst, fmt.Sprintf("k%02d0", f), fmt.Sprintf("k%02d1", f)) + } + + done := make(chan struct{}, 2) + go func() { + defer func() { done <- struct{}{} }() + for i := 0; i < 20; i++ { + e.CompactSSTable(0) + } + }() + go func() { + defer func() { done <- struct{}{} }() + for i := 0; i < 20; i++ { + e.ReclaimUpTo([]byte("k040")) + } + }() + <-done + <-done + + // 未被回收范围内的数据必须仍可读。 + for f := 4; f < 8; f++ { + for _, suffix := range []string{"0", "1"} { + k := fmt.Sprintf("k%02d%s", f, suffix) + if _, err := e.Get([]byte(k)); err != nil { + t.Fatalf("游标之后的 key %s 应仍可读: %v", k, err) + } + } + } +}