Graph 编译与执行
图搭好了怎么跑?今天讲 Compile 把图变成可执行的 runner、两种执行引擎(Pregel/DAG)、节点间用 channel 传值、一步步调度,以及 State 状态。这是 Eino 引擎最核心的一天。
为什么要"编译"这一步
Compile(ctx) 做一次性的重活:选引擎、类型检查、建 channel、算调度顺序,最后包成 Day 06 的 Runnable。之后你 Invoke/Stream 多少遍,都直接跑这个 runner,不再重复检查。"编译期做检查和准备、运行期只管快跑"是本讲贯穿的思想。你 AddNode/AddEdge 得到的只是一张"描述图"(谁连谁的静态结构),它不能直接跑。Compile(ctx) 把这张描述图转换成一个真正能执行的 runner,并顺手做一堆检查。
类型检查就像备菜时验食材:这道菜要牛肉,冰箱只有鸡肉——当场喊停,而不是炒到一半才发现(见 L03)。
channel 就像后厨的传菜台:上一个灶台做好放上去,下一个灶台从上面取(见 L05)。
super-step 就像"一轮一轮出菜":这一轮所有灶台同时开火,出完再看下一轮做什么(见 L06)。
Compile 五步(真实代码)
先建立三阶段的全局观——一张图的一生分三段,检查和准备全挤在中间那段:
generic_graph.go:158 的 Compile 大致做五件事:
选引擎:根据图类型选 Pregel(Graph,可循环)或 DAG(Workflow,不可循环)。
类型检查:检查每条边"上游输出类型"能否喂给"下游输入类型"(L03)。
建 channel:为节点间的数据传递准备"管道"(L05)。
算调度顺序:确定谁先跑谁后跑(拓扑),检测有无非法环。
包成 Runnable:生成 composableRunnable → 转成 Day 06 的 4 范式 Runnable 返回。
*runner(graph_run.go),它实现了 Invoke/Stream/Collect/Transform 四个方法——所以编译完你直接拿到一个 Day 06 讲的完整 Runnable。"编译期做检查和准备、运行期只管跑"是贯穿全篇的思想。类型检查前移:编译期就报错
Eino 是强类型的。每条边编译时都检查上游输出类型 → 下游输入类型能否对上(graph.go 的类型对齐逻辑)。对不上就在 Compile 时报错,而不是等运行到一半才崩。
Pregel vs DAG:两种执行引擎
Eino 有两套执行引擎(pregel.go / dag.go),Compile 时根据你用的 API 自动选:
Pregel(Graph 默认)
- 触发模式:AnyPredecessor——任一前驱到达就触发
- 允许环(循环)✅
- 按"一轮一轮"(super-step)推进
- 适合 Agent 的模型↔工具循环
DAG(Workflow 默认)
- 触发模式:AllPredecessor——所有前驱到齐才触发
- 不允许环(有向无环图)❌
- 按依赖拓扑并行推进
- 适合精细的字段级数据流水线
节点之间怎么传值:channel
每个节点有个"收件箱"叫 channel(channel.go,别和 Go 原生 chan 混——这是 Eino 自己的数据容器)。上游节点算完,把结果写进下游节点的 channel;下游节点执行前,从自己的 channel 读齐输入。
// 概念示意
type channel interface {
update(values []any) error // 上游把值写进来(可能多个上游)
get(...) (any, error) // 下游取值当输入
}
① B 算完
update(["天气数据"]) 投进 D 的 channel;② C 算完 update(["新闻数据"]) 也投进来;③ 引擎确认 D 的所有上游都到齐了,D 执行前 get() 一次拿到合并后的输入 ["天气数据","新闻数据"]。D 完全不用关心"两个上游谁先谁后到"——channel 这个收件箱替它兜住并合并了。
一步一步跑:super-step 调度
Pregel 引擎按"轮"(super-step)推进(graph_run.go 的运行循环):
把输入写进 START 的下游 channel,得到"本轮要跑的节点集合"。
本轮的节点并发执行(各自从 channel 取输入、算、把结果写进下游 channel)。
根据"谁的 channel 被写了",算出下一轮要跑哪些节点。
回到第 2 步,直到 END 的 channel 被写入(拿到最终结果)或超过 MaxRunSteps。
WithMaxRunSteps(n) 限制最多多少轮,超了报 ErrExceedMaxSteps 防死循环。本轮多个节点是并发跑的(goroutine),这就是 Graph 局部并行的底层。"北京天气?")被写进 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 |
| 1 | chat | 模型判断"要查工具",输出 tool_calls | 分支 → tools |
| 2 | tools | 查天气,得 "北京晴 25℃" | chat(回头边) |
| 3 | chat(第二次) | 模型输出最终答案,无 tool_calls | 分支 → END |
| 4 | (END) | END 的 channel 有值,运行结束 | — 返回结果 |
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
})
今日小结 + 动手
🧠 今天你应该能回答
- 为什么要"编译"?(检查+准备放编译期,运行期只管快跑)
- 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