2026/9/10 2:27:09

Milvus 源码解读:pkoracle 包如何用 Bloom Filter 集合预测 Segment 中主键的存在性

Milvus 源码解读:pkoracle 包如何用 Bloom Filter 集合预测 Segment 中主键的存在性 Milvus 源码解读pkoracle 包如何用 Bloom Filter 集合预测 Segment 中主键的存在性【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus导读pkoraclePrimary Key Oracle主键预言机是 Milvus 数据写入链路中一个小而关键的基础包它为 flushcommon 的 metacache 定义了“段内主键是否存在”的判定接口与两类实现。本文以 pkoracle/README.md 为主线结合BloomFilterSet、LazyPkStats的源码与测试讲清 Bloom Filter 集合在段Segment主键预测中的底层原理、异步懒加载设计以及相关配置参数的完整含义。读完你将理解 Milvus 如何在不精确匹配的前提下用极低的内存开销快速定位哪些段可能包含某主键并知道如何在你的代码中正确选择NewBloomFilterSet与NewBloomFilterSetWithBatchSize。一、pkoracle 包定位metacache 的主键存在性判定接口在 Milvus 的写入链路中internal/flushcommon/metacache负责缓存 DataNode 各 Segment 的元信息状态、级别、分区、以及 Bloom Filter 集合等是写入路径判断某条主键数据该去哪个段的前置关卡。而 pkoracle 正是 metacache 中与主键统计相关的核心抽象层。原文档的定义非常凝练This package defines the interface and implementations for segments bloom filter sets of flushcommon metacache.也就是说pkoracle 解决的是一个存在性预测问题给定一个主键Primary Key简称 PK判断某个 Segment 中是否可能存在该 PK。由于 Bloom Filter 的特性这里的回答天然是可能存在hit或必然不存在miss不可能出现确定存在因此接口命名PkExists和注释里大量使用 predict、may exist 这类谨慎措辞。pkoracle 目录下共有四个 Go 源文件结构非常精简internal/flushcommon/metacache/pkoracle/ ├── README.md # 包级说明本文依据 ├── pk_stats.go # PkStat 接口定义 ├── bloom_filter_set.go # BloomFilterSet 实现基本实现 ├── lazy_pk_stats.go # LazyPkStats 实现懒加载包装 └── bloom_filter_set_test.go # 单元测试二、PkStat 接口存在性判定能力的统一契约pk_stats.go 定义了整个包的基石——PkStat接口它由 7 个方法组成type PkStat interface { PkExists(lc *storage.LocationsCache) bool BatchPkExist(lc *storage.BatchLocationsCache) []bool BatchPkExistWithHits(lc *storage.BatchLocationsCache, hits []bool) []bool UpdatePKRange(ids storage.FieldData) error Roll(newStats ...*storage.PrimaryKeyStats) GetHistory() []*storage.PkStatistics }方法职责可以划分为三组方法职责PkExists单个主键存在性查询返回 boolBatchPkExist/BatchPkExistWithHits批量查询前者返回新分配的[]bool后者把命中结果累积写进调用方传入的hits切片可用于多段合并命中UpdatePKRange/Roll/GetHistory写入与滚动维护向当前 Bloom Filter 追加 PK、把当前状态滚入历史、读取历史统计接口层面已经能看出设计意图查询走 Bloom Filter 位运算写入/滚动负责维护当前段与历史段两层统计。两个实现BloomFilterSet与LazyPkStats都在文件头部通过var _ PkStat (*Xxx)(nil)编译期断言强制实现了该接口防止接口漂移。三、BloomFilterSet基于段 statslog 的基本实现3.1 数据结构与构造方式bloom_filter_set.go 中的BloomFilterSet用sync.RWMutex保护三个核心字段type BloomFilterSet struct { mut sync.RWMutex batchSize uint // 新 Bloom Filter 的初始化容量 current *storage.PkStatistics // 当前进行中的统计写路径追加中 history []*storage.PkStatistics // 历史统计集合已滚动的批次 }current正在被新写入 PK 更新的统计对象服务于growing增量段的实时判定history已经滚动Roll完成的统计对象列表服务于flushed已落盘段。两种构造函数对应两种使用场景这一点务必区分是原文档反复强调的重点// 仅供已 flush 的段使用若用于 growing 段应使用 NewBloomFilterSetWithBatchSize func NewBloomFilterSet(historyEntries ...*storage.PkStatistics) *BloomFilterSet { return BloomFilterSet{ batchSize: paramtable.Get().CommonCfg.BloomFilterSize.GetAsUint(), history: historyEntries, } } // batchSize 用于初始化新的 bloom filter // 它应当是段每次同步sync对应的预估行数estimated row count per batch func NewBloomFilterSetWithBatchSize(batchSize uint, historyEntries ...*storage.PkStatistics) *BloomFilterSet { return BloomFilterSet{ batchSize: batchSize, history: historyEntries, } }关键差异在于batchSize的来源NewBloomFilterSet直接读取全局参数common.bloomFilterSize默认 100000而NewBloomFilterSetWithBatchSize由调用方显式传入。为什么 growing 段需要显式 batchSize因为 Bloom Filter 的大小需要与每个同步批次的行数匹配才能把误报率控制在可接受范围——用 10 万容量去承载数百万行的段误报率会急剧上升。3.2 查询路径current history 双栈命中单键查询 PkExists 的逻辑是先查current再遍历history任何一个统计对象命中即返回truefunc (bfs *BloomFilterSet) PkExists(lc *storage.LocationsCache) bool { bfs.mut.RLock() defer bfs.mut.RUnlock() if bfs.current ! nil bfs.current.TestLocationCache(lc) { return true } for _, bf : range bfs.history { if bf.TestLocationCache(lc) { return true } } return false }批量查询 BatchPkExist 则把命中结果累积进同一个hits切片等价于当前统计 全部历史统计的逐层叠加func (bfs *BloomFilterSet) BatchPkExist(lc *storage.BatchLocationsCache) []bool { hits : make([]bool, lc.Size()) if bfs.current ! nil { bfs.current.BatchPkExist(lc, hits) } for _, bf : range bfs.history { bf.BatchPkExist(lc, hits) } return hits }BatchPkExistWithHits与之一致只是复用调用方传入的hits切片避免高频路径上的重复分配——这也是 go 生态中常见的调用方预分配、实现方累积优化模式。3.3 写入路径UpdatePKRange 与 Roll 的滚动机制写入与滚动是保证当前/历史两层结构正确运转的关键UpdatePKRange若current为空则依据batchSize、全局误报率上限common.maxBloomFalsePositive与过滤器类型common.bloomFilterType创建一个新的storage.PkStatistics然后把传入的 FieldData 主键追加进去底层会同时维护 PK 的 min/max 区间和 Bloom Filter 位图。Roll当有新的PrimaryKeyStats传入时将其转换为PkStatistics携带 PkFilter、MaxPK、MinPK追加到history并清空current——即当前批次结束滚入历史。若调用Roll()时不传任何参数则仅保留现有历史、不做任何变更。这样一个写入→滚动→再写入的循环恰好对应段从 growing 到 flushed 生命周期中主键统计的累积方式每个滚动的批次成为历史中的一个PkStatistics查询时逐一判命。3.4 测试验证写入后必命中、滚动后进入历史bloom_filter_set_test.go 用 testify suite 覆盖了三条关键路径可直接作为使用范例TestWriteReadL60-L78先断言写入前 PK 不存在调用UpdatePKRange写入{1,2,3,4,5}后逐个PkExists与批量BatchPkExist均应命中——验证写入→查询闭环。TestBatchPkExistL80-L109用 10 万容量构造NewBloomFilterSetWithBatchSize(uint(capacity))写入后以 1000 为批次循环批量查询同时验证BatchPkExist与BatchPkExistWithHits两种批量接口——这是 growing 段高频查询的典型形态。TestRollL111-L130新建集合GetHistory()为空 → 写入数据后Roll(newEntry)history长度变为 1 → 再调用空参Roll()历史长度保持不变——验证滚动语义。四、LazyPkStats异步加载 PkStats 的懒加载包装4.1 设计动机与结构LazyPkStats 解决的是统计尚未就绪时的查询语义问题。当段的 PkStats 需要从外部异步加载例如从对象存储读取 statslog时查询方不能阻塞等待因此需要一个占位包装器type LazyPkStats struct { inner atomic.Pointer[PkStat] } func NewLazyPkstats() *LazyPkStats { return LazyPkStats{} } func (s *LazyPkStats) SetPkStats(pk PkStat) { if pk ! nil { s.inner.Store(pk) } }inner使用go.uber.org/atomic的atomic.Pointer[PkStat]SetPkStats可在任意时刻通常由异步协程注入真实统计读侧无锁即可看到最新值。4.2 未就绪时的 fail-open 语义关键设计原文档特别强调LazyPkStats不能用于 growing 段。原因在查询实现中一目了然func (s *LazyPkStats) PkExists(lc *storage.LocationsCache) bool { inner : s.inner.Load() if inner nil { return true // 统计未就绪 → 保守返回存在宁可误报不可漏报 } return (*inner).PkExists(lc) }BatchPkExist/BatchPkExistWithHits同理当inner为 nil 时返回全部为 true的结果切片lo.RepeatBy(lc.Size(), ...)。这是典型的fail-open开放失败策略PK 判定结果只用于缩小候选段范围漏报说不存在但实际存在会导致数据被错误跳过、造成不可恢复的丢失而误报最多带来一次多余的下沉检查因此宁可全部命中。4.3 写操作被显式禁止与 fail-open 读语义相对LazyPkStats对写/维护操作一律拒收UpdatePKRange返回merr.WrapErrServiceInternal(UpdatePKRange shall never be called on LazyPkStats)Roll返回同样的内部错误GetHistory直接返回 nil并注释 shall never be called。这三个方法通过运行时报错而非编译期缺失来强制调用方遵守懒加载包装器不可写入、不可滚动的约束一旦误用会立即暴露在日志中。五、在 metacache 中的实际调用PredictSegments 候选段预测pkoracle 接口在 metacache 中的消费点在 meta_cache.go 的PredictSegmentsfunc (c *metaCacheImpl) PredictSegments(pk storage.PrimaryKey, filters ...SegmentFilter) ([]int64, bool) { var predicts []int64 lc : storage.NewLocationsCache(pk) segments : c.GetSegmentsBy(filters...) for _, segment : range segments { if segment.GetBloomFilterSet().PkExists(lc) { predicts append(predicts, segment.segmentID) } } return predicts, len(predicts) 0 }流程非常清晰为单个 PK 构造LocationsCache缓存该 PK 的 Bloom Filter 哈希位置避免重复计算→ 遍历满足过滤条件的 Segment → 命中PkExists为 true的 Segment ID 进入候选列表。这里的命中结果会被写入路径用于决定数据下沉flush/sync的目标段正是预测一词的落地场景。LocationsCache与BatchLocationsCache定义在 internal/storage/pk_statistics.go其核心优化是哈希位置缓存对BasicBF基本 Bloom Filterk个哈希位置在首次计算后缓存不同段可复用对BlockedBF分块 Bloom Filter只需缓存 1 个哈希结果k值变化无需重算注释明确提醒该 helper 非并发安全应在同一 goroutine 内使用。PkStatistics.BatchPkExistpk_statistics.go的判定顺序也值得一提先做位图运算便宜再做 PK min/max 区间过滤贵一点二者取且从而在误报率与 CPU 开销之间取得平衡。六、可调参数从 paramtable 到 milvus.yamlBloomFilterSet的容量、类型与误报率均来自全局参数表 component_param.go并在 configs/milvus.yaml 中暴露为可配置项配置键milvus.yaml默认值引入版本说明common.bloomFilterSize1000002.3.2Bloom Filter 初始大小即NewBloomFilterSet的默认 batchSizecommon.bloomFilterTypeBlockedBloomFilter2.4.3过滤器类型支持BasicBloomFilter与BlockedBloomFiltercommon.maxBloomFalsePositive0.0012.3.2Bloom Filter 允许的最大误报率决定哈希函数个数与位图尺寸common.bloomFilterApplyBatchSize100002.4.5批量向 Bloom Filter 追加 PK 时的批大小common.bloomFilterApplyParallelFactor2—追加 PK 时的并行因子默认 2×CPU 核数其中前三个参数直接参与UpdatePKRange创建PkStatistics的过程见 bloom_filter_set.gobfs.current storage.PkStatistics{ PkFilter: bloomfilter.NewBloomFilterWithType(bfs.batchSize, paramtable.Get().CommonCfg.MaxBloomFalsePositive.GetAsFloat(), paramtable.Get().CommonCfg.BloomFilterType.GetValue()), }三项参数均可热刷新refreshable:true。bloomFilterSize与maxBloomFalsePositive之间是经典的 Bloom Filter 权衡容量越大、误报率上限越宽松位图越大、内存开销越高误报率越低——需要结合段的大小与可用内存酌情调整。七、横向参照querynodev2 中的姊妹实现值得注意的是Milvus 在查询节点QueryNode的 delegator 侧还有一份同名的独立实现internal/querynodev2/pkoracle/bloom_filter_set.go。两者解决的问题相同用段级 Bloom Filter 预测 PK 是否可能存在但面向不同生命周期查询侧BloomFilterSet实现Candidate接口MayPkExist按segmentID / partitionID / segType追踪用于查询路由时过滤不可能包含目标 PK 的 Segment写入侧本包的BloomFilterSet实现PkStat接口聚焦写入时数据下沉的目标段预测。两份实现共享internal/storage.PkStatistics与LocationsCache作为底层能力从侧面印证了段级主键 Bloom Filter 预测是 Milvus 写入与查询两条链路共同依赖的基础设施。八、小结何时用哪个实现回到原文档的三句话给出工程决策清单场景选择原因已 flush 的段NewBloomFilterSet(...)batchSize 取自全局common.bloomFilterSize直接由历史统计构建growing增量段NewBloomFilterSetWithBatchSize(batchSize, ...)batchSize 需与每批次同步行数匹配控制误报率PkStats 异步加载中、查询不能阻塞LazyPkStats 就绪后SetPkStats(...)未就绪时 fail-open 返回全命中宁误报不漏报严禁用于 growing 段且不可调用UpdatePKRange/Roll/GetHistorypkoracle 用不到 200 行 Go 代码把主键存在性预测这一写入路径的高频操作抽象成接口 双实现 无锁原子懒加载的清晰形态。理解它你就掌握了 Milvus 写入链路中 Segment 候选集裁剪的第一层筛子。【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考