跳转至

cloudwego/eino:字节跳动开源的 Go AI 编排框架

如果你要用 Go 写一个"调用大模型 + 按需调用工具 + 支持多路径分支判断"的 AI 应用,直接手写一堆嵌套 if/else 和函数调用很快会失控。cloudwego/eino 就是字节跳动开源的一套 Go 编排框架:你像画流程图一样,把"读 Prompt 模板 → 调大模型 → 按需调用工具 → 把结果喂回模型再判断"这些步骤定义成一张图(Graph)或一条链(Chain),框架负责按图调度执行、处理流式输出、支持中断与恢复。它要解决的核心问题是:多步骤 AI 智能体的控制流,怎么写得既灵活又不失控。

本文所有代码引用均来自 github.com/cloudwego/eino 仓库源码(clone 后逐文件核对,标注了具体文件路径),不再使用编造的 API。


0. 一个最小例子

在深入类型定义之前,先看 eino 最基础的用法:把两个步骤(拼 Prompt、调模型)串成一张图并执行。

package main

import (
    "context"
    "fmt"

    "github.com/cloudwego/eino/compose"
    "github.com/cloudwego/eino/schema"
)

func main() {
    ctx := context.Background()

    // compose/generic_graph.go: NewGraph[I, O any](opts ...NewGraphOption) *Graph[I, O]
    graph := compose.NewGraph[string, *schema.Message]()

    // Lambda 节点:一个普通函数,负责拼 Prompt
    _ = graph.AddLambdaNode("prompt", compose.InvokableLambda(
        func(ctx context.Context, in string) (string, error) {
            return fmt.Sprintf("System: Answer concisely.\nUser: %s", in), nil
        }))

    // ChatModel 节点:调用大模型,chatModel 需实现 model.BaseChatModel 接口
    _ = graph.AddChatModelNode("model", chatModel)

    _ = graph.AddEdge(compose.START, "prompt")
    _ = graph.AddEdge("prompt", "model")
    _ = graph.AddEdge("model", compose.END)

    runnable, err := graph.Compile(ctx)
    if err != nil {
        panic(err)
    }

    out, err := runnable.Invoke(ctx, "How is the weather?")
    if err != nil {
        panic(err)
    }
    fmt.Println(out.Content)
}

这几行代码已经用到了 eino 的核心机制:AddLambdaNode/AddChatModelNode 往图里加节点,AddEdge 把节点连起来,Compile 把这张图编译成一个可以直接 InvokeRunnable。下面展开讲这套拓扑到底是怎么定义和执行的,以及流式、分支、中断恢复这些进阶能力。

eino 的核心设计思路(人话版)

抛开术语,eino 要解决的问题也很朴素:

  1. 怎么把"调用模型 → 判断要不要调用工具 → 调用工具 → 把结果喂回模型"这套多步骤逻辑,用画图的方式表达出来,而不是写一堆嵌套 if/else?—— eino 用 Graph(可以带分支、甚至带环)或更简单的 Chain(纯线性)描述这套流程,你只管往图里加节点、连边,编译后拿到一个可以直接调用的 Runnable
  2. 大模型的输出往往是流式的(一个个 token 蹦出来),怎么让"流式输出"和"一次性拿完整结果"这两种用法在框架层面统一,而不是让每个组件自己维护两套代码?—— eino 定义了 Invoke/Stream/Collect/Transform 四个方法,只要实现其中一个,框架自动帮你推导出另外三个。
  3. 长时间运行的智能体中途出错,或者需要人工介入怎么办?—— eino 提供了 Callback 和 Interrupt/Resume 机制,可以在指定节点前后暂停执行、落 checkpoint,之后再从中断点恢复。

下面按"拓扑怎么定义 → 怎么流式执行 → 实际怎么写完整代码 → 踩过的坑"的顺序展开。


1. 怎么把多步骤逻辑画成图:Graph / Chain 的真实类型定义

Eino 的核心设计思想是"配置即拓扑,运行即状态":先把节点和边搭好(配置阶段),编译成一个 Runnable 之后再真正执行(运行阶段)。图的入口是 compose.NewGraph[I, O](opts ...NewGraphOption) *Graph[I, O]compose/generic_graph.go),返回的 Graph[I, O] 只是对内部 *graph 的一层泛型包装:

// compose/graph.go
type graph struct {
    nodes        map[string]*graphNode
    controlEdges map[string][]string
    dataEdges    map[string][]string
    branches     map[string][]*GraphBranch
    startNodes   []string
    endNodes     []string
    stateType      reflect.Type
    stateGenerator func(ctx context.Context) any
    ...
}

值得注意的是:Eino 没有一个万能的 AddNode 方法,而是为每一种组件接口单独提供了强类型的注册方法(compose/graph.go 304-451 行):

