系列「企业级 AI Agent 实现拆解」E49 篇,Part 10 生产工程篇第七章。上一篇 讲了 Graph 如何在运行时调度节点。这篇往前一步:graph.Compile() 被调用的那一刻,Eino 在做什么——类型校验、DAG/Pregel 选择、节点编译、环检测,直到交出一个可运行的 composableRunnable

读完这篇你会知道

  • graph.Compile() 的 10 个阶段分别做什么
  • DAG 模式 vs Pregel 模式是怎么在编译期确定的
  • addNode / addEdge 阶段就在做类型推断,不是等 Compile 才检查
  • validateDAG 用 Kahn 算法检测环的实现
  • toGraphInfo 是给 devops callback 用的——Graph 结构的快照
  • 为什么 Compile 之后就不能再改 Graph

Compile 的整体时序

一个 Eino Graph 的生命周期分两个阶段:构建期(AddXxxNode + AddEdge)和编译期(Compile)。

[构建期]
graph.AddLLMNode("agent", llm)
graph.AddToolsNode("tools", toolCallers)
graph.AddEdge(START, "agent")
graph.AddEdge("agent", "tools")
graph.AddEdge("tools", END)

[编译期]
runnable, err := graph.Compile(ctx, opts...)
// 之后 g.compiled = true,再调用 AddXxxNode 返回 ErrGraphCompiled

compile() 是一次性的转化:把"图描述"变成"可运行的执行器"。


阶段 0:快速失败

func (g *graph) compile(ctx context.Context, opt *graphCompileOptions) (*composableRunnable, error) {
    if g.buildError != nil {
        return nil, g.buildError  // 构建期就出错了,直接失败
    }
    // ...
}

buildErroraddNode / addEdge 设置的——每次构建期出错都会把错误存在这里。Compile 开头就检查,省得后续步骤白跑。


阶段 1:决定运行模式(DAG vs Pregel)

// 默认 Pregel(支持有环图)
runType := runTypePregel
cb := pregelChannelBuilder

// 满足以下条件之一 → 切换到 DAG
if (opt != nil && opt.nodeTriggerMode == AllPredecessor) || isWorkflow(g.cmp) {
    runType = runTypeDAG
    cb = dagChannelBuilder
}

这里 cb 是一个函数指针,决定后续 channel 实例化用哪种 builder:

条件 运行模式 chanBuilder
Graph(默认)/ Chain Pregel pregelChannelBuilder
nodeTriggerMode == AllPredecessor DAG dagChannelBuilder
Workflow(ComponentOfWorkflow DAG dagChannelBuilder

Chain 和 Workflow 不允许自定义 nodeTriggerMode,传了会直接报错:

if isChain(g.cmp) || isWorkflow(g.cmp) {
    if opt != nil && opt.nodeTriggerMode != "" {
        return nil, errors.New(fmt.Sprintf("%s doesn't support node trigger mode option", g.cmp))
    }
}

阶段 2:确定 Eager 模式

eager := false
if isWorkflow(g.cmp) || runType == runTypeDAG {
    eager = true
}
if opt != nil && opt.eagerDisabled {
    eager = false
}

Eager 模式的含义:节点就绪后立刻把输入预加载出来,不等最后一刻。Workflow 和 DAG 默认开启,可以被 eagerDisabled 关掉。这是一个运行时的调度细节,影响节点输入的准备时机。


阶段 3:基础完整性校验

if len(g.startNodes) == 0 {
    return nil, errors.New("start node not set")
}
if len(g.endNodes) == 0 {
    return nil, errors.New("end node not set")
}

没有连 START 或 END 的图,Compile 直接拒绝。

类型推断校验

// toValidateMap 不为空 = 有节点的类型无法推断
for _, v := range g.toValidateMap {
    if len(v) > 0 {
        return nil, fmt.Errorf("some node's input or output types cannot be inferred: %v", g.toValidateMap)
    }
}

toValidateMapaddEdge 时维护的——每条边连接时,都会尝试推断两端节点的输入/输出类型。如果到 Compile 时还有无法推断的类型,报错。

FieldMapping 重复校验

for key := range g.fieldMappingRecords {
    toMap := make(map[string]bool)
    for _, mapping := range g.fieldMappingRecords[key] {
        if _, ok := toMap[mapping.to]; ok {
            return nil, fmt.Errorf("duplicate mapping target field: %s of node[%s]", mapping.to, key)
        }
        toMap[mapping.to] = true
    }
    // 校验通过 → 注册 FieldMapping 转换器为 preNode handler
    g.handlerPreNode[key] = append(g.handlerPreNode[key], g.getNodeGenericHelper(key).inputFieldMappingConverter)
}

同一个目标字段不能被两条 FieldMapping 同时映射。校验通过后,转换器被注册为 preNode handler(E48 里讲过的 preProcessor)。


阶段 4:递归编译子图

key2SubGraphs := g.beforeChildGraphsCompile(opt)
chanSubscribeTo := make(map[string]*chanCall)

for name, node := range g.nodes {
    node.beforeChildGraphCompile(name, key2SubGraphs)

    r, err := node.compileIfNeeded(ctx)  // 子图节点在这里递归编译
    if err != nil {
        return nil, err
    }

    chCall := &chanCall{
        action:        r,
        writeTo:       g.dataEdges[name],
        controls:      g.controlEdges[name],
        preProcessor:  node.nodeInfo.preProcessor,
        postProcessor: node.nodeInfo.postProcessor,
    }
    // 挂上条件分支
    branches := g.branches[name]
    if len(branches) > 0 {
        branchRuns := make([]*GraphBranch, 0, len(branches))
        branchRuns = append(branchRuns, branches...)
        chCall.writeToBranches = branchRuns
    }

    chanSubscribeTo[name] = chCall
}

Graph 支持嵌套(一个 Graph 的某个节点本身也是一个 Graph)。compileIfNeeded 会递归触发子图的 compile()

每个节点编译完之后,都会被包装成 chanCall——这就是 E48 里 taskManager 执行的基本单元:

  • action:节点本体(composableRunnable
  • writeTo:输出数据流向哪些节点
  • controls:控制流流向哪些节点
  • preProcessor / postProcessor:前后处理器

阶段 5:构建前驱映射

dataPredecessors    := make(map[string][]string)
controlPredecessors := make(map[string][]string)

// 反转 controlEdges:start → [end1, end2] 变成 end → [start1, start2]
for start, ends := range g.controlEdges {
    for _, end := range ends {
        controlPredecessors[end] = append(controlPredecessors[end], start)
    }
}
// 同样反转 dataEdges
for start, ends := range g.dataEdges {
    for _, end := range ends {
        dataPredecessors[end] = append(dataPredecessors[end], start)
    }
}
// 分支也要处理
for start, branches := range g.branches {
    for _, branch := range branches {
        for end := range branch.endNodes {
            controlPredecessors[end] = append(controlPredecessors[end], start)
            if !branch.noDataFlow {
                dataPredecessors[end] = append(dataPredecessors[end], start)
            }
        }
    }
}

原始的 controlEdgesdataEdges 是"从上游看下游"的方向。这里把它们翻转成"从下游看上游"的前驱映射,供 dagChannel 的依赖追踪使用(E48 的 ControlPredecessors)。


阶段 6:构造 runner

r := &runner{
    chanSubscribeTo:     chanSubscribeTo,       // 节点 → chanCall 映射
    controlPredecessors: controlPredecessors,
    dataPredecessors:    dataPredecessors,
    inputChannels:       inputChannels,          // START 节点的出边
    eager:               eager,
    chanBuilder:         cb,                     // DAG 或 Pregel builder
    inputType:           g.inputType(),
    outputType:          g.outputType(),
    genericHelper:       g.genericHelper,

    preBranchHandlerManager: ...,
    preNodeHandlerManager:   ...,
    edgeHandlerManager:      ...,
    mergeConfigs:            mergeConfigs,
}

// 构建 successors 映射(每个节点的所有下游)
successors := make(map[string][]string)
for ch := range r.chanSubscribeTo {
    successors[ch] = getSuccessors(r.chanSubscribeTo[ch])
}
r.successors = successors

runner 是 Graph 运行时的核心,它持有所有执行所需的数据。chanBuilder 是在这里绑定的——运行时每次主循环创建新 channel 实例时,调用的就是这个 builder。

如果 Graph 有 State(stateGenerator != nil),还会设置 r.runCtx——每次 run 调用时自动把 State 注入 Context:

if g.stateGenerator != nil {
    r.runCtx = func(ctx context.Context) context.Context {
        var parent *internalState
        if p, ok := ctx.Value(stateKey{}).(*internalState); ok {
            parent = p
        }
        return context.WithValue(ctx, stateKey{}, &internalState{
            state:  g.stateGenerator(ctx),
            parent: parent,
        })
    }
}

阶段 7:DAG 环检测(Kahn 算法)

只有 DAG 模式才做这步——Pregel 模式允许有环(ReAct 循环就是环),所以不检查。

if runType == runTypeDAG {
    err := validateDAG(r.chanSubscribeTo, controlPredecessors)
    if err != nil {
        return nil, err
    }
    r.dag = true
}

validateDAGKahn 算法(拓扑排序的 in-degree 版本):

func validateDAG(chanSubscribeTo map[string]*chanCall, controlPredecessors map[string][]string) error {
    // 1. 初始化每个节点的"等待中的前驱数量"
    m := map[string]int{}
    for node := range chanSubscribeTo {
        if edges, ok := controlPredecessors[node]; ok {
            m[node] = len(edges)
            for _, pre := range edges {
                if pre == START {
                    m[node] -= 1  // START 不算真实前驱(永远就绪)
                }
            }
        } else {
            m[node] = 0
        }
    }

    // 2. 反复扫描:找到入度为 0 的节点,删除它(把下游的入度都减 1)
    hasChanged := true
    for hasChanged {
        hasChanged = false
        for node := range m {
            if m[node] == 0 {
                hasChanged = true
                for _, subNode := range chanSubscribeTo[node].controls {
                    if subNode == END { continue }
                    m[subNode]--
                }
                for _, subBranch := range chanSubscribeTo[node].writeToBranches {
                    for subNode := range subBranch.endNodes {
                        if subNode == END { continue }
                        m[subNode]--
                    }
                }
                m[node] = -1  // 标记为"已处理"
            }
        }
    }

    // 3. 还有 count > 0 的节点 → 它们在环里
    var loopStarts []string
    for k, v := range m {
        if v > 0 {
            loopStarts = append(loopStarts, k)
        }
    }
    if len(loopStarts) > 0 {
        return fmt.Errorf("%w: %s", DAGInvalidLoopErr, formatLoops(findLoops(loopStarts, chanSubscribeTo)))
    }
    return nil
}

为什么是 Kahn 而不是 DFS:Kahn 天然给出所有参与环的节点列表(count > 0 的那些),报错信息可以包含完整的环路径,比 DFS 的"第一个发现的环节点"更有用。


阶段 8:CheckPoint 和 HITL 配置

if opt != nil {
    // 为每个节点收集流转换对(用于序列化/反序列化 CheckPoint)
    inputPairs  := make(map[string]streamConvertPair)
    outputPairs := make(map[string]streamConvertPair)
    for key, c := range r.chanSubscribeTo {
        inputPairs[key]  = c.action.inputStreamConvertPair
        outputPairs[key] = c.action.outputStreamConvertPair
    }
    r.checkPointer = newCheckPointer(inputPairs, outputPairs, opt.checkPointStore, opt.serializer)

    r.interruptBeforeNodes = opt.interruptBeforeNodes
    r.interruptAfterNodes  = opt.interruptAfterNodes
}

checkPointer 需要知道每个节点的输入/输出类型信息(用于序列化),所以必须在节点都编译完之后才能创建。


阶段 9:MaxRunSteps 默认值

if r.dag && r.options.maxRunSteps > 0 {
    return nil, fmt.Errorf("cannot set max run steps in dag mode")
} else if !r.dag && r.options.maxRunSteps == 0 {
    r.options.maxRunSteps = len(r.chanSubscribeTo) + 10  // Pregel 默认步数上限
}

DAG 模式有限步保证(拓扑排序),不需要 maxRunSteps。Pregel 模式(有环)必须有上限,默认是节点数 + 10,防止 Agent 死循环。


阶段 10:onCompileFinish 与 GraphInfo

g.compiled = true
g.onCompileFinish(ctx, opt, key2SubGraphs)
return r.toComposableRunnable(), nil

onCompileFinish 调用 toGraphInfo 生成 GraphInfo 快照,然后调用所有注册的 compile callbacks:

type GraphInfo struct {
    Nodes         map[string]GraphNodeInfo  // 每个节点的元信息
    Edges         map[string][]string        // 控制边
    DataEdges     map[string][]string        // 数据边
    Branches      map[string][]GraphBranch   // 条件分支
    InputType     reflect.Type
    OutputType    reflect.Type
    Name          string
    GenStateFn    func(ctx context.Context) any
    CompileOptions any
}

这个 GraphInfo 是给框架上层(比如 DevOps 面板、可观测性系统)用的——Graph 编译完成后,外部系统可以拿到 Graph 的完整结构图来可视化或记录。

子图(嵌套 Graph)的 GraphInfo 通过 key2SubGraphs 被嵌入进父节点的 gNodeInfo.GraphInfo,所以整棵 Graph 树的结构都在一次 callback 里拿得到。


addEdge 期间就在做类型检查

一个容易忽略的点:类型推断不是等 Compile 才做的,addEdge 调用时就在做:

func (g *graph) addEdgeWithMappings(startNode, endNode string, ...) error {
    // ...
    g.addToValidateMap(startNode, endNode, mappings)
    err = g.updateToValidateMap()  // 立刻尝试推断类型
    if err != nil {
        return err
    }
    g.dataEdges[startNode] = append(g.dataEdges[startNode], endNode)
    return nil
}

updateToValidateMap 会尝试匹配上下游节点的 I/O 类型。如果类型不匹配,addEdge 就报错——不用等到 Compile。只有那些"暂时还推断不出来"的类型才留到 toValidateMap 里,等 Compile 阶段再检查。

这意味着大部分类型错误在构建期就被抓到了,开发体验更好。


整体流程图

graph.Compile(ctx, opts...)
    │
    ├─ [0] buildError? → fail
    │
    ├─ [1] DAG vs Pregel → chanBuilder 函数指针
    │
    ├─ [2] eager 模式确定
    │
    ├─ [3] 基础校验
    │       ├─ startNodes / endNodes 不为空
    │       ├─ toValidateMap 空(类型全推断完)
    │       └─ FieldMapping 无重复目标
    │
    ├─ [4] 遍历节点
    │       ├─ 子图递归 compile
    │       └─ 包装成 chanCall{action, writeTo, controls, pre/postProcessor}
    │
    ├─ [5] 构建 dataPredecessors + controlPredecessors(反转边图)
    │
    ├─ [6] 构造 runner(chanBuilder, successors, stateGenerator…)
    │
    ├─ [7] DAG 模式 → validateDAG(Kahn 算法检环)
    │
    ├─ [8] CheckPoint + HITL 配置
    │
    ├─ [9] maxRunSteps 默认值
    │
    ├─ [10] g.compiled = true
    │         onCompileFinish → toGraphInfo → callbacks
    │
    └─ return runner.toComposableRunnable()

小结

compile() 做的事本质上是两件

  1. 验证:类型、结构、依赖——构建期能验证的在 addEdge 就做了,剩下的在 Compile 里统一验证。
  2. 组装:把图描述变成执行器——chanCall 数组 + 前驱映射 + runner

拿到 composableRunnable 之后,你就有了一个可以反复调用的执行单元,内部封装了 E48 讲过的主循环。g.compiled = true 之后任何 AddXxxNode/AddEdge 调用都返回 ErrGraphCompiled——图是不可变的,执行时不会有并发修改问题。


代码来源:eino/compose/graph.go

Logo

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

更多推荐