Indexer + Retriever 源码:从 Eino 接口到 pgvector SQL(第73篇-E59)
系列「企业级 AI Agent 实现拆解」E59 篇,Part 13 RAG 篇第八章。上一篇 把 pgvector 的 SQL 层摸清了。这篇把它接进 Eino:两个接口、一个类型、约 200 行 Go,跑通写入和检索。顺带把
eino-ext现有几个后端的共同套路提炼出来——所有后端实现都是同一个六步模板。
读完这篇你会知道
indexer.Options/retriever.Options里有什么,谁该提供默认值eino-ext所有后端实现的六步套路(redis / es8 / milvus 全一样)database/sql不认识[]float64,向量参数怎么传ScoreThreshold是相似度不是距离,换算搞错会返回空列表embCtx:为什么 embedding 的打点要单独换RunInfo(实测三层嵌套输出)- 一个类型同时实现
Indexer+Retriever的实际好处- 参数化查询顶不顶得住注入(实测把
DROP TABLE当查询扔进去)
一、两个接口的 Options 各有什么
接口本身各一个方法(第 66 篇《最简 RAG》已经列过),值得看的是 Options:
// indexer.Options
type Options struct {
Index *string // 写到哪个索引/表
SubIndexes []string // 逻辑子分区
Embedding embedding.Embedder // 用哪个模型算向量
}
// retriever.Options
type Options struct {
Index *string
SubIndex *string
TopK *int
ScoreThreshold *float64
Embedding embedding.Embedder
DSLInfo map[string]any // 后端专属过滤表达式
}
三个设计点:
**① Embedding 在 Options 里,不只在 Config 里。**意味着可以调用时换模型:
store.Retrieve(ctx, query, retriever.WithEmbedding(otherModel))
配合第 71 篇《Embedder 缓存层》讲的「key 里带 model」,A/B 测两个模型不用建两套 store。
② 全是指针类型。TopK *int 而不是 TopK int——这样实现方能区分「调用方没传」和「调用方传了 0」。nil 才是「用你的默认值」。
③ DSLInfo map[string]any 是逃生门。过滤条件千差万别(ES 的 query DSL、Milvus 的表达式、SQL 的 WHERE),框架不可能抽象统一,索性开个 map[string]any 让后端自己定义结构。它跟第 67 篇《Document 组件源码》讲的 implSpecificOptFn 是两种不同的扩展机制:前者是数据(可序列化,能跨进程传),后者是函数(只能进程内用)。
二、eino-ext 后端实现的六步套路
读 eino-ext/components/{indexer,retriever}/redis 的实现,模式非常固定。以 Retrieve 为例:
func (r *Retriever) Retrieve(ctx context.Context, query string,
opts ...retriever.Option) (docs []*schema.Document, err error) {
// ① 通用选项:config 作默认值,调用时可覆盖
co := retriever.GetCommonOptions(&retriever.Options{
Index: &r.config.Index,
TopK: &r.config.TopK,
ScoreThreshold: r.config.DistanceThreshold,
Embedding: r.config.Embedding,
}, opts...)
// ② 实现专属选项
io := retriever.GetImplSpecificOptions(&implOptions{}, opts...)
// ③ callback 进场
ctx = callbacks.EnsureRunInfo(ctx, r.GetType(), components.ComponentOfRetriever)
ctx = callbacks.OnStart(ctx, &retriever.CallbackInput{
Query: query,
TopK: *co.TopK,
Filter: io.FilterQuery,
ScoreThreshold: co.ScoreThreshold,
})
// ④ defer 里报错
defer func() {
if err != nil {
callbacks.OnError(ctx, err)
}
}()
// ⑤ embedding + 双重校验
emb := co.Embedding
if emb == nil {
return nil, fmt.Errorf("[redis retriever] embedding not provided")
}
vectors, err := emb.EmbedStrings(r.makeEmbeddingCtx(ctx, emb), []string{query})
if err != nil {
return nil, err
}
if len(vectors) != 1 {
return nil, fmt.Errorf("[redis retriever] invalid return length of vector, got=%d, expected=1", len(vectors))
}
// ⑥ 拼后端查询 → 执行 → 转 Document → OnEnd
// ...
}
六步照抄就行。两个细节容易漏:
**GetCommonOptions 的第一个参数是 base,要把 config 填进去。**这是「构造时配默认值、调用时可覆盖」的实现方式。传 nil 的话调用方不传 TopK 就变成 nil 指针,解引用直接 panic。
len(vectors) != 1 这个校验不是多余的。EmbedStrings 的契约是「一条文本一个向量、顺序一致」,但这只是约定——某个 adapter 出 bug 少返一条,不校验的话下一行 vectors[0] 就是越界或者错位。第 71 篇的缓存层里也有同样的校验。
三、写 pgvector 版:向量参数怎么传
eino-ext 没有官方 pgvector 组件,所以自己写。第一个卡点就在传参:
_, err := db.ExecContext(ctx, "INSERT INTO chunks (embedding) VALUES ($1)", vector)
// vector 是 []float64 → database/sql 不认识这个类型
database/sql 只支持有限的几种 driver value(int64 / float64 / bool / []byte / string / time.Time)。[]float64 不在里面。
两条路:
- 引
pgvector-go,它注册了自定义类型,可以直接传pgvector.NewVector(v) - 手动转成 pgvector 的文本格式,让 PostgreSQL 自己转型
我选第二条,少一个依赖:
func vec2str(v []float64) string {
var sb strings.Builder
sb.Grow(len(v) * 8)
sb.WriteByte('[')
for i, f := range v {
if i > 0 {
sb.WriteByte(',')
}
sb.WriteString(strconv.FormatFloat(f, 'f', -1, 64))
}
sb.WriteByte(']')
return sb.String()
}
然后在 SQL 里显式转型:
INSERT INTO chunks (id, content, metadata, embedding)
VALUES ($1, $2, $3::jsonb, $4::vector)
$4::vector 那个转型不能省。传进去的是 text,PostgreSQL 需要知道往哪个类型转。
strconv.FormatFloat(f, 'f', -1, 64) 里的 -1 是精度:用最少的位数保证往返不丢精度。写 %f 会固定 6 位小数,1536 维的向量白扔一半精度还多占字节。
幂等写入
INSERT INTO chunks (...) VALUES (...), (...), ...
ON CONFLICT (id) DO UPDATE SET
content = EXCLUDED.content,
metadata = EXCLUDED.metadata,
embedding = EXCLUDED.embedding
实测同一批数据写两遍:
################ 1. Store:8 片,BatchSize=5 ################
写入 8 片:[annual sick overtime reimburse approval resign insurance device]
################ 2. 幂等:同样的数据再写一遍 ################
再写 8 片,表里总行数:8
**重跑不产生重复行。**索引任务经常要重试(网络抖动、embedding 超时、进程被杀),没有 ON CONFLICT 就得先 DELETE 再 INSERT,多一次往返还容易留下中间态。
多值 INSERT 的占位符编号
批量写要拼出 ($1,$2,$3,$4),($5,$6,$7,$8),...,编号是最容易错的地方:
for j, d := range batch {
if j > 0 {
sb.WriteByte(',')
}
n := j * 4 // 每行 4 个参数
fmt.Fprintf(&sb, "($%d,$%d,$%d::jsonb,$%d::vector)", n+1, n+2, n+3, n+4)
args = append(args, d.ID, d.Content, metaJSON, vec2str(vectors[j]))
}
n := j * 4 里的 4 必须和 args 里 append 的个数一致。加一列忘了改这个数字,报错信息是 bind message supplies N parameters, but prepared statement requires M——不难修但很费时间。
顺带提醒:PostgreSQL 的参数上限是 65535 个。每行 4 个参数的话,单批最多 16383 行。BatchSize 设太大会撞上这个墙。
四、ScoreThreshold 是相似度,不是距离
这是接 pgvector 时最容易搞反的地方。
retriever.Options.ScoreThreshold 的注释写得很清楚:
// ScoreThreshold is the score threshold for the retriever,
// eg 0.5 means the score of the document must be greater than 0.5.
分数,越大越相关。而 pgvector 的 <=> 返回的是距离,越小越相关(第 72 篇《pgvector 入门》讲过)。所以要换算:
相似度 >= t ⇔ 1 - 距离 >= t ⇔ 距离 <= 1 - t
代码里两处都要用相似度:
// SELECT 出来的 score 用相似度
fmt.Fprintf(&sb, `SELECT id, content, metadata,
1 - (embedding <=> $1::vector) AS score FROM %s WHERE embedding IS NOT NULL`, table)
// 阈值过滤也用相似度
if co.ScoreThreshold != nil {
args = append(args, *co.ScoreThreshold)
fmt.Fprintf(&sb, ` AND 1 - (embedding <=> $1::vector) >= $%d`, len(args))
}
// 但排序用距离(升序),这样才能走 HNSW 索引
fmt.Fprintf(&sb, ` ORDER BY embedding <=> $1::vector LIMIT $%d`, len(args))
**注意排序仍然用距离升序。**写成 ORDER BY 1 - (embedding <=> $1) DESC 逻辑上等价,但表达式变了,走不到 vector_cosine_ops 索引。
阈值不能拍脑袋定
实测四个阈值:
################ 5. ScoreThreshold(相似度阈值)################
阈值 0.0 → 返回 8 条
阈值 0.3 → 返回 0 条
阈值 0.5 → 返回 0 条
阈值 0.9 → 返回 0 条
0.3 就全空了。
因为这个 demo 用的是字袋 embedder,实际最高分只有 0.1796:
1. overtime score=0.1796
2. resign score=0.1217
3. reimburse score=0.0833
换成真实 embedding 模型,同一批文本的分数分布会完全不同(通常 0.5~0.9 之间,因为真实模型的向量不是稀疏计数)。
**所以「相似度低于 0.7 就丢弃」这种阈值不能抄。**必须先跑一遍看你的模型在你的语料上的实际分布,再定。第 70 篇《Embedding 选型》那套评测集顺手就能出这个直方图。
一个更稳的做法:先不设阈值,靠 TopK 控制数量。等积累了真实查询日志,再看「答对的查询」和「答错的查询」的分数分布,找一个能分开两者的值。
WHERE 过滤和 HNSW 后过滤是两码事
第 72 篇讲过 pgvector 的过滤是后过滤(LIMIT 5 可能只返回 3 行)。但要区分两种情况:
- 走 HNSW 索引时:先按向量取
ef_search个候选,再应用Filter,可能不够数 - 优化器放弃 HNSW 时:变成先
WHERE过滤再精确排序,结果精确、数量足
我这个实现里 metadata @> $n::jsonb 属于前者。租户/部门过滤要是很窄,还是得走第 72 篇说的物理隔离或分区。
五、embCtx:让 embedding 的打点归位
eino-ext 每个后端实现里都有个 makeEmbeddingCtx,我照抄了一个:
func (s *PGStore) embCtx(ctx context.Context, emb embedding.Embedder) context.Context {
typ := "Embedding"
if t, ok := components.GetType(emb); ok {
typ = t
}
return callbacks.ReuseHandlers(ctx, &callbacks.RunInfo{
Name: typ,
Type: typ,
Component: components.ComponentOfEmbedding,
})
}
作用是换掉 RunInfo,保留 handler。
不换会怎样?此刻 ctx 里的 RunInfo 是 {Component: Retriever, Type: PGVector}——embedding 的打点会被归到 Retriever 名下。做成本统计时,你会看到「Retriever 消耗了 N 个 token」,而实际那是 embedding 花的。
components.GetType(emb) 就是第 67 篇讲的 Typer 接口,让打点里显示 CharBag 而不是反射出来的 *main.charBag。
实测三层嵌套(我给假 embedder 也加了自己的打点):
################ 1. Store:8 片,BatchSize=5 ################
> [Indexer /PGVector ] OnStart
> [Embedding /CharBag ] OnStart
< [Embedding /CharBag ] OnEnd
> [Embedding /CharBag ] OnStart
< [Embedding /CharBag ] OnEnd
< [Indexer /PGVector ] OnEnd
Embedding 的打点出现了两次——8 片、BatchSize=5,分 2 批,每批一次 EmbedStrings。分批逻辑不用看日志,从 callback 就能数出来。
检索时是一次:
> [Retriever /PGVector ] OnStart
> [Embedding /CharBag ] OnStart
< [Embedding /CharBag ] OnEnd
< [Retriever /PGVector ] OnEnd
有个前提要说清楚:这个嵌套输出能出来,是因为 embedder 自己调了 callbacks.OnStart/OnEnd。embCtx 只负责准备 RunInfo,不会代打。第 67 篇讲的 IsCallbacksEnabled 就是这个约定——组件自己打点,框架不重复包。我最初的假 embedder 没打点,跑出来只有 Indexer 和 Retriever 两层,一开始还以为 embCtx 没生效。
六、传错类型的 option 会被静默忽略
第 67 篇讲过 GetImplSpecificOptions 里那个 if ok 的宽松取舍。实测一下:
hits, err := store.Retrieve(ctx, "请假需要什么材料",
retriever.WithTopK(2),
retriever.WrapImplSpecificOptFn(func(o *struct{ Foo string }) { o.Foo = "bar" }),
)
################ 7. 传错类型的 option 会怎样 ################
err=<nil>,返回 2 条 → 类型不匹配被静默忽略(第 67 篇讲的宽松取舍)
**不报错,正常返回。**那个类型不匹配的 option 被 if ok 挡掉了。
对使用者是个隐患:WithMetaFilter 写成别的后端的同名函数,过滤条件会静默失效——检索结果看起来正常,只是范围不对。
自查办法:可疑的时候在实现里把解析出来的 impl options 打出来看一眼。
七、参数化查询顶不顶得住
Retrieve 的 query 是用户输入,会被送去 embedding,本身不进 SQL。但 metadata 过滤条件可能来自上游。实测把 SQL 注入直接当查询扔进去:
evil := "'; DROP TABLE e73demo.chunks; --"
hits, err := store.Retrieve(ctx, evil, retriever.WithTopK(2))
################ 9. SQL 注入:query 里带引号 ################
err=<nil>,返回 2 条;表还在,行数 8
**表还在,8 行没少。**因为所有值都走 $n 占位符,由驱动做参数绑定,不参与 SQL 文本拼接。
要守住这条线,有一个地方必须小心:表名不能参数化。
fmt.Fprintf(&sb, `SELECT ... FROM %s WHERE ...`, s.cfg.Table)
%s 直接拼进 SQL。Table 来自 Config(代码里写死的常量),所以安全。但如果哪天做成「表名从 HTTP 请求里读」,这就是注入点。表名要动态,就得用白名单校验或者 pq.QuoteIdentifier。
八、一个类型同时实现两个接口
type PGStore struct {
cfg Config
}
var _ indexer.Indexer = (*PGStore)(nil)
var _ retriever.Retriever = (*PGStore)(nil)
eino-ext 是分成两个包的(indexer/redis 和 retriever/redis),因为它们要独立版本化。自己写的话合成一个类型更省事:
**① 共享连接池。**两个包各自持有 *sql.DB 的话,连接数翻倍,还得配两次。
② 维度校验只写一次。
if got := len(vectors[0]); got != s.cfg.VectorDim {
return nil, fmt.Errorf("[pgvector indexer] dim mismatch: table=%d model=%d",
s.cfg.VectorDim, got)
}
实测:
################ 8. 维度不匹配的报错 ################
err=[pgvector indexer] dim mismatch: table=128 model=64
**在 Go 层就拦下了,没走到数据库。**PostgreSQL 也会拦(第 72 篇实测过 expected 128 dimensions, not 3),但错误信息里没有「表期望多少、模型给了多少」这个对照,排查时要多翻两层。
**③ 写和读用同一个 Embedding 配置。**第 66 篇强调过存查必须同模型。分成两个对象,配错的机会就多一倍。
代价是这个类型同时被两个接口约束,改动时要照顾两边。规模大了再拆。
小结
- Options 全用指针,为了区分「没传」和「传了零值」;
Embedding也在 Options 里,所以能调用时换模型 DSLInfo map[string]any是数据形态的逃生门,implSpecificOptFn是函数形态的,前者可序列化跨进程- 六步套路:
GetCommonOptions(config 作 base)→GetImplSpecificOptions→EnsureRunInfo+OnStart→defer OnError→ embedding 双重校验 → 拼查询/转 Document/OnEnd database/sql不认识[]float64:转成[1,2,3]文本 +$n::vector显式转型,FormatFloat(f,'f',-1,64)保精度ON CONFLICT DO UPDATE让索引任务可重跑,实测 8 片写两遍还是 8 行ScoreThreshold是相似度:sim >= t ⇔ dist <= 1-t;排序仍用距离升序才能走 HNSW 索引- 阈值不能抄别人的:实测字袋模型最高分 0.1796,阈值 0.3 就返回 0 条。先看你自己的分数分布
embCtx只换RunInfo不代打点,打点得组件自己做;换了之后 embedding 的成本才不会算到 Retriever 头上- 类型不匹配的 option 静默忽略,过滤条件可能悄悄失效
- 参数化守住了注入,但表名走的是
%s拼接,别让它来自用户输入
下一篇(第 74 篇《多查询 + 重排序》)回到效果:检索结果不理想时,多查询改写(MultiQuery)和重排序(Reranker)各自解决什么问题,以及它们分别在什么情况下不起作用。
代码状态说明:
PGStore(约 200 行,同时实现indexer.Indexer+retriever.Retriever)在 PostgreSQL 18.4 + pgvector 0.8.2 + eino v0.9.13 + lib/pq v1.12.3 上真机运行,本文九组实测输出原样粘贴。实验用临时 schemae73demo,跑完已清理。embedding 用的是本地确定性字袋实现(不需要 API Key),所以分数的绝对值没有参考意义——它只用来验证接口链路、参数传递和阈值换算是否正确。第二节eino-ext的六步套路引自components/retriever/redis/retriever.go。
更多推荐
所有评论(0)