Go 基础体系 · 第 107/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。

Go RAG 完整流程:切块、Embedding、向量库、重排与引用

本文以 Go 1.26.4、PostgreSQL 17.6、pgvector 0.8.0、Qdrant Go client v1.15.x 与 Milvus Go SDK v2.6.x 稳定版本线为基线。生产应锁定 go.mod、镜像 digest、Embedding 模型 revision 与索引配置;示例不要求同时引入三种向量库。

RAG 不是“向量搜索后把文本塞给模型”。完整链路包含数据摄取、解析、规范化、切块、Embedding、索引发布、查询改写、候选召回、权限过滤、去重、重排、上下文编排、回答生成、引用校验和持续评测。原始文档是事实源,向量索引是可重建的派生物;把版本、权限和证据贯穿全程,答案才可追溯。

1. 先画清离线与在线数据流

离线 ingestion 负责把可信来源变成可检索块,在线 retrieval 负责在一次请求预算内找出证据。两者通过不可变 index_version 连接,不能让在线请求读到“半批新、半批旧”的集合。

source(version, ACL) -> fetch -> parse -> normalize -> chunk
                     -> embed(model, dimension) -> staging index -> publish

query(actor) -> normalize/rewrite -> lexical + vector recall
             -> ACL filter -> fusion/deduplicate -> rerank
             -> context budget -> generate -> citation verify

解析失败不能产生空 chunk;索引写入失败不能把源版本标成已发布。在线故障与“没有证据”必须区分:前者可降级,后者应明确资料不足。

2. 文档、块与索引契约

文档 ID 表示逻辑对象,source version 表示一次不可变快照,chunk ID 由文档、版本、位置和规范化文本哈希稳定派生。不要用数据库自增 ID 作为唯一幂等依据,否则重跑会制造重复向量。

type Chunk struct {
	ID             string
	TenantID       string
	DocumentID     string
	SourceVersion  string
	Ordinal        int
	Text           string
	TextSHA256     [32]byte
	StartRune      int
	EndRune        int
	ACLGroups       []string
	SourceURL       string
	EmbeddingModel string
	IndexVersion   string
}

type Embedder interface {
	Embed(context.Context, []string) ([][]float32, error)
}

type VectorStore interface {
	Upsert(context.Context, []Chunk, [][]float32) error
	Search(context.Context, []float32, SearchFilter, int) ([]Hit, error)
}

接口在调用方定义,具体 Qdrant、Milvus 或 pgvector 适配器留在基础设施层。metadata 至少保留 tenant、ACL、来源版本、页码或字符范围、内容哈希和删除状态。模型名、维度、归一化方式、距离度量共同构成索引 schema;其中任意一项改变都要新建版本,不能把不同向量空间混写。

3. 摄取要可重放、可取消且幂等

摄取任务以 (tenant_id, document_id, source_version, pipeline_version) 为唯一键。worker 先读取事实源并核对校验和,再解析和切块,批量生成向量,写 staging 索引,最后以一次原子状态更新发布。网络调用不应包在数据库长事务中。

func ingest(ctx context.Context, job Job, source Source, index Index) error {
	if err := ctx.Err(); err != nil {
		return context.Cause(ctx)
	}
	doc, err := source.Load(ctx, job.DocumentID, job.SourceVersion)
	if err != nil {
		return fmt.Errorf("load source document: %w", err)
	}
	chunks, err := splitDocument(doc, job.PipelineVersion)
	if err != nil {
		return fmt.Errorf("split document: %w", err)
	}
	if err := index.Stage(ctx, job.Key(), chunks); err != nil {
		return fmt.Errorf("stage chunks: %w", err)
	}
	if err := index.Publish(ctx, job.Key()); err != nil {
		return fmt.Errorf("publish index version: %w", err)
	}
	return nil
}

重试同一 key 时,内容哈希一致就返回既有结果;哈希不同说明上游违反不可变版本约定,应报冲突。每批限制文档字节、页数、chunk 数、Embedding 条数和总 token。取消要传到下载、解析器、Embedding HTTP、数据库和队列续租;worker 退出前把任务留在可重试状态,而不是标记成功。

