Day 09 / 共 20 天 · 第 2 周 编排引擎

Graph 编译与执行

图搭好了怎么跑?今天讲 Compile 把图变成可执行的 runner、两种执行引擎(Pregel/DAG)、节点间用 channel 传值、一步步调度,以及 State 状态。这是 Eino 引擎最核心的一天。

📍 你在整门课的位置 · 第 2 周 编排引擎(共 4 周 · 20 天)
D6 流式范式 D7 Chain 链 D8 Graph 图 D9 编译执行 D10 Workflow
L01

为什么要"编译"这一步

🤔 痛点:图搭好了,为什么不能直接跑,非要"编译"一下? 你 AddNode/AddEdge 出来的只是一堆"谁连谁"的记录(一张静态描述),里面没人检查过类型对不对、执行顺序怎么排、数据从哪个管道流。要是直接边跑边算这些,每次调用都重复算一遍,慢,而且错误要等跑到那一步才暴露。
💡 本质:编译 = 把"描述图"一次性转成"可高速执行的 runner",检查全部前移 Compile(ctx) 做一次性的重活:选引擎、类型检查、建 channel、算调度顺序,最后包成 Day 06 的 Runnable之后你 Invoke/Stream 多少遍,都直接跑这个 runner,不再重复检查。"编译期做检查和准备、运行期只管快跑"是本讲贯穿的思想。

AddNode/AddEdge 得到的只是一张"描述图"(谁连谁的静态结构),它不能直接跑Compile(ctx) 把这张描述图转换成一个真正能执行的 runner,并顺手做一堆检查。

🍳 用「后厨」来理解本讲(后面统一用这个类比) 建图就像写菜谱(先放油、再下葱、然后……),编译就像开工前的备菜与验菜:检查食材齐不齐、步骤顺序合不合理,生成一套"可执行的做菜流程"。编译只做一次,之后能反复高速出菜(Invoke/Stream 很多遍)。
类型检查就像备菜时验食材:这道菜要牛肉,冰箱只有鸡肉——当场喊停,而不是炒到一半才发现(见 L03)。
channel 就像后厨的传菜台:上一个灶台做好放上去,下一个灶台从上面取(见 L05)。
super-step 就像"一轮一轮出菜":这一轮所有灶台同时开火,出完再看下一轮做什么(见 L06)。
L02

Compile 五步(真实代码)

先建立三阶段的全局观——一张图的一生分三段,检查和准备全挤在中间那段:

① 声明期 AddNode / AddEdge AddBranch 只是登记"谁连谁" 不能跑 ② 编译期 Compile 选引擎·类型检查 建 channel·算调度 包成 Runnable 重活只做一次 ⚙️ ③ 运行期 Invoke / Stream 按拓扑一轮轮跑 可反复调、只管快跑 🚀 ↻ 编译一次,运行无数次
一张图的三阶段:声明期只登记 → 编译期做完所有检查与准备 → 运行期只管高速跑

generic_graph.go:158Compile 大致做五件事:

1

选引擎:根据图类型选 Pregel(Graph,可循环)或 DAG(Workflow,不可循环)。

2

类型检查:检查每条边"上游输出类型"能否喂给"下游输入类型"(L03)。

3

建 channel:为节点间的数据传递准备"管道"(L05)。

4

算调度顺序:确定谁先跑谁后跑(拓扑),检测有无非法环。

5

包成 Runnable:生成 composableRunnable → 转成 Day 06 的 4 范式 Runnable 返回。

编译产物是 *runnergraph_run.go),它实现了 Invoke/Stream/Collect/Transform 四个方法——所以编译完你直接拿到一个 Day 06 讲的完整 Runnable。"编译期做检查和准备、运行期只管跑"是贯穿全篇的思想。
L03

类型检查前移:编译期就报错

Eino 是强类型的。每条边编译时都检查上游输出类型 → 下游输入类型能否对上(graph.go 的类型对齐逻辑)。对不上就在 Compile 时报错,而不是等运行到一半才崩。