func (g *graph) AddChatModelNode(key string, node model.BaseChatModel, opts ...GraphAddNodeOpt) error
func (g *graph) AddChatTemplateNode(key string, node prompt.ChatTemplate, opts ...GraphAddNodeOpt) error
func (g *graph) AddToolsNode(key string, node *ToolsNode, opts ...GraphAddNodeOpt) error
func (g *graph) AddRetrieverNode(key string, node retriever.Retriever, opts ...GraphAddNodeOpt) error
func (g *graph) AddLambdaNode(key string, node *Lambda, opts ...GraphAddNodeOpt) error

也就是说,节点的输入输出类型是由组件接口(model.BaseChatModelprompt.ChatTemplate……)在编译期锁死的,而不是像 Python 框架那样先塞进一个通用容器再在运行时反射猜测类型。唯一的"自由节点"是 Lambdacompose/types_lambda.go),只能通过 compose.InvokableLambda / StreamableLambda / CollectableLambda / TransformableLambda 四个构造函数生成,签名同样受泛型约束。

1.1 图执行时允不允许"绕回去":Pregel 与 DAG 两种引擎

图定义好了,接下来是怎么执行它。这里要解决的问题是:允不允许"某个节点的结果绕回前面的节点",也就是允不允许出现循环(比如"工具调用结果要重新喂回模型")?compose/graph.go 里为此定义了两种运行模式:

const (
    runTypePregel graphRunType = "Pregel" // 可以有环,兼容 AnyPredecessor 触发
    runTypeDAG    graphRunType = "DAG"    // 严格有向无环,兼容 AllPredecessor 触发
)

对应到 compose/types.goNodeTriggerModeAnyPredecessor(默认,Pregel 语义——任一前驱完成即触发,因此允许"工具结果回填后再次进入模型节点"这类循环)和 AllPredecessor(所有前驱都完成才触发,用于严格 DAG/Workflow 场景)。切换方式是 compose.WithNodeTriggerMode(...)compose/graph_compile_options.go)。之前理解的"状态转移回溯 (Loop Back)"确有依据,但底层机制是 Pregel 的超步(super-step)触发模型,而不是什么特殊的"回溯"标记。

1.2 怎么让图"根据模型输出走不同的路":分支路由的真实签名

除了顺序执行,很多场景需要"根据上一个节点的输出决定接下来走哪条路"(比如模型要不要调用工具)。这靠分支(Branch)实现:

// compose/branch.go
type GraphBranchCondition[T any] func(ctx context.Context, in T) (endNode string, err error)

func NewGraphBranch[T any](condition GraphBranchCondition[T], endNodes map[string]bool) *GraphBranch

endNodes 是一个白名单:condition 返回的目标节点如果不在这个 map 里,NewGraphMultiBranch 内部会直接报错(fmt.Errorf("branch invocation returns unintended end node: %s", end)),这是编译期之外唯一的一层运行时兜底校验。


2. 流式输出是怎么处理的:StreamReader 与四个方法互相推导

大模型的响应经常是一边生成一边吐 token 出来的(流式),但有些场景你只想要一次性的完整结果(非流式)。eino 不想让每个组件都写两套逻辑,所以设计了一套"实现一个方法、框架自动补全其余三个"的机制。下面先看流本身是什么,再看这套自动补全怎么工作。

2.1 流的物理实现

schema.Pipe[T](cap int) (*StreamReader[T], *StreamWriter[T])schema/stream.go)本质是一个带容量的 channel 封装,StreamReader 提供 Recv() / Close() / Copy(n int)(用于给多个下游/多个 Callback Handler 各发一份)等方法。没有什么"零拷贝级联通道"的黑魔法,就是标准的 Go channel + 引用计数关闭。

2.2 Runnable 的四个方法:为什么只用实现一个,其余三个会自动补全

Eino 组件对外统一暴露 compose.Runnable[I, O] 接口(compose/runnable.go),包含四个方法:

type Runnable[I, O any] interface {
    Invoke(ctx context.Context, input I, opts ...Option) (output O, err error)
    Stream(ctx context.Context, input I, opts ...Option) (output *schema.StreamReader[O], err error)
    Collect(ctx context.Context, input *schema.StreamReader[I], opts ...Option) (output O, err error)
    Transform(ctx context.Context, input *schema.StreamReader[I], opts ...Option) (output *schema.StreamReader[O], err error)
}

之前说的"流的自动降级/级联"并非虚构,真实机制是:用户只需要实现其中一个方法,框架用另外三个函数把它补全。例如只实现了 Transform(流进流出)时,Invoke 会这样被自动派生(compose/runnable.go):

func invokeByTransform[I, O, TOption any](t Transform[I, O, TOption]) Invoke[I, O, TOption] {
    return func(ctx context.Context, input I, opts ...TOption) (output O, err error) {
        srInput := schema.StreamReaderFromArray([]I{input})
        srOutput, err := t(ctx, srInput, opts...)
        if err != nil {
            return output, err
        }
        return defaultImplConcatStreamReader(srOutput)
    }
}