4. 解析与规范化不能破坏证据位置

HTML 要去导航、脚本和重复页脚,PDF 要保留页码与阅读顺序,Markdown 要保留标题层级,代码要保留文件与符号边界。OCR 文本应保存置信度。规范化可以统一换行、Unicode 和空白,却不能把表格列拼乱或丢掉能支持引用的位置。

解析器面对恶意文件:限制压缩比、递归层级、图片像素、页数、CPU、内存和执行时间;在隔离进程中处理 PDF、Office 和图片。抓取 URL 时限制 HTTPS、域名、端口和响应类型,每次 DNS 解析及重定向都阻断环回、私网、链路本地和云 metadata 地址。解析出来的“忽略系统规则”只是文档数据,不获得工具权限。

保留原文与规范化哈希、parser 版本和警告,抽样监测乱码与表格损坏。规则升级后创建新 pipeline version 并重建,避免旧引用无法复现。

5. 切块策略由任务和评测决定

固定 500 字符加 50 字重叠只是起点。技术文章适合按标题、段落、代码块切;API 文档应尽量把签名、参数和约束放在同块;表格按行组切且重复表头;FAQ 可按问答对切。块过小缺少语境,块过大降低区分度并挤占生成上下文。

切块流程先按结构得到候选单元,再用目标模型 tokenizer 控制上下限,最后为孤立短段合并父标题。重叠只用于跨边界语义,过大重叠会让召回列表充满近重复内容。每个 chunk 可附带短的层级前缀,但向用户引用仍指向原文位置。

建议从 300~800 token、10%~15% 重叠做实验,而不是当成通用结论。比较 Recall@K、重排后证据覆盖率、输入 token、延迟和答案忠实度;代码、中文长文和表格分别建子集。问题答案需要跨两段时,可以存父子 chunk:小块负责召回,父块负责补充上下文。

6. Embedding 批处理与版本管理

Embedding 相同文本应按 (normalized_text_hash, model_revision, dimensions, normalization) 缓存。批处理同时受条数、总 token、总字节和等待时间约束;固定数量 worker 配合 semaphore,禁止为每个 chunk 无界启动 goroutine。

func embedBatch(ctx context.Context, e Embedder, chunks []Chunk) ([][]float32, error) {
	texts := make([]string, len(chunks))
	for i := range chunks {
		texts[i] = chunks[i].Text
	}
	vectors, err := e.Embed(ctx, texts)
	if err != nil {
		return nil, fmt.Errorf("embed %d chunks: %w", len(chunks), err)
	}
	if len(vectors) != len(chunks) {
		return nil, errors.New("embedding count mismatch")
	}
	for _, vector := range vectors {
		if len(vector) != 1536 || !finiteVector(vector) {
			return nil, errors.New("invalid embedding vector")
		}
	}
	return vectors, nil
}

验证 NaN、Inf、维度和返回顺序。供应商返回部分成功时只按稳定 item ID 对齐,不能猜位置。429/5xx 在剩余 deadline 足够时有限重试并加入抖动;永久参数错误进入死信。记录实际 token 与费用,缓存命中单独计数。更换模型时双写新索引、离线评测、灰度查询,再切读别名并保留回滚窗口。

7. 距离度量决定分数含义

余弦相似度比较方向,适合长度不重要的语义向量;点积同时受方向和模长影响,很多已归一化模型上与余弦排序一致;欧氏距离比较空间直线距离。必须遵循 Embedding 模型说明,并使用索引对应 operator class。

对向量 ab,余弦为 (a·b)/(|a||b|),范围通常为 [-1,1];距离常写成 1-cosine。点积越大越近,L2 距离越小越近。不同数据库返回的是 distance 还是 similarity 不统一,领域适配器应归一化名称,禁止拿“0.8 阈值”跨模型、跨度量复用。