为什么这很重要? 想象一个 Agent 图,跑到第 8 步才发现"这个节点输出 string,但下游要 []Message,类型不匹配"——运行时才崩,浪费了前 7 步的 token 和时间,还可能在生产环境半夜炸。Eino 把这类错误提前到编译期:你 Compile 的那一刻就报错,根本跑不起来。这就是 Day 01 讲的"fail-closed / 错误前移"哲学在引擎层的体现。
类型对齐还有点智能:如果上游输出的具体类型 实现了下游要的接口,也算匹配(Go 的接口赋值);还支持一些自动的 map 字段映射(Workflow 用)。检查不了的边缘情况会留一个运行时兜底检查。
L04

Pregel vs DAG:两种执行引擎

Eino 有两套执行引擎(pregel.go / dag.go),Compile 时根据你用的 API 自动选:

Pregel(Graph 默认)

  • 触发模式:AnyPredecessor——任一前驱到达就触发
  • 允许环(循环)
  • 按"一轮一轮"(super-step)推进
  • 适合 Agent 的模型↔工具循环

DAG(Workflow 默认)

  • 触发模式:AllPredecessor——所有前驱到齐才触发
  • 不允许环(有向无环图)❌
  • 按依赖拓扑并行推进
  • 适合精细的字段级数据流水线
名字来历:Pregel 是 Google 的图计算模型,特点就是"一轮轮迭代、允许回头",天然适合循环。关键差别一句话:Pregel"任一前驱到就跑"(所以能循环),DAG"等所有前驱到齐才跑"(所以适合汇聚型任务但不能有环)。Agent 用 Pregel,Workflow 用 DAG。Day 08 那个"回头边"能成立,正是因为 Pregel 的 AnyPredecessor。
L05

节点之间怎么传值:channel

每个节点有个"收件箱"叫 channelchannel.go,别和 Go 原生 chan 混——这是 Eino 自己的数据容器)。上游节点算完,把结果写进下游节点的 channel;下游节点执行前,从自己的 channel 读齐输入

// 概念示意
type channel interface {
    update(values []any) error   // 上游把值写进来(可能多个上游)
    get(...) (any, error)        // 下游取值当输入
}
为什么要 channel 这个中间层? 因为一个节点可能有多个上游(比如汇聚节点,L06 的并行合并)。channel 就是它的"收件箱":所有上游把各自结果投进来,节点执行时统一取。有多个值时,channel 负责按规则合并(比如流式场景把多个流合并)。这层抽象让"多入一出""并行汇聚"变得统一简单。
📝 最小例子:汇聚节点的 channel 收了两个上游 节点 D 有两个上游 B、C(并行跑)。
① B 算完 update(["天气数据"]) 投进 D 的 channel;② C 算完 update(["新闻数据"]) 也投进来;③ 引擎确认 D 的所有上游都到齐了,D 执行前 get() 一次拿到合并后的输入 ["天气数据","新闻数据"]
D 完全不用关心"两个上游谁先谁后到"——channel 这个收件箱替它兜住并合并了。
L06

一步一步跑:super-step 调度

Pregel 引擎按"轮"(super-step)推进(graph_run.go 的运行循环):

1

把输入写进 START 的下游 channel,得到"本轮要跑的节点集合"。

2

本轮的节点并发执行(各自从 channel 取输入、算、把结果写进下游 channel)。

3

根据"谁的 channel 被写了",算出下一轮要跑哪些节点。

4

回到第 2 步,直到 END 的 channel 被写入(拿到最终结果)或超过 MaxRunSteps