即:把单个输入包成一个只有一个元素的"伪流",调用 Transform,再把输出流拼接(concat)回单体值。反过来 streamByInvoke 则是把 Invoke 的单个输出包装成一个单元素的流。这套机制是纯函数式的组合,不涉及运行时反射黑魔法。

2.3 流式场景下的一个具体坑:工具调用 Chunk 的真实拼接逻辑

流式场景下有一个容易被忽略的细节:模型返回的"要调用哪个工具、参数是什么"这条信息本身也是被切碎分片吐出来的,需要在框架层面拼接回完整的一条。当 ChatModel 流式输出时,ToolCall.Function.Arguments 就是被切碎的 JSON 片段。真实拼接函数是 concatToolCallsschema/message.go:1284),核心逻辑:按 ToolCall.Index 分组,对同一 Index 的多个 chunk 做 ID/Type/Function.Name 一致性校验(不一致直接返回 error),并用 strings.Builder 顺序追加 Arguments 字符串;最后按 Index 排序输出。这一层之上是 schema.ConcatMessages(msgs []*Message) (*Message, error)schema/message.go:1644)和消费流的 schema.ConcatMessageStream(s *StreamReader[*Message]) (*Message, error)schema/message.go:1842)。对于非 *Message 类型的自定义流类型,compose.RegisterStreamChunkConcatFunc[T]compose/stream_concat.go)允许用户注册自己的拼接函数,供只实现了 Stream 却被以 Invoke 方式调用的组件使用。

2.4 长流程怎么暂停和恢复:Callback 与 Interrupt/Resume

之前担心的"Callback 与 Interrupt 中断恢复机制是不是杜撰"——查证后这两个机制都是真实存在的,而且实现相当完整:

  • Callback:callbacks/interface.go 定义了统一的 Handler 接口(OnStart/OnEnd/OnError/OnStartWithStreamInput/OnEndWithStreamOutput),通过 RunInfo(节点名、组件类型、组件分类)区分触发来源,文档里特别强调"不同 Handler 之间没有执行顺序保证,流式回调拿到的 StreamReader 是复制品,用完必须 Close(),否则原始流无法释放"。
  • Interrupt/Resume:compose/interrupt.go 提供 compose.Interrupt(ctx, info) errorcompose.StatefulInterrupt(...),配合 compose.WithInterruptBeforeNodes([]string) / WithInterruptAfterNodes([]string)(编译选项)在指定节点前后中断执行并落一个 checkpoint;compose/checkpoint.go 提供 CheckPointStore 接口和 compose.WithCheckPointStore(store) / WithCheckPointID(id)compose/resume.go 提供 compose.GetInterruptState[T](ctx)GetResumeContext[T](ctx) 供节点在恢复执行时判断"我是不是被中断过、有没有拿到外部注入的续跑数据"。整个机制基于 internal/core.InterruptSignal,用 Address(节点路径)定位到具体是图里哪一个节点触发的中断,支持嵌套子图。

这套东西的实际定位类似 LangGraph 的 Human-in-the-loop / Checkpoint,不是文档拍脑袋编出来的概念。


3. 完整实践:在最小例子基础上加上工具调用和分支路由

前面第 0 节的最小例子只有"拼 Prompt → 调模型"两步。真实的智能体还需要"模型判断要不要调用工具、调用完再把结果喂回去"这一层分支逻辑。以下示例的每一行 API 都能在源码里找到对应签名(compose.NewGraphAddLambdaNodeAddChatModelNodeAddToolsNodeAddBranchNewGraphBranchWithMaxRunSteps 均为真实导出符号),组装方式参考了仓库 README.md 中 Composition 一节的写法:

package main

import (
    "context"
    "fmt"

    "github.com/cloudwego/eino/compose"
    "github.com/cloudwego/eino/schema"
)