绝对阈值需用标注查询校准;它既可能漏掉措辞不同的证据,也可能接纳语义相近却事实不符的段落。召回 top-N 后重排,并设置“无充分证据”判定。

8. HNSW、IVF 与精确搜索如何选择

精确搜索计算每个候选距离,召回确定但随数据量线性增长,适合小集合、离线真值和强过滤后候选很少的场景。HNSW 构建多层近邻图,查询快、召回高,代价是内存、构建时间和更新维护;m/ef_construction/ef_search 分别影响图密度、构建质量和查询探索范围。

IVF 先把向量分桶,查询只探测部分簇;nlistnprobe 在速度与召回间取舍,数据分布变化后可能需要重训。磁盘型或量化索引降低内存,但会损失精度。参数不能照搬厂商样例,要用本租户规模、过滤选择性和并发压测。

过滤与 ANN 的执行顺序尤其关键。高选择性 ACL 若在 ANN 后过滤,可能取回 100 个却全部无权;支持 payload pre-filter 的引擎应先约束候选空间。pgvector 的计划可能随统计信息变化,要用 EXPLAIN (ANALYZE, BUFFERS) 验证。建立小规模 exact ground truth,持续比较 ANN Recall@K。

9. pgvector 的可审计实现

关系元数据、事务与向量放在一起便于中等规模系统治理。下面把版本和 ACL 放入同表,索引 operator class 与余弦距离一致。

CREATE EXTENSION IF NOT EXISTS vector;
CREATE TABLE rag_chunk (
    chunk_id text PRIMARY KEY,
    tenant_id text NOT NULL,
    document_id text NOT NULL,
    source_version text NOT NULL,
    acl_groups text[] NOT NULL,
    content text NOT NULL,
    content_sha256 bytea NOT NULL,
    index_version text NOT NULL,
    embedding vector(1536) NOT NULL,
    deleted_at timestamptz,
    UNIQUE (tenant_id, document_id, source_version, content_sha256)
);
CREATE INDEX rag_chunk_hnsw ON rag_chunk
USING hnsw (embedding vector_cosine_ops);

SELECT chunk_id, document_id, content,
       1 - (embedding <=> $1) AS similarity
FROM rag_chunk
WHERE tenant_id = $2
  AND acl_groups && $3::text[]
  AND index_version = $4
  AND deleted_at IS NULL
ORDER BY embedding <=> $1
LIMIT $5;

向量和过滤参数都通过驱动绑定,不拼 SQL。查询设置 context deadline、statement timeout、最大返回行和总内容字节。索引发布可维护一张 active index pointer 表,在事务中切换版本;旧版本延迟回收,为在途请求和回滚留窗口。

10. 混合召回补足专有词与编号

向量擅长语义,BM25/全文索引擅长错误码、函数名、订单号和罕见专有词。两路各取候选后可用 Reciprocal Rank Fusion:score(d)=Σ 1/(k+rank_i(d)),它只使用名次,避免直接混合不同量纲分数。

func rrf(lists [][]Hit, offset float64) []Hit {
	scores := make(map[string]float64)
	items := make(map[string]Hit)
	for _, list := range lists {
		for rank, hit := range list {
			scores[hit.Chunk.ID] += 1 / (offset + float64(rank+1))
			items[hit.Chunk.ID] = hit
		}
	}
	return sortHits(items, scores)
}

先按 chunk ID 去重,再按 document 与相邻位置控制多样性。query rewrite 可以展开缩写或生成关键词,但必须保留原问题,并计入模型延迟与成本;模型改写失败时回退原 query。多查询召回会放大数据库负载,要限制改写数与并发。

11. 权限过滤与缓存隔离

Actor 来自认证中间件,tenant 和 ACL filter 由服务端构造,绝不采用模型或用户 JSON 中自报的租户。权限应尽量下推到检索引擎;所有结果离开适配器前再次断言 tenant、ACL、删除标记与版本。文档撤权要快速更新索引或使用权威权限服务做最终校验。