循环就是在这里发生的:Agent 图里,某轮"tools"节点执行完,把结果写进"chat"的 channel(回头边),于是下一轮"chat"又被安排执行——一轮轮转,直到某轮分支走了 END。WithMaxRunSteps(n) 限制最多多少轮,超了报 ErrExceedMaxSteps 防死循环。本轮多个节点是并发跑的(goroutine),这就是 Graph 局部并行的底层。
🎫 第一人称:假如「你」是那个输入值,被引擎一轮轮往前推 第 0 轮,我("北京天气?")被写进 START 下游的 channel → 引擎说"本轮该 chat 跑" → chat 从 channel 取到我,算完把带 tool_calls 的结果写进"分支"的 channel → 引擎算出下一轮该 tools 跑 → tools 取走、查完、把结果写回 chat 的 channel(回头边)→ 引擎又安排 chat 跑(第二次)→ chat 这次输出无 tool_calls,分支把结果写进 END 的 channel → 引擎看到 END 的 channel 有值了,把我作为最终结果吐出来。我在传菜台上被取放了好几次——每一次"取放"就是一轮 super-step。

把上面的第一人称走查,落成一张单步走查表(跟踪 "北京天气?" 从头到尾):

super-step本轮跑的节点此刻发生什么写进了谁的 channel
0(START)输入写入 START 下游chat
1chat模型判断"要查工具",输出 tool_calls分支 → tools
2tools查天气,得 "北京晴 25℃"chat(回头边)
3chat(第二次)模型输出最终答案,无 tool_calls分支 → END
4(END)END 的 channel 有值,运行结束— 返回结果
⚠️ 小白常误以为:一轮 super-step 只跑一个节点。其实:一轮里"所有当前该跑的节点"是并发跑的(多个 goroutine 一起)——只是 Agent 那个例子每轮恰好只有一个节点该跑,才看起来像一个个来。有并行分支时,一轮能同时跑好几个。
L07

State:跨节点共享的"便签本"

数据顺着边流,但有时你想要一个所有节点都能读写的共享状态(比如累计计数、缓存、整轮对话历史)。这就是 State——编译时用 WithGenLocalState 提供一个初始 state:

r, _ := g.Compile(ctx,
    compose.WithGenLocalState(func(ctx) *MyState { return &MyState{} }))
// 节点里通过 ProcessState 安全读写:
compose.ProcessState(ctx, func(ctx, s *MyState) error {
    s.Count++          // 加锁保护,并发安全
    return nil
})
State vs 数据流,怎么选? 数据流(边):一个节点的输出正好是下一个的输入——用边传,清晰。State(便签本):需要"跨越多个节点、贯穿整个运行"的共享信息——用 State。比如"记录这次对话已经调用了几次工具""缓存第一步查到的用户档案供后面所有节点用"。State 每次运行独立(local state),并发读写自动加锁。别滥用 State——能用数据流表达的就用边,State 是给真正的全局共享用的。
L08

今日小结 + 动手

🧠 今天你应该能回答

  • 为什么要"编译"?(检查+准备放编译期,运行期只管快跑)
  • Compile 五步?(选引擎/类型检查/建channel/算调度/包Runnable)
  • Pregel 和 DAG 的关键差别?(Any vs All 前驱触发;能否有环)
  • channel 是干嘛的?(节点的收件箱,处理多上游合并)
  • super-step 怎么实现循环?State 什么时候用?

✋ 动手:对着真实代码读一遍

# 1. Compile 入口(L02)
sed -n '158,230p' compose/generic_graph.go | head -50

# 2. 两种引擎
ls compose/ | grep -E 'pregel|dag'
sed -n '1,60p' compose/pregel.go

# 3. channel(L05)
sed -n '1,80p' compose/channel.go | head -50

# 4. 运行循环 super-step(L06)
grep -n 'MaxRunSteps\|ErrExceedMaxSteps\|super' compose/*.go
明天预告 · Day 10:第三种编排 Workflow——按字段连线(A 的 output.Name → B 的 input.UserName),DAG 引擎不可循环,最精细。Day 10 讲字段映射 FieldMapping 的真实代码,收官第 2 周编排引擎。
← Day 08 Graph Day 10 · Workflow →