WeKnora 向量库可插拔抽象:四层设计与跨后端评分归一化
第一篇讲架构全景时提到,WeKnora 横向用「接口 + 注册表」让向量库、模型、存储都可插拔;第三篇讲检索管线时又带过混合召回。这篇番外把镜头对准「可插拔」这三个字本身:它到底怎么做到换一个向量库像改一行配置一样简单。
先说清楚这件事的分量。一个生产级 RAG 系统,向量库选型经常会变:起步用 Postgres 加 pgvector 图省事,数据量上来了想换 Milvus 或 Qdrant,要做全文检索想上 Elasticsearch,部署在腾讯云又想用腾讯云向量库。如果每换一次都要把索引、检索、删除的代码重写一遍,迁移成本会高到没人敢动。WeKnora 用一套分层抽象,把这件事变成了「加一个后端实现 + 改一个环境变量」。
四层抽象:从「裸存取」到「可插拔引擎」
WeKnora 的向量检索不是一个接口打天下,而是四层叠起来的。先看全景:
flowchart TD
R["RetrieveEngineRegistry 注册表层<br/>byEngineType 加 byStoreID,支持按需重建"]
C["CompositeRetrieveEngine 组合层<br/>读按类型路由,写广播全后端"]
H["KeywordsVectorHybridRetrieveEngineService 通用服务层<br/>embedding 加批量加退避"]
P["postgres / pgvector"]
M["milvus"]
Q["qdrant"]
E["elasticsearch v7/v8"]
O["weaviate / doris / 腾讯云 / sqlite ..."]
R -. 按需重建 .-> C
C --> H
H --> P
H --> M
H --> Q
H --> E
H --> O
从下往上,每一层的职责是:
- Repository(后端仓储层):每个后端各自实现,只管把向量和关键词「存进去、查出来」,贴着具体存储。
- Hybrid Service(通用服务层):一个通用装饰器,包住任意 Repository,把 embedding 生成、批量、退避这些「每个后端都一样」的逻辑抽出来只写一次。
- Composite(组合层):把多个 Service 组合成一个统一引擎,读请求按检索类型路由,写请求广播给所有后端。
- Registry(注册表层):管理所有引擎实例,支持按 storeID 懒加载重建。
下面逐层拆开看。
第一层 Repository:后端只实现「裸存取」
接口 RetrieveEngineRepository 定义在 internal/types/interfaces/retriever.go,它要求后端实现的方法都是贴着存储的:Save、BatchSave、Retrieve、DeleteByXXX、Support、EngineType。以 pgvector 为例(internal/application/repository/retriever/postgres/repository.go):
func (r *pgRepository) EngineType() types.RetrieverEngineType {
return types.PostgresRetrieverEngineType
}
func (r *pgRepository) Support() []types.RetrieverType {
// pgvector 同时支持向量检索和关键词检索
return []types.RetrieverType{types.VectorRetrieverType, types.KeywordsRetrieverType}
}
func (r *pgRepository) Retrieve(ctx context.Context, params types.RetrieveParams) ([]*types.RetrieveResult, error) {
// 直接写 SQL,其中 <=> 是 pgvector 的余弦距离算子
// SELECT ..., (1 - embedding <=> ?) AS score FROM ... ORDER BY embedding <=> ? LIMIT ?
...
}
关键在于:这一层只关心「我这个存储引擎怎么存、怎么查」,完全不碰 embedding 怎么算、批量怎么切、并发怎么控。Milvus 的实现用的是 milvus-sdk 的 gRPC 客户端,Elasticsearch 用的是官方 typed client,但它们对上层暴露的都是同一个 RetrieveEngineRepository 接口。Support() 声明自己支持哪些检索类型(向量 / 关键词),这个信息在第三层 Composite 路由时会用到。
第二层 Hybrid Service:embedding 和批量只写一次
如果让每个后端都自己实现「把文本转成向量、批量插入、失败重试」,那就会有近十份几乎一样的代码。WeKnora 用一个装饰器 KeywordsVectorHybridRetrieveEngineService 把这些通用逻辑收拢到一处(keywords_vector_hybrid_indexer.go):
func NewKVHybridRetrieveEngine(
indexRepository interfaces.RetrieveEngineRepository,
engineType types.RetrieverEngineType,
) interfaces.RetrieveEngineService {
return &KeywordsVectorHybridRetrieveEngineService{indexRepository, engineType}
}
它包住任意一个 Repository,对外提供更完整的 RetrieveEngineService。所有「每个后端都一样」的重活都在这层:
- embedding 生成:
Index时先判断retrieverTypes含不含 vector,含就调embedder.Embed把内容转向量,再把向量塞进 params 交给底层 Repository 去存。 - 输入净化:
sanitizeForEmbedding先用正则把 base64 内联图片替换成[image](避免把一大坨图片编码喂给 embedding API),再按字符数封顶 20000,超了就截断并告警。 - 批量加退避:
BatchIndex用batchEmbedWithBackoff做指数退避(5 次,从 200ms 翻倍到 3200ms),再把结果按 40 条一批切开,用信号量把并发限制在 5,避免打爆后端。
这一层是整套「可插拔」设计里最省力的一环:新增一个后端,embedding 和批量逻辑一行都不用写,它自动继承这套装饰器。后端作者只需要专注实现第一层那个贴着存储的 Repository。
第三层 Composite:读路由,写广播
一个知识库可能同时绑了多个引擎(比如向量走 Milvus、关键词走 Elasticsearch)。CompositeRetrieveEngine 负责把它们组合成一个统一引擎(composite.go)。有意思的是,它对读和写用了两种不同的 fan-out 形状。
读路径——按类型路由。Retrieve 收到一批 RetrieveParams,每个 param 并发地去找到第一个「支持该 RetrieverType」的引擎来查:
for _, engineInfo := range c.engineInfos {
if slices.Contains(engineInfo.retrieverType, param.RetrieverType) {
result, err := engineInfo.retrieveEngine.Retrieve(ctx, param)
// 加锁收集结果 ...
break // 每种检索类型只由一个引擎处理
}
}
写路径——广播全后端。Index、BatchIndex、DeleteByXXX、CopyIndices 这些操作用 concurrentExecWithError 给每个引擎都开一个 goroutine,把同一份数据写进所有绑定的存储:
func (c *CompositeRetrieveEngine) concurrentExecWithError(ctx context.Context, fn func(context.Context, *engineInfo) error) error {
var wg sync.WaitGroup
errCh := make(chan error, len(c.engineInfos))
for _, engineInfo := range c.engineInfos {
wg.Add(1)
go func() {
defer wg.Done()
if err := fn(ctx, eng); err != nil {
errCh <- err
}
}()
}
wg.Wait()
close(errCh)
// 返回第一个错误(如果有)
...
}
为什么读要路由、写要广播?因为语义不同:查的时候,一种检索类型找一个最合适的引擎就够了;写的时候必须保证每个绑定的存储都拿到这份索引,否则哪天换个引擎去查就会漏。BatchIndex 还会先按 SourceID 去重,避免重复写入。
第四层 Registry:两种注册维度加按需重建
RetrieveEngineRegistry(registry.go)管着所有引擎实例,内部维护两张 map:
byEngineType:由环境变量RETRIEVE_DRIVER注册的「类型级」引擎(向后兼容的老路径)。byStoreID:由VectorStore数据表注册的「实例级」引擎。这个区别很关键——同一种引擎类型可以注册多个实例,比如两个不同的 Elasticsearch 集群,靠 storeID 区分。
最见功力的是 GetOrLoadByStoreID。注册表是每进程独立的:一个实例上注册好的引擎,在另一个实例上是缺失的,直到那个实例重启;启动时创建失败的引擎,也会一直缺着。与其等运维重新发布,不如按需从数据库重建。这段代码把并发安全考虑得非常细:
singleflight.Group:并发的多个请求命中同一个缺失的 store,只触发一次重建,其余共享结果。storeGen代数计数器:重建开始前采样一次代数,只有代数没变才发布结果,避免一个慢重建覆盖掉期间发生的注册或注销(乐观并发)。failedUntil冷却(30 秒):后端一直挂着时,不让每个请求都白等一次构建超时。EngineBuildTimeout(10 秒)加context.WithoutCancel:重建用独立超时,且从发起请求的 ctx 里剥离出来,否则第一个调用者取消会连累所有共享这次重建的调用者。- panic 恢复:第三方客户端构造函数可能 panic,而 singleflight 会在它自己的 goroutine 上重新抛出,绕过了 HTTP 恢复中间件,所以这里必须自己 recover,不然一个坏后端能拖垮整个进程。
还有一处安全细节:重建失败的原始错误不回传给调用者,因为错误里含后端 endpoint 甚至凭据,只记录到结构化日志,对外统一返回净化过的 sentinel。singleflight 的 key 还特意带上 tenantID,防止一个没做归属校验的调用者「搭便车」加入别的租户的重建。
工厂与归属校验:谁能用哪个 store
拿到引擎的入口是两个工厂函数(factory.go):
CreateRetrieveEngineForKB:同步路径,从 ctx 里读租户信息,用于应用服务里约 23 处调用。CreateRetrieveEngineFromPayload:异步任务路径,租户 ID 从反序列化的任务载荷里显式传入(异步 handler 不填 ctx)。
两者最终都走 resolveBoundEngine,它做两件事:归属校验(StoreOwnedBy 确认这个 store 属于当前租户,防跨租户 IDOR)加解析引擎。绑定关系分两种:
if vectorStoreID == nil || *vectorStoreID == "" {
// 未绑定具体 store:回退到租户的 effective engines(由 RETRIEVE_DRIVER 驱动)
return NewCompositeRetrieveEngine(registry, tenantInfo.GetEffectiveEngines())
}
// 绑定了具体 store:校验归属,懒加载引擎,用它支持的全部检索类型
return resolveBoundEngine(ctx, registry, ownership, tenantID, *vectorStoreID)
错误处理也分级得很清楚,用四个 sentinel 区分:ErrVectorStoreNotFound(不存在,异步任务不重试)、ErrVectorStoreUnavailable(暂时不可用,异步任务应重试)、ErrVectorStoreForbidden(跨租户,不重试)、ErrTenantInfoMissing。这个区分不是洁癖——异步 worker 靠它决定「这个任务是丢弃还是重试」。把一次临时的数据库抖动误判成「永久不存在」,会导致任务被静默丢掉。所有这些 sentinel 都故意不含 store UUID,避免被拿去枚举探测。
跨后端评分归一化:把不同的尺子对齐
到这里引擎能插了,但还有个隐蔽的问题:不同后端返回的相似度分数根本不在一个量纲上。同样是余弦相似度——
- Milvus 的 COSINE 直接返回原始余弦,范围是
[-1, 1]。 - Weaviate 返回 certainty,定义为
(2 - distance) / 2,天然落在[0, 1]。 - pgvector 算
(1 - distance),OpenSearch 的 k-NN 插件用 SpaceType 做了平移,Elasticsearch 的 Lucene script_score 不允许负分……
如果不归一化就把它们混进同一个排序列表,分数就没有可比性。EngineAwareNormalizer(normalizer.go)按引擎类型套各自的公式,统一映射到 [0, 1]:
func (EngineAwareNormalizer) Normalize(_ context.Context, score float64,
retrieverType types.RetrieverType, engineType types.RetrieverEngineType) float64 {
if retrieverType != types.VectorRetrieverType {
return score // BM25 等非向量分数不归一化,直接透传
}
switch engineType {
case types.MilvusRetrieverEngineType:
return clamp01((score + 1) / 2) // [-1,1] 映射到 [0,1]
// 其余引擎到达这里时已经在 [0,1],用 clamp01 兜底
default:
return clamp01(score)
}
}
有两个设计判断特别值得学:
- 只归一化向量分数,关键词(BM25)分数透传。因为 BM25 是无上界的正数,强行缩放会压垮长尾;而下游用的是 RRF 融合,基于排名而非绝对分值,天然对量纲免疫。
clamp01顺手处理 NaN 和 Inf。这不是多余的——下游slices.SortFunc要求严格弱序,而 NaN 跟谁比都不满足「大于」或「小于」,漏进去会破坏排序不变式。
归一化之后,向量和关键词两路结果用 RRF(Reciprocal Rank Fusion) 融合(knowledgebase_search_fusion.go):
// RRF 分数 = vectorWeight/(k+vectorRank) + keywordWeight/(k+keywordRank)
rrfK := retrievalCfg.GetEffectiveRRFK() // 默认 60
vectorWeight, keywordWeight := retrievalCfg.GetEffectiveRRFWeights() // 默认 0.7 / 0.3
RRF 只看排名不看原始分,这正是它适合融合异构来源的原因。k=60 是个平滑常数,用来压低头部名次的绝对优势,让两路结果更均衡地贡献分数。
多存储 fan-out:并发查加按需归一化
最后一块拼图,是知识库绑了多个 store 时的并发检索(knowledgebase_search_fanout.go)。retrieveFromStores 用 errgroup 把并发限制在 4,每个 store 组带独立超时:
if len(groups) == 1 {
// 单个 store 走快路径,不起 goroutine
return groups[0].Engine.Retrieve(ctx, paramsWithTopK(groups[0]))
}
g, gctx := errgroup.WithContext(ctx)
g.SetLimit(defaultMultiStoreFanoutLimit) // = 4
// 每组一个 goroutine,带 per-group 超时,结果加锁汇总
这里有个精巧的判断:归一化只在结果跨越多个引擎类型时才做(hasMixedEngineTypes)。同一个引擎产出的原始分数本来就可比,不必归一化;只有要把 Milvus 的分和 Elasticsearch 的分放一起排时,才需要 EngineAwareNormalizer 出手。遇到未知引擎类型,退化成 clamp01 兜底,并对每种未知类型只告警一次(告警前还用 SanitizeForLog 清洗,防止 CR/LF 日志注入)。
任何一个 store 组失败,都不会把原始错误抛给用户,而是收敛成一个统一的「向量检索不可用」错误——内部细节只进日志,响应体不泄露。
小结
WeKnora 这套向量库可插拔抽象,给我的启发是它把「变化」和「不变」切得非常干净:
- 不变的部分(embedding、批量、退避、并发控制、评分归一化、RRF 融合)收拢到 Hybrid Service 和归一化器里,只写一次,所有后端共享。
- 变化的部分(每个存储引擎怎么存、怎么查)隔离在薄薄的 Repository 层,加一个后端就实现一个 Repository,其余全自动继承。
- 运行时的复杂性(多实例、按需重建、并发安全、跨租户隔离、错误分级)沉淀在 Registry 和 Factory 里,用 singleflight、代数计数器、冷却、sentinel 这些手段逐一化解。
回头看第一篇说的「横向用接口加注册表让向量库可插拔」,这句话背后是四层抽象、两种 fan-out、一套跨后端评分对齐机制在撑着。所谓「可插拔」,从来不是定义一个接口那么简单,而是把每一处会因后端不同而不同的东西,都找到合适的层去安放。
系列导航:
- (一)架构全景:Go 与 Python 双引擎
- (二)自适应分片机制:三层策略与父子分块
- (三)RAG 检索管线:12 阶段 Pipeline 剖析
- (四)ReAct Agent 与三大能力
- (五·番外)向量库可插拔抽象 ← 本文
