系列「企业级 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.submitwait:同步 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;全部 ReadySkipped 才合并值返回。

条件分支跳过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 可以带 preProcessorpostProcessor

type chanCall struct {
    action        *composableRunnable  // 节点本体
    preProcessor  *composableRunnable  // 输入预处理(submit 前执行)
    postProcessor *composableRunnable  // 输出后处理(waitOne 后执行)
}
  • preProcessor:在 submit 阶段执行,可以修改节点输入(比如 FieldMapping 转换)
  • postProcessor:在 waitOne 后执行(runPostHandler),可以修改节点输出

这两层 Hook 是 Graph API 里 AddInputConverterAddOutputConverter 等功能的底层实现。


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 做的事:

  1. 把当前 channel 状态序列化进 CheckPoint
  2. 记录"下次从哪些节点恢复"(RerunNodes
  3. 返回 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

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