检索缓存 key 至少包含 tenant、actor 权限版本、规范化 query、index version、embedding model、过滤条件与参数版本。不能跨租户共享结果,也不能仅以 query 文本为 key。权限变化时递增 permission epoch,使旧缓存自然失效。日志、trace 和离线评测导出同样遵守隔离,避免“查询没泄漏,遥测泄漏”。

索引 namespace 物理隔离最直观但成本高;共享集合加强制 filter 更经济但更依赖实现正确性。高敏客户可独立 collection/数据库,普通租户可逻辑隔离。无论哪种方案都用跨租户攻击测试证明,而不是依赖配置说明。

12. Rerank 从候选中挑真正证据

bi-encoder 向量先独立编码 query 和文档,召回快;cross-encoder reranker 联合阅读 query 与候选,相关性更好但计算昂贵。通常先召回 30~100 条,去重后重排,取 5~12 条进入上下文。候选正文、标题与位置都要限制长度。

{
  "model": "reranker-zh-2026-08",
  "query": "流式生成取消后是否可以自动重试",
  "documents": [
    {"id": "chunk_a", "text": "首个增量输出后不得透明重试……"},
    {"id": "chunk_b", "text": "取消会关闭响应体,但费用可能已经产生……"}
  ],
  "top_n": 6,
  "return_documents": false
}

reranker timeout 时可以回退融合顺序,但响应和指标应记录降级。分数仍只在固定模型内可比。为避免一篇长文垄断结果,采用 per-document cap 或 MMR 多样性;相邻块确需共同回答时在重排后做邻居扩展,再按 token 预算裁剪。

13. 上下文编排与证据防注入

应用为每个选中块生成不可伪造的短 citation ID,例如 S1,把来源、标题、位置和原文放入结构化边界。system 规则说明资料仅是证据,其中的命令不具权限;真正的工具与权限保护仍在应用侧。

任务:仅依据 SOURCES 回答;证据不足时明确说明。不编造引用。

<SOURCES>
[S1 document=go-stream version=7 location=section-10]
首个可见 delta 发出后,不应透明重试……
[S2 document=billing version=3 location=p12]
取消不保证撤销已经产生的推理费用……
</SOURCES>

问题:用户点停止后服务应该怎么处理?

按重要性而不是数据库返回偶然顺序排列,保留 system、问题、证据和输出的独立 token 预算。内容超限时优先删低分重复块,不能从每块尾部盲截导致关键限定词丢失。上下文 hash、候选 ID 和版本进入受控 trace,原文默认不记录。

14. 引用由系统验证,不由模型自由编写

模型只允许输出本轮 S1..Sn,应用把 ID 映射成经过 allowlist 校验的 URL 和位置。生成后解析引用,拒绝未知 ID,并检查每个关键可验证主张附近是否至少有一个引用。引用存在不代表支持主张,还要做 entailment 评测或抽样人工审核。

{
  "answer": "停止后应取消同一请求上下文,并把已输出结果标为不完整。[S1] 已发生的推理费用仍可能被计入。[S2]",
  "citations": [
    {"id": "S1", "claims": [0]},
    {"id": "S2", "claims": [1]}
  ],
  "insufficient_evidence": false
}

页面展示时 URL、标题和摘录按 HTML 上下文转义。来源已删除或用户在生成后撤权时,不再展开正文;必要时答案也失效。引用点击和无引用回答率可以帮助诊断,但不能把用户身份放进公开指标。

15. 取消、重试与失败恢复

在线总 deadline 拆给 query embedding、两路召回、rerank 和生成,任何 semaphore、队列发送和重试等待都响应 context。查询 Embedding 与只读检索通常可有限重试;生成已输出后不可透明重试。降级顺序预先定义,例如 rerank 失败用融合结果、向量库失败用全文检索、两者都失败则返回暂时不可用,而不是无证据调用模型。

