目录 · 第 25 / 28 章
EinoPart V · 综合实战与工程化

25RAG 全流程 Capstone

answer_from_document 工具:一张 load→chunk→score→filter→answer 的 Workflow,字段级扇入 + BatchNode 并发打分。

quickstart/chatwitheino/rag/rag.go:122quickstart/chatwitheino/rag/rag.go:158quickstart/chatwitheino/rag/rag.go:124eino-examples/adk/common/tool/graphtool/graph_tool.go:42

一个真实的 RAG 管线,竟然是一张 Workflow

前面所有讲编排的章节,到这里一起下场。我们读透 chatwitheino 里的 answer_from_document 工具——一个”从上传文档里检索并合成带引用答案”的 RAG 管线。别人的 RAG 常常是一坨胶水代码:读文件、切块、打分、排序、拼提示词,顺手用 for 循环串起来。而 Eino 把它写成一张 compose.Workflow:load → chunk → score → filter → answer 五个节点。读它的价值,不在任何单个 API,而在于看 Part IV 的机制——字段映射、BatchNode 并发、扇入扇出——如何在一个真实管线里咬合,并最终被整体包成 Agent 的一个工具。

工具的契约:两个 struct 定住输入输出

先看这个工具对外长什么样。它的输入输出各是一个 struct,JSON tag 上的 jsonschema 描述会被 utils.InferTool 读出来,自动生成工具的参数 schema——你不用手写 JSON Schema(quickstart/chatwitheino/rag/rag.go:64):

// Input 的 JSON tag 被 utils.GoStruct2ToolInfo 用来自动生成参数 schema
type Input struct {
FilePath string `json:"file_path" jsonschema:"description=Absolute path to the uploaded document file"`
Question string `json:"question" jsonschema:"description=The question to answer from the document"`
}
// Output 是工具返回的结构化结果
type Output struct {
Answer string `json:"answer"`
Sources []string `json:"sources"` // 生成答案时引用的关键片段
}

对模型而言,这个工具就叫 answer_from_document,收 {file_path, question}、还 {answer, sources}。管线内部有多复杂,对调用方完全透明——这正是”把编排包装成工具”(Ch24)的意义:一整张图,对外只是一个可调用的工具

五个节点:把 RAG 拆成一条 Workflow

buildWorkflow 把这五步显式画成一张 Workflow(quickstart/chatwitheino/rag/rag.go:122)。注意每个节点后面挂的 AddInput——它声明这个节点的数据从哪来:

// buildWorkflow constructs the RAG compose.Workflow (uncompiled).
// graphtool.NewInvokableGraphTool compiles it per invocation.
func buildWorkflow[M adk.MessageType](cm model.BaseModel[M]) *compose.Workflow[Input, Output] {
scoreWF := newScoreWorkflow(cm)
scorer := batch.NewBatchNode(&batch.NodeConfig[scoreTask, scoredChunk]{
Name: "ChunkScorer",
InnerTask: scoreWF,
MaxConcurrency: 5,
})
wf := compose.NewWorkflow[Input, Output]()
// load: read file from disk, emit a single Document.
wf.AddLambdaNode("load", compose.InvokableLambda(
func(ctx context.Context, in Input) ([]*schema.Document, error) {
data, err := os.ReadFile(in.FilePath)
if err != nil {
return nil, fmt.Errorf("read %q: %w", in.FilePath, err)
}
return []*schema.Document{{Content: string(data)}}, nil
},
)).AddInput(compose.START)
// chunk: split each Document into ~800-char pieces.
wf.AddLambdaNode("chunk", compose.InvokableLambda(
func(ctx context.Context, docs []*schema.Document) ([]*schema.Document, error) {
var out []*schema.Document
for _, d := range docs {
out = append(out, splitIntoChunks(d.Content, 800)...)
}
return out, nil
},
)).AddInput("load")
// score: score each chunk against the question in parallel via BatchNode.
// Chunks comes from "chunk"; Question comes directly from START.
// Both use WithNoDirectDependency because the execution order is already
// established by the direct edges START→load→chunk→score.
wf.AddLambdaNode("score", compose.InvokableLambda(
func(ctx context.Context, in scoreIn) ([]scoredChunk, error) {
tasks := make([]scoreTask, len(in.Chunks))
for i, c := range in.Chunks {
tasks[i] = scoreTask{Text: c.Content, Question: in.Question}
}
return scorer.Invoke(ctx, tasks)
},
)).
AddInputWithOptions("chunk",
// … 省略 53 行;完整声明 L120–220,点击上方「浏览完整文件」
wf := compose.NewWorkflow[Input, Output]()
// load:读文件,产出一个 Document
wf.AddLambdaNode("load", /* os.ReadFile */).AddInput(compose.START)
// chunk:按段落切成 ~800 字的块
wf.AddLambdaNode("chunk", /* splitIntoChunks */).AddInput("load")
// score:对每块并发打分(下面详解字段映射)
wf.AddLambdaNode("score", /* BatchNode 打分 */). /* 字段映射扇入 */
// filter:按分数排序,取 Top-3
wf.AddLambdaNode("filter", /* sort + top-k */).AddInput("score")
// answer:从 Top-K 合成带引用的答案(字段映射扇入)
wf.AddLambdaNode("answer", /* synthesize */). /* 字段映射扇入 */
wf.End().AddInput("answer")

load → chunk → filter 是最朴素的线性依赖:一个节点的整块输出,喂给下一个节点。真正值得学的是 scoreanswer 这两个节点——它们的输入不是”上一个节点的整块输出”,而是从多个上游拼装出来的

字段级扇入:让不相邻的节点共享 Question

score 节点要打分,需要两样东西:待打分的块(来自 chunk)和原始问题(来自最开头的 START)。问题是,chunk 的输出类型是 []*schema.Document,里面根本没有 Question——如果硬要让 Question 一路穿过 load、chunk 每个节点的输出类型,会污染整条管线。

Workflow 的字段映射(field mapping)解决了这个问题:一个节点的输入 struct,可以按字段从不同上游分别取(quickstart/chatwitheino/rag/rag.go:158):

wf.AddLambdaNode("score", compose.InvokableLambda(
func(ctx context.Context, in scoreIn) ([]scoredChunk, error) {
// in.Chunks 来自 chunk;in.Question 来自 START
...
},
)).
AddInputWithOptions("chunk",
[]*compose.FieldMapping{compose.ToField("Chunks")}, // chunk 整块输出 → scoreIn.Chunks
compose.WithNoDirectDependency()).
AddInputWithOptions(compose.START,
[]*compose.FieldMapping{compose.MapFields("Question", "Question")}, // START.Question → scoreIn.Question
compose.WithNoDirectDependency())

scoreIn{Chunks, Question} 这个 struct 就是这样被”拼”出来的:Chunks 字段接 chunk 的整块输出,Question 字段直接从 STARTQuestion 字段拉过来。answer 节点同理——它的 synthIn{TopK, Question} 里,TopKfilter,Question 又一次从 START 直投(quickstart/chatwitheino/rag/rag.go:198)。于是原始问题只在需要它的两个节点上按字段接线,而不是被迫穿过整条管线。

这里的 WithNoDirectDependency 是点睛之笔:它表示”我要这个上游的数据,但不额外加一条执行依赖边”。因为执行顺序已经由 START → load → chunk → score 这条主链保证了,再给 score 加一条”必须等 START”的边纯属冗余——数据流和控制流在这里被干净地拆开

🔑 本章的设计钥匙

一条 RAG 管线,不是把步骤堆在一起,而是用 Workflow 的字段映射把”数据依赖”精确地表达到字段粒度scoreanswer 都要用到最开头的 Question,但 Question 不必污染 loadchunk 的输出类型——它只在真正需要的节点上,通过 MapFields 单独接一根线。WithNoDirectDependency 更进一步:把数据依赖(要谁的值)和控制依赖(等谁跑完)拆成两件事。这就是全书反复出现的钥匙在数据层的体现:图的形状,就是数据流的形状——只不过这次精确到了 struct 的单个字段。

并发打分:BatchNode 把扇出扇入收进一个节点

score 内部藏着本书最核心的一个动作:扇出扇入。文档切成几十上百块,每块都要独立地问一次模型”这块和问题有多相关(0–10 分)“。串行打分会慢到无法忍受,所以这里用 batch.NewBatchNode 把它并发化(quickstart/chatwitheino/rag/rag.go:124):

scoreWF := newScoreWorkflow(cm) // 单块打分的内层 workflow
scorer := batch.NewBatchNode(&batch.NodeConfig[scoreTask, scoredChunk]{
Name: "ChunkScorer",
InnerTask: scoreWF,
MaxConcurrency: 5, // 最多 5 个块同时打分
})

BatchNode 就是 Part IV “扇出扇入”的一个封装:输入一个 []scoreTask,它扇出成 N 个并发子任务(每个跑一遍 InnerTask),再把 N 个 scoredChunk 结果扇入收回成一个切片。MaxConcurrency: 5 是背压闸门——再多的块也只放 5 个同时打模型,避免打爆下游 QPS。这正是本章配套交互演示里”一份产出扇出成多路、再归并回一路”的真实落点。

打分函数本身还演示了一条工程铁律:局部失败不拖垮整批。某块的模型调用报错、或返回的 JSON 解析不了,都被降级成”0 分”而非中断整个 batch(quickstart/chatwitheino/rag/rag.go:158)——一块打分失败,不该让整篇文档的问答崩掉。

filter 与 answer:从”全部块”收敛到”一个答案”

filter 把打完分的块按分数降序排,取分数 ≥ 3 的前 3 块(quickstart/chatwitheino/rag/rag.go:175)——这是 RAG 里”召回后重排、只留最相关”的一步。answer 则做最后合成:如果 TopK 为空(文档里压根没相关内容),直接返回一句”没找到”;否则把 Top-K 片段拼成提示词,让模型生成一段带引用编号的答案,并把用到的片段原样放进 Output.Sources(quickstart/chatwitheino/rag/rag.go:198)。

至此,一条完整的 RAG 管线——读、切、并发打分、重排、合成——全部用 Workflow 的节点和边表达完毕,没有一行游离的编排胶水。

把整张 Workflow 包成一个工具

最后一步,是把这张 Workflow 变成模型能调用的工具。BuildToolgraphtool.NewInvokableGraphTool 完成这一步(quickstart/chatwitheino/rag/rag.go:109):

func BuildTool[M adk.MessageType](ctx context.Context, cm model.BaseModel[M]) (tool.BaseTool, error) {
wf := buildWorkflow(cm)
return graphtool.NewInvokableGraphTool[Input, Output](
wf, "answer_from_document",
"Search a large uploaded document ... synthesize a cited answer ...",
)
}

NewInvokableGraphTool 做了两件事(eino-examples/adk/common/tool/graphtool/graph_tool.go:42):一是用 utils.GoStruct2ToolInfo[Input]Input 的 tag 生成工具 schema;二是每次调用时按需编译并运行这张图,还内置了 checkpoint store——这意味着这张 RAG 图可以在执行中途被中断、之后再恢复(下一章的人工审批就靠它)。于是”一整张 Workflow”和”一个 Agent 工具”之间,只隔着这一层 graphtool 适配器。

func NewInvokableGraphTool[I, O any](compilable Compilable[I, O],
name, desc string,
opts ...compose.GraphCompileOption,
) (*InvokableGraphTool[I, O], error) {
tInfo, err := utils.GoStruct2ToolInfo[I](name, desc)
if err != nil {
return nil, err
}
return &InvokableGraphTool[I, O]{
compilable: compilable,
compileOptions: opts,
tInfo: tInfo,
}, nil
}
flowchart TB
  START(["START · {FilePath, Question}"]) --> LOAD["load · 读文件"]
  LOAD --> CHUNK["chunk · 段落切块"]
  CHUNK -->|"Chunks"| SCORE["score · BatchNode 并发打分 (≤5)"]
  SCORE --> FILTER["filter · 取 Top-3 (score≥3)"]
  FILTER -->|"TopK"| ANSWER["answer · 合成带引用答案"]
  ANSWER --> STOP(["END · {Answer, Sources}"])
  START -.->|"Question 字段直投"| SCORE
  START -.->|"Question 字段直投"| ANSWER

answer_from_document:一张 RAG Workflow

本章小结

  • answer_from_document 是一条完整 RAG 管线(读→切→打分→重排→合成),但它不是胶水代码,而是一张 compose.Workflow(quickstart/chatwitheino/rag/rag.go:122)。
  • 工具的输入输出用两个带 jsonschema tag 的 struct 定死,schema 由 utils.InferTool 自动生成(quickstart/chatwitheino/rag/rag.go:64)。
  • 字段级扇入:scoreanswerAddInputWithOptions + MapFields 从多个上游按字段拼装输入,让 Question 只接到需要它的节点,而非污染整条管线(quickstart/chatwitheino/rag/rag.go:158)。
  • WithNoDirectDependency数据依赖控制依赖拆开:要值,但不加冗余的执行边。
  • BatchNode(MaxConcurrency: 5)是 Part IV 扇出扇入的封装:一批块扇出成并发打分、再归并回一个切片;局部失败降级为 0 分而不拖垮整批(quickstart/chatwitheino/rag/rag.go:124)。
  • 整张 Workflow 经 graphtool.NewInvokableGraphTool 包成一个可调用、可中断/恢复的工具(eino-examples/adk/common/tool/graphtool/graph_tool.go:42)。

下一章,我们给这个 Agent 装上”眼睛”:callback 驱动的可观测性、Eino Dev 可视化调试,以及如何部署上线。

源码

正在读取完整文件…