Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,8 @@ BanDB 坐在数据仓库(ClickHouse、Doris 等)的写入入口之前,把
- **高并发写入吸收**:突发高频写入平稳落地,内存占用有界、不会被写入打爆。
- **落盘前数据清洗**:写入进系统的一刻即校验、脱敏、丢弃畸形帧,脏数据不进缓冲。
- **崩溃恢复与断点重续**:进程崩溃重启自动恢复数据;投递从上次已提交位点续传,已投数据不重投。
- **缓冲按已投递位点回收**:开启保留期后,整份已投递完的数据文件被丢弃,本地缓冲不再只增不减
(`RetentionEnabled`,默认关闭——开启即意味着已投递的数据不再能从本地读回)。
- **高并发限流**:过载时自适应限流、主动拒绝多余请求,保护系统不被压垮。
- **可靠投递下游**:按位点批量投递、失败自动重试;下游故障时熔断隔离、恢复后自动探测放行,至少一次送达。
- **横向分片扩展**:数据量增大时按分片扩展到多节点,多副本容错;读请求自动在副本间择优、分摊负载。
Expand Down
8 changes: 7 additions & 1 deletion config/global.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 // 一致性哈希每节点的虚拟节点数
Expand Down Expand Up @@ -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, // 一致性哈希默认虚拟节点数
Expand Down
17 changes: 15 additions & 2 deletions service/delivery/source.go
Original file line number Diff line number Diff line change
Expand Up @@ -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:<ts>)这类天然单调的摄入场景。
//
// 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 升序,
Expand All @@ -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
Expand Down
2 changes: 1 addition & 1 deletion service/delivery/source_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
29 changes: 28 additions & 1 deletion service/delivery_bootstrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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
}
37 changes: 33 additions & 4 deletions service/fsm.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package service

import (
"bytes"
"encoding/json"
"log/slog"
"sync"
Expand All @@ -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"
)

Expand Down Expand Up @@ -184,6 +186,25 @@ func (k *KVServer) Checkpoint() {
}
}

// ReclaimDelivered 丢弃已整体投递完的 SSTable,返回回收的文件数。
//
// bound 取自投递已提交的游标,但会先被压到 offset 保留前缀之下——游标本身就以
// `__offset__/<sink>` 为 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 {
Expand Down Expand Up @@ -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) {
Expand All @@ -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
Expand Down
4 changes: 2 additions & 2 deletions service/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 基础路由处理器
Expand Down Expand Up @@ -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))
}

Expand Down
4 changes: 2 additions & 2 deletions service/scan_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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}`,
Expand Down Expand Up @@ -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))
}
Expand Down
4 changes: 3 additions & 1 deletion service/shard_routing_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
55 changes: 55 additions & 0 deletions storage/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,11 @@ type Engine struct {

// opts 是构造时传入的参数,引擎此后只认它,不再读全局配置。
opts Options

// fileMu 串行化一切「删除 SSTable 文件」的动作:compaction 与保留期回收。
// 二者若交错,compaction 正在读的源文件可能被回收删除——POSIX 下已打开的 fd 仍可读,
// 于是那批已回收的数据会被写进合并输出文件,即「已回收的数据复活」。
fileMu sync.Mutex
}

// SkipNode 跳表节点
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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++ {
Expand All @@ -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
}
Loading
Loading