func main() {
    ctx := context.Background()

    // compose/generic_graph.go: NewGraph[I, O any](opts ...NewGraphOption) *Graph[I, O]
    graph := compose.NewGraph[string, *schema.Message]()

    // Lambda 节点必须由 compose.InvokableLambda 等构造函数生成(compose/types_lambda.go)
    _ = graph.AddLambdaNode("prompt", compose.InvokableLambda(
        func(ctx context.Context, in string) (string, error) {
            return fmt.Sprintf("System: Answer concisely.\nUser: %s", in), nil
        }))

    // AddChatModelNode 要求真实实现 model.BaseChatModel 接口的组件(如 openai.ChatModel)
    _ = graph.AddChatModelNode("model", chatModel)

    // AddToolsNode 要求 *compose.ToolsNode,由 compose.NewToolNode(ctx, conf) 构造
    _ = graph.AddToolsNode("tools", toolsNode)

    _ = graph.AddEdge(compose.START, "prompt")
    _ = graph.AddEdge("prompt", "model")

    // compose/branch.go: NewGraphBranch[T](condition GraphBranchCondition[T], endNodes map[string]bool) *GraphBranch
    _ = graph.AddBranch("model", compose.NewGraphBranch(
        func(ctx context.Context, msg *schema.Message) (string, error) {
            if len(msg.ToolCalls) > 0 {
                return "tools", nil
            }
            return compose.END, nil
        },
        map[string]bool{"tools": true, compose.END: true},
    ))

    _ = graph.AddEdge("tools", compose.END)

    // compose/graph_compile_options.go: WithMaxRunSteps 是真实的图级"最大执行步数"护栏
    runnable, err := graph.Compile(ctx, compose.WithMaxRunSteps(10))
    if err != nil {
        panic(err)
    }

    out, err := runnable.Invoke(ctx, "How is the weather?")
    if err != nil {
        panic(err)
    }
    fmt.Println(out.Content)
}

chatModelmodel.BaseChatModel 实现)和 toolsNode(通过 compose.NewToolNode(ctx, &compose.ToolsNodeConfig{...}) 构造,见 compose/tool_node.go)需要具体的模型/工具实现才能跑通,这里略去初始化代码,重点是图的拓扑组装方式与之前版本编造的 compose.NewNode 完全不同——Eino 里根本不存在 compose.NewNode 这个符号


4. 生产级失效排查与性能抖动防护(护栏名称已核对为真实符号)

故障模式 底层诱因 系统现象 防御与排查手段
拓扑死循环 (Pregel 无限超步) AnyPredecessor 触发模式下注册了环形链路,且分支条件永远走不到 compose.END 协程/内存缓慢增长,Invoke 长时间不返回。 编译时用 compose.WithMaxRunSteps(n)compose/graph_compile_options.go)强制设定超步上限,超过即返回 error,而不是无限跑下去。
流式内存堆积 schema.Pipe[T](cap)cap 设置过大,或下游消费速度长期慢于上游生产。 内存占用上升、GC 抖动。 合理设置 Pipe 的容量参数;用完的 StreamReader.Copy(n) 副本必须逐个 Close(),否则底层引用计数不会归零,见 callbacks/interface.go 对流式 Handler 的强制要求。
节点注册期类型不匹配 传给 AddChatModelNode/AddLambdaNode 等方法的实现和该方法要求的接口(如 model.BaseChatModel)或 InvokableLambda[I, O] 的泛型参数不一致。 编译期直接报错(Go 泛型 + 接口约束),根本到不了运行时。 这正是 Eino 用泛型替代 map[string]any 的核心收益:把"节点接错线"从运行时 panic 前移到 go build 阶段。
中断恢复丢状态 使用了 compose.Interrupt 中断执行,但没有配置 compose.WithCheckPointStore,或者恢复时忘记调用 compose.GetResumeContext 读取续跑数据。 恢复后节点拿不到中断前的状态,行为等价于重新跑了一遍。 严格按 compose/checkpoint.go + compose/resume.go 的三件套使用:WithCheckPointStore 落盘、WithCheckPointID 定位、GetInterruptState/GetResumeContext 在节点内部读取。

5. 资深系统架构师面试表达方案

面试提问:在构建大规模智能体(Agent)系统时,你们为什么要选用 Eino 这套 DAG 框架,相比于 Python 派系的 LangChain 它的优势在哪?

回答模版: 选 Eino 而不是直接抄一套 Python 风格的 Agent 框架,很大程度是被一次线上事故逼的——早期我们用类似 LangChain 的思路直接用 map[string]any 传参,某次某个节点漏填了一个 key,一路传到很后面的节点才崩,日志里完全看不出问题出在哪一层。换成 Eino 之后,每种节点都有专门的注册方法(AddChatModelNodeAddToolsNodeAddLambdaNode……),输入输出类型在 AddNode 那一行就被 Go 泛型和组件接口锁死了,这类问题基本在编译期就能暴露,省了不少运行时排查的功夫。

流式这块也踩过一次坑:早期我们没注意到 Eino 的 Runnable 四个方法(Invoke/Stream/Collect/Transform)是互相推导的——只要实现一个,其余三个框架会自动用"包一层单元素流再拼接回去"的方式补全——结果我们自己在下游又手动做了一次 Chunk 聚合,相当于聚合了两遍,内存占用莫名其妙涨了不少。后来翻了 compose/runnable.go 才搞清楚这套自动补全的实现方式,才知道自己没必要重复实现。整体感受是,Eino 用编译期类型换取了运行时的确定性,但也意味着写自定义 Node 时得老老实实按泛型接口来,不能像 Python 那样随手塞一个函数进去就能跑。