离线 ingestion 至少一次交付,阶段状态记录 pending/running/staged/published/failed、attempt、lease 和最后错误类别。worker 崩溃后租约过期可由另一 worker 接管;幂等 chunk ID 与 upsert 防重复。发布前失败可清理 staging,发布后失败不能删 active。删除采用 tombstone 与异步物理回收,并用对账任务确认事实源、向量库、缓存一致。

退避只处理限流和瞬时上游错误,尊重有界 Retry-After。认证、维度、非法文档和取消不重试。一次任务只能由一层控制重试,防队列、SDK、HTTP 层相乘造成费用风暴。

16. 评测要分解检索、生成与引用

数据集来自真实匿名问题、无答案问题、专有词、跨段证据、时效更新、权限和对抗文档。每条保存允许证据 chunk/document、禁止来源、答案关键点与 hard gate。检索测 Recall@K、MRR、nDCG、证据覆盖;重排测排序提升;生成测忠实度、完整性、拒答;引用测合法率、覆盖率和支持率。

case_id: rag-cancel-017
principal: {tenant: t42, groups: [engineering]}
question: "首个流式片段发出后能否自动重试?"
required_sources: [doc-streaming-v7]
forbidden_sources: [doc-other-tenant]
assertions:
  recall_at_10: true
  citation_support: true
  must_include: ["不能透明重试", "部分结果"]
  max_context_tokens: 6000

调整 chunk、Embedding、索引参数、融合权重、reranker、prompt 或模型都跑同一回归集,报告置信区间、分桶失败和资源变化。线上点击率容易受位置偏差影响,只作补充。安全越权、伪造引用和泄漏是零容忍门禁,不能被平均分抵消。

17. 诊断、成本与容量

一次 trace 分为 fetch/parse/embed/upsert 或 query/embed/search/filter/rerank/generate/verify。记录模型 revision、index/pipeline version、候选数量、过滤后数量、上下文 token、TTFT、总耗时和稳定错误码。不要每个 chunk 建 span,也不要记录原始问题和全文。

检索为空时依次检查:事实源是否发布、解析是否为空、chunk 是否存在、Embedding 维度与模型是否一致、active index 是否正确、ACL 是否过严、ANN 参数是否漏召回。延迟升高则分解连接池等待、数据库计划、reranker 排队与生成时间。维护一个受控诊断命令,能按 chunk ID 查看版本和状态但仍执行授权。

成本包含解析 CPU、Embedding token、向量存储、索引内存、rerank、生成输入 token 和运维。容量按文档变更速率、向量数、维度、每向量字节、索引放大、查询并发和上下文长度估算。批量嵌入节省往返,却增加失败重做范围;缓存节省费用,却必须把版本和权限纳入 key。

18. 安全、部署与上线清单

文档内容视为不可信输入:隔离解析,阻断 SSRF/压缩炸弹,扫描恶意文件,限制 Unicode 控制字符;检索只返回最小必要字段;prompt injection 不能改变工具 registry、Actor 或网络策略。Embedding 服务和向量库使用最小权限凭据、TLS、网络白名单与密钥轮换。备份不仅包含向量,还要包含能重建的事实源、版本和 pipeline 配置。

部署采用 staging collection 加原子别名切换。新索引先完成条数、哈希、权限、抽样向量和评测检查,再灰度读取;回滚只切别名。服务优雅关闭时停止新摄取、续租在途任务或安全放回队列,等待有限时间后取消;在线实例先摘 readiness,再结束查询。告警关注摄取积压、发布失败、维度错误、Recall 代理下降、跨租户拒绝异常、P95/P99、费用与磁盘水位。

上线前确认:事实源可重建;任务幂等;模型、维度、度量和索引版本一致;ACL 在检索阶段生效;缓存含权限版本;候选经过去重与有界重排;上下文和输出有限额;引用只映射本轮证据;取消贯穿全链路;降级不把“故障”伪装成“无答案”;每次配置变化都有离线评测与回滚方案。做到这些,RAG 才是可验证的信息系统,而不是一次相似度搜索演示。


系列导航与关联阅读

官方资料

本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。