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 把这张图编译成一个可以直接 Invoke 的 Runnable。下面展开讲这套拓扑到底是怎么定义和执行的,以及流式、分支、中断恢复这些进阶能力。
eino 的核心设计思路(人话版)¶
抛开术语,eino 要解决的问题也很朴素:
- 怎么把"调用模型 → 判断要不要调用工具 → 调用工具 → 把结果喂回模型"这套多步骤逻辑,用画图的方式表达出来,而不是写一堆嵌套 if/else?—— eino 用 Graph(可以带分支、甚至带环)或更简单的 Chain(纯线性)描述这套流程,你只管往图里加节点、连边,编译后拿到一个可以直接调用的
Runnable。 - 大模型的输出往往是流式的(一个个 token 蹦出来),怎么让"流式输出"和"一次性拿完整结果"这两种用法在框架层面统一,而不是让每个组件自己维护两套代码?—— eino 定义了
Invoke/Stream/Collect/Transform四个方法,只要实现其中一个,框架自动帮你推导出另外三个。 - 长时间运行的智能体中途出错,或者需要人工介入怎么办?—— 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.BaseChatModel、prompt.ChatTemplate……)在编译期锁死的,而不是像 Python 框架那样先塞进一个通用容器再在运行时反射猜测类型。唯一的"自由节点"是 Lambda(compose/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.go 的 NodeTriggerMode:AnyPredecessor(默认,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 片段。真实拼接函数是 concatToolCalls(schema/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) error和compose.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.NewGraph、AddLambdaNode、AddChatModelNode、AddToolsNode、AddBranch、NewGraphBranch、WithMaxRunSteps 均为真实导出符号),组装方式参考了仓库 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)
}
chatModel(model.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 之后,每种节点都有专门的注册方法(AddChatModelNode、AddToolsNode、AddLambdaNode……),输入输出类型在 AddNode 那一行就被 Go 泛型和组件接口锁死了,这类问题基本在编译期就能暴露,省了不少运行时排查的功夫。
流式这块也踩过一次坑:早期我们没注意到 Eino 的 Runnable 四个方法(Invoke/Stream/Collect/Transform)是互相推导的——只要实现一个,其余三个框架会自动用"包一层单元素流再拼接回去"的方式补全——结果我们自己在下游又手动做了一次 Chunk 聚合,相当于聚合了两遍,内存占用莫名其妙涨了不少。后来翻了 compose/runnable.go 才搞清楚这套自动补全的实现方式,才知道自己没必要重复实现。整体感受是,Eino 用编译期类型换取了运行时的确定性,但也意味着写自定义 Node 时得老老实实按泛型接口来,不能像 Python 那样随手塞一个函数进去就能跑。