Graph 执行引擎:DAG Channel、TaskManager 与 CheckPoint 源码拆解(第62篇-E48)
系列「企业级 AI Agent 实现拆解」E48 篇,Part 10 生产工程篇第六章。上一篇 讲了 StreamReader 的流式传输。这篇往上一层:Eino Graph 是怎么调度节点执行的——Channel 如何追踪依赖、TaskManager 如何并发执行、CheckPoint 如何实现 HITL 中断恢复。
读完这篇你会知道
- Graph 执行的"主循环"长什么样:submit → wait → calculateNext,循环直到 END
- DAG Channel 和 Pregel Channel 的核心区别
dagChannel如何用三态(waiting/ready/skipped)追踪依赖channelManager.updateAndGet:一次 channel 更新 + 就绪查询的完整流程taskManager.submit和wait:同步 vs 异步执行的判断逻辑- preProcessor / postProcessor:节点级别的输入输出变换
- HITL 中断和 CheckPoint:Graph 如何在任意节点前后暂停并恢复
主循环:Graph 运行的骨架
所有 Graph 执行最终走到 runner.run,核心是一个三步主循环:
while true:
1. submit(nextTasks) // 把待执行节点提交给 taskManager
2. completedTasks = wait() // 等待(至少一个)节点完成
3. nextTasks = calculateNextTasks(completedTasks) // 算出下一轮任务
if isEnd: return result
// compose/graph_run.go(简化)
for step := 0; ; step++ {
if !r.dag && step >= maxSteps {
return nil, ErrExceedMaxSteps // Pregel 防死循环
}
tm.submit(nextTasks)
completedTasks, canceled, _ := tm.wait()
nextTasks, result, isEnd, err = r.calculateNextTasks(ctx, completedTasks, isStream, cm, optMap)
if isEnd { return result, nil }
}
DAG 模式不限步数(r.dag == true),因为拓扑排序保证无环。Pregel 模式(有环图,用于 ReAct 循环)有 maxSteps 上限,防止 Agent 无限循环。
DAG Channel vs Pregel Channel
Eino 有两种 channel 类型,通过 chanBuilder 函数指针在编译时确定:
// channel 接口
type channel interface {
reportValues(map[string]any) error // 上游节点写入输出
reportDependencies([]string) // 上游节点标记"我完成了"
reportSkip([]string) bool // 条件分支跳过某节点
get(bool, string, *edgeHandlerManager) (any, bool, error) // 检查并取值
convertValues(fn func(map[string]any) error) error
load(channel) error // CheckPoint 恢复用
}
dagChannel:静态依赖 + 三态追踪
DAG Channel 负责追踪哪些前驱已完成,只有所有前驱就绪,节点才能执行。
type dagChannel struct {
ControlPredecessors map[string]dependencyState // 哪些前驱节点必须完成
DataPredecessors map[string]bool // FieldMapping 的间接依赖
Values map[string]any // 已收到的值
Skipped bool // 整个节点是否被跳过
}
type dependencyState uint8
const (
dependencyStateWaiting dependencyState = iota // 0:还在等
dependencyStateReady // 1:已完成
dependencyStateSkipped // 2:被条件分支跳过
)
get 方法的逻辑:扫描所有前驱状态,只要还有一个 Waiting,返回 ready=false;全部 Ready 或 Skipped 才合并值返回。
条件分支跳过(reportSkip)也会级联传播:
func (ch *dagChannel) reportSkip(keys []string) bool {
// 把传入的 key 标为 Skipped
// 如果 ControlPredecessors 全部是 Skipped → ch.Skipped = true
// 返回 true 表示这个节点也被跳过了
}
channelManager.reportBranch 调用 reportSkip 后,会递归把下游节点也标为 Skipped——一个分支跳过,整条跳过路径都不会执行。
pregelChannel:无静态依赖,随时就绪
Pregel Channel 没有 ControlPredecessors,任何时候 Values 非空就算就绪:
func (ch *pregelChannel) get(isStream bool, name string, edgeHandler *edgeHandlerManager) (any, bool, error) {
if len(ch.Values) == 0 {
return nil, false, nil // 没有值 → 不就绪
}
// 有值就合并返回,消费掉
defer func() { ch.Values = map[string]any{} }()
// ...
}
func (ch *pregelChannel) reportSkip(_ []string) bool { return false } // 永不跳过
func (ch *pregelChannel) reportDependencies(_ []string) { return } // 不追踪
ReAct 循环的 Agent 用 Pregel 模式——LLM 每轮输出之后,不管前面是否有"未完成的前驱",直接根据 LLM 的输出决定下一步。
channelManager:数据总线
channelManager 管理所有节点的 channel,是 Graph 的"数据总线":
type channelManager struct {
isStream bool
channels map[string]channel // 每个节点一个 channel
successors map[string][]string
dataPredecessors map[string]map[string]struct{}
controlPredecessors map[string]map[string]struct{}
edgeHandlerManager *edgeHandlerManager // 边上的转换函数
preNodeHandlerManager *preNodeHandlerManager // 节点前的转换函数
}
每轮循环里最关键的调用是 updateAndGet:
func (c *channelManager) updateAndGet(ctx, values, dependencies) (map[string]any, error) {
// 1. 把各节点的输出写进对应的 channel
c.updateValues(ctx, values)
// 2. 更新控制依赖(谁完成了)
c.updateDependencies(ctx, dependencies)
// 3. 扫描所有 channel,找出就绪的节点,返回它们的输入值
return c.getFromReadyChannels(ctx)
}
边转换(edgeHandler):两个节点之间的连边可以带转换函数,比如 FieldMapping——从上游节点的输出里提取特定字段再送给下游。这个转换在 channel.get 内部调用 edgeHandlerManager.handle(from, to, value) 执行。
fan-in 合并:当一个节点有多个上游,channel 收到多份值时,mergeValues 把它们合并(invoke 模式合并结构体字段,stream 模式用 MergeStreamReaders)。
taskManager:并发调度
type taskManager struct {
runWrapper runnableCallWrapper // invoke 或 transform(流式)
num uint32 // 当前运行中的任务数
done *internal.UnboundedChan[*task] // 完成通知 chan
runningTasks map[string]*task
cancelCh chan *time.Duration // HITL 中断信号
canceled bool
deadline *time.Time
}
submit:决定同步还是异步
func (t *taskManager) submit(tasks []*task) error {
// ...preProcessor 处理...
var syncTask *task
if t.num == 0 && (len(tasks) == 1 || t.needAll) && t.cancelCh == nil {
// 以下条件满足时,同步执行(不启 goroutine):
// 1. 当前没有其他在跑的任务(t.num == 0)
// 2. 只有一个新任务,或者 needAll 模式
// 3. 没有 HITL 中断通道(可中断的 Graph 不能同步跑)
syncTask = tasks[0]
tasks = tasks[1:]
}
for _, task := range tasks {
t.num += 1
go t.execute(task) // 异步执行
}
if syncTask != nil {
t.num += 1
t.execute(syncTask) // 同步执行(当前 goroutine,省切换开销)
}
}
这个优化很重要:线性 Graph(每轮只有一个节点可跑)不开 goroutine,避免不必要的调度开销。只有真正有并行节点时才异步。
execute:节点执行
func (t *taskManager) execute(currentTask *task) {
defer func() {
if panicInfo := recover(); panicInfo != nil {
currentTask.err = safe.NewPanicErr(panicInfo, debug.Stack())
}
t.done.Send(currentTask) // 无论成功失败都通知
}()
// 注入 RunInfo + Callback Handler(E45 讲过)
ctx := initNodeCallbacks(currentTask.ctx, currentTask.nodeKey, ...)
// 执行节点(invoke 或 transform)
currentTask.output, currentTask.err = t.runWrapper(ctx, currentTask.call.action, currentTask.input, ...)
}
节点 panic 会被 recover 捕获,包装成 PanicErr 返回给上层——Graph 不会因为单个节点 panic 而整体崩溃。
wait:等待完成 + 处理取消
func (t *taskManager) wait() (tasks []*task, canceled bool, canceledTasks []*task) {
if t.needAll {
// needAll 模式:等所有并行节点都完成再继续(如并行 Tool 调用)
tasks, canceledTasks = t.waitAll()
return tasks, t.canceled, canceledTasks
}
// 默认:等第一个完成就处理(流水线模式,尽早启动下游)
ta, success, canceled := t.waitOne()
// ...
}
needAll vs 默认:默认模式拿到一个完成就继续计算下一步,适合流水线。needAll 等所有节点完成,适合并行工具调用(必须拿到所有工具结果才能让 LLM 继续)。
preProcessor / postProcessor
每个节点的 chanCall 可以带 preProcessor 和 postProcessor:
type chanCall struct {
action *composableRunnable // 节点本体
preProcessor *composableRunnable // 输入预处理(submit 前执行)
postProcessor *composableRunnable // 输出后处理(waitOne 后执行)
}
- preProcessor:在
submit阶段执行,可以修改节点输入(比如 FieldMapping 转换) - postProcessor:在
waitOne后执行(runPostHandler),可以修改节点输出
这两层 Hook 是 Graph API 里 AddInputConverter、AddOutputConverter 等功能的底层实现。
HITL 中断:在任意节点前后暂停
interruptBeforeNodes / interruptAfterNodes 让 Graph 在特定节点前后暂停,等待人工确认:
// 配置:在某节点执行前暂停
graph.AddEdge("agent", "human_review",
compose.WithInterruptBefore("human_review"))
// 运行时:Graph 检测到命中 interruptBeforeNodes
if keys := getHitKey(nextTasks, r.interruptBeforeNodes); len(keys) > 0 {
return nil, r.handleInterrupt(ctx, tempInfo, nextTasks, ...)
}
handleInterrupt 做的事:
- 把当前 channel 状态序列化进 CheckPoint
- 记录"下次从哪些节点恢复"(
RerunNodes) - 返回
InterruptError(不是真正的错误,是暂停信号)
恢复:用 CheckPoint ID 重新调用 Graph,主循环在开头检测到 CheckPoint 存在,直接 restoreCheckPointState 恢复 channel 状态,跳过重新初始化:
if cp := getCheckPointFromCtx(ctx); cp != nil {
// 恢复 channel 状态
ctx, err = r.restoreCheckPointState(ctx, path, stateModifier, cp, isStream, cm)
// 从中断点继续
nextTasks, err = r.restoreTasks(ctx, cp.Inputs, cp.SkipPreHandler, cp.RerunNodes, ...)
}
整个循环继续,就像从未中断过。
流式模式的差异
run(ctx, isStream=true, ...) 和非流式的主循环结构相同,区别在 runWrapper:
// 非流式:runWrapper = runnableInvoke
func runnableInvoke(ctx, r, input, opts) (any, error) {
return r.i(ctx, input, opts...) // 调用 .Invoke
}
// 流式:runWrapper = runnableTransform
func runnableTransform(ctx, r, input, opts) (any, error) {
return r.t(ctx, input.(streamReader), opts...) // 调用 .Transform
}
流式模式下,节点的输入输出都是 streamReader。channel 合并多个上游时,用 MergeStreamReaders(E47 讲过的 fan-in)而不是 struct merge。
小结
Eino Graph 执行引擎的核心是channel-based 依赖追踪 + goroutine 池并发:
| 组件 | 职责 |
|---|---|
dagChannel |
静态依赖三态追踪(waiting/ready/skipped),支持条件分支跳过 |
pregelChannel |
无静态依赖,有值即就绪,用于有环图(ReAct 循环) |
channelManager.updateAndGet |
数据总线:写入输出 + 更新依赖 + 返回就绪节点 |
taskManager.submit |
智能调度:线性路径同步执行,并行节点开 goroutine |
taskManager.wait |
默认等第一个(流水线);needAll 等全部(并行工具调用) |
| preProcessor / postProcessor | 节点级输入输出变换,AddInputConverter 的底层 |
| HITL + CheckPoint | 在任意节点前后暂停,序列化 channel 状态,可恢复继续 |
理解了这套调度机制,就能理解为什么 Eino Graph 可以同时支持流水线式的 RAG 链路、有环的 ReAct Agent 循环,以及需要人工确认的 HITL 工作流——底层是同一个主循环,靠 channel 类型和 taskManager 配置区分。
代码来源:eino/compose/graph_run.go · eino/compose/graph_manager.go · eino/compose/dag.go · eino/compose/pregel.go
更多推荐



所有评论(0)