系列「企业级 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 不在里面。

两条路:

  1. pgvector-go,它注册了自定义类型,可以直接传 pgvector.NewVector(v)
  2. 手动转成 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/OnEndembCtx 只负责准备 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/redisretriever/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)→ GetImplSpecificOptionsEnsureRunInfo+OnStartdefer 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 上真机运行,本文九组实测输出原样粘贴。实验用临时 schema e73demo,跑完已清理。embedding 用的是本地确定性字袋实现(不需要 API Key),所以分数的绝对值没有参考意义——它只用来验证接口链路、参数传递和阈值换算是否正确。第二节 eino-ext 的六步套路引自 components/retriever/redis/retriever.go

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