0%

Eino 运行时分析:从编译期单例到协程级调度

1. 缘起

上一篇《Agent Orchestration - 智能体编排》写的是 Eino 怎么实现多智能体,讲的是用法。

Eino 是 CloudWeGo 出的 Go 语言 LLM 应用编排框架。
它把 ChatModel、Tool、Retriever 这些组件连成图,编译成一个 Runnable 来跑。
README 里写明参考了 LangChain 和 Google ADK,所以熟悉那两个框架的人会觉得眼熟。

通过分析运行时,来了解 Eino 内部调度、状态管理、中间件、可观测等方面。

2. 全景:一次请求的完整坐标

编译期的产物被所有请求共享,运行期的状态每请求一份。这条边界是后面每一节都要落回去的地方。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
┌─ 编译期:进程内一份,只读,被所有请求共享
│ Graph / Workflow / Chain 声明
│ │
│ ▼ compose.Compile(ctx)

│ Compiled Runnable
│ ├── Nodes map[string]*nodeRunner,各节点的无状态执行闭包
│ ├── Dependency Table DAG 入度表 / Pregel 邻接表
│ ├── Converters 类型适配管道,Stream ↔ Value、FieldMapping
│ ├── Middleware Chain 编译期串好的装饰器链
│ └── Callback Chain 编译期固化的监听链

│ ▸ Compile() 做掉四件事:类型检查、拓扑分析、转换层注入、中间件接线
└─

│ 拓扑与接线在编译期固化,运行时不再改


┌─ 运行期:每请求一份,互不可见
│ HTTP / RPC 请求进入,携带独立的 req.Context
│ Ephemeral Execution Instance
│ ├── 独立 Context TraceID / Deadline / 取消信号
│ ├── 独立 State Table 请求级私有状态,无跨请求污染
│ └── Goroutine Tree 调度协程树,取消时整棵退出(§9)
└─

3. 编译期:Compile 做的四件事

Eino 的图是声明式的。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
┌─ 声明期:只是在描述一张图,不能跑
│ graph := compose.NewGraph[In, Out]()
│ graph.AddChatModelNode("model", m)
│ graph.AddToolsNode("tools", t)
│ graph.AddBranch("model", router)
│ graph.AddEdge(START, "model")
│ graph.AddEdge("tools", "model") ← 回边,形成环
└─


▼ compose.Compile(ctx)

┌─ 编译期:做掉四件事
│ 1. 类型检查 相邻节点的输出类型能不能接上
│ 2. 拓扑分析 有没有环,决定走 Pregel 还是 DAG
│ 3. 转换层注入 流式节点接到非流式节点时,插入 drain 适配器
│ 4. 中间件接线 Middleware 按接口变体串成装饰器链

│ 产物是 Compiled Runnable,此后只读
└─

编译期做掉的事情,运行时就不做了。这是它跟「每个请求反射建图」类框架最大的差别。

compose/graph_compile_options.go 里有个选项跟运行时形态直接相关:

1
2
WithMaxRunSteps(maxSteps int)       // 循环图的最大步数,超过返回 ErrExceedMaxSteps
WithEagerExecution(...) // 是否启用 Eager 执行,见 §5

还有一个运行期的 WithRuntimeMaxSteps,可以在单次调用时覆盖编译期的设置。

4. 请求级状态:ProcessState 与节点 Hook

编译产物共享,请求级状态不能共享。Eino 的边界划在 context.Context 上。

读写都走 ProcessState

1
2
3
4
compose.ProcessState[S](ctx, func(ctx context.Context, s S) error {
s.Counter++ // 同一时刻只有一个 handler 在跑
return nil
})

v0.4.0 之前是 GetState[S](ctx),返回状态对象本身,并发保护要调用方自己保证;改成回调之后,这个责任收到框架里。

配套还有四个节点级 Hook,用来在节点执行前后读写状态:

1
2
3
4
StatePreHandler[I, S]              // 普通节点,执行前
StatePostHandler[O, S] // 普通节点,执行后
StreamStatePreHandler[I, S] // 流式节点,执行前
StreamStatePostHandler[O, S] // 流式节点,执行后

以及 GenLocalState[S],给每个请求生成独立的初始状态。

5. 调度:Pregel 与 DAG 两条路径

Eino 有两种触发模式,由 NodeTriggerMode 区分。源码里 compose/pregel.gocompose/dag.go 是两个独立文件。

模式常量适用触发条件
PregelAnyPredecessor循环图(ReAct)任一前置节点在上一步完成
DAGAllPredecessor无环图所有前置节点都完成

Pregel:超步屏障

带环的图,比如思考 → 调工具 → 再思考,需要严格的步数同步,否则会出现节点读到上一轮半成品的情况。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
┌─ Step N:本轮活跃节点并发跑
│ ├── Model 协程 产出 ToolCall 决策
│ ├── 日志协程 记录本轮输入输出
│ └── 校验协程 检查参数合法性
│ │
│ ▼ 所有产出写入 Mailbox,不直接投递给下游
├─ ══ 同步屏障 ══
│ 等 Step N 内所有活跃协程全部退出,Mailbox 里的数据才算本轮定稿
└─




┌─ 路由评估
│ 是否调 Tool ? ── 是 ──→ 进入 Step N+1
│ 是否结束 ? ── 是 ──→ 走向 END
│ 步数 > MaxRunSteps ? ── 是 ──→ ErrExceedMaxSteps
└─

屏障的意义是把并发的网状依赖压成一轮一轮的批次,杜绝同一轮内节点互相看到中间态。

DAG:Eager 触发

无环图走 AllPredecessor,哪个节点的入边到齐了就立刻触发,不等整步。v0.4.0 起改成 Eager,早期 DAG 复用 Pregel 的超步屏障,A、B 跑完后 C 明明可以跑,却要等整步同步。

1
2
3
4
5
6
7
┌─ DAG 的触发
│ A ──┐
│ B ──┼──→ D 谁的入边到齐就立刻跑,不等整步
│ C ──┘

│ 注意:Eager 与 AnyPredecessor 不兼容,循环图仍走 Pregel
└─

所以「超步屏障消除竞态」这句话,只对 Pregel 模式成立。DAG 模式是 Eager 的,同步开销被拿掉了。

6. 数据流转:值、流、转换、背压

节点之间流转的形态有三种:普通值、流,以及流被拼装回的值。转换层在编译期接好,背压交给 Go channel 的容量。

流降级为值:隐式屏障

当一个流式节点连到一个只接受普通值的节点时,编译器会在中间插一个转换层。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
┌─ 分支汇聚,接线在编译期完成

├─ 分支 A:Model 节点,流式输出
│ │ StreamReader[*Message]
│ ▼
│ drain 适配器
│ ├─ 起一个协程循环 Recv()
│ ├─ 调 schema.ConcatMessages() 拼装
│ └─ 等另一个分支的普通值也到达
│ ▲
│ │ DocList
├─ 分支 B:检索节点,普通值

├─ 两个字段都就绪 → FieldMapping 组装

└─ 传给总结节点,非流式

运行时表现为:流被 drain 完,拼装成完整对象,再连同另一个分支的普通值一起传给下游。

字段映射的代码在 compose/field_mapping.go。好处是编译期就能发现字段名写错,比用 map[string]any 传数据、跑到运行时才崩要好。

Pipe 的容量就是背压阈值

1
sr, sw := schema.Pipe[string](3)

容量 3 表示通道里最多缓存 3 个元素,写第 4 个时 Send 阻塞,直到有人读走。Eino 没有另造一套背压机制,就是把 Go channel 的语义封成了 StreamReader / StreamWriter

StreamReader[T] 本身是联合类型,内部按 typ 分五种:

类型来源
stream[T]Pipe[T]() 创建的通道流
arrayReader[T]StreamReaderFromArray[T](),纯索引,零开销
multiStreamReader[T]MergeStreamReaders(),多流合并
streamReaderWithConvert[T]转换或过滤
childStreamReader[T]Copy() 扇出的子流

背压的传导与流的关闭

背压不是一条传到底的链条。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
┌─ 第一段:本地传导,机制保证
│ [ 下游消费慢 ]
│ │
│ ▼
│ [ Pipe 通道被塞满,Send 阻塞 ]
│ │
│ ▼
│ [ 发送协程被 Go 调度器挂起,gopark ]
│ │
│ ▼
│ [ 本地停止从上游读取 ]

│ 到这里为止是确定的
└─


▼ 再往上传不动,取决于上游是什么

┌─ 第二段:分两种情况

├─ 情况 A:上游是网络流,比如直连模型服务的 SSE
│ 本地不读
│ → socket 接收缓冲区填满
│ → TCP 通告零窗口,Window Size = 0
│ → 服务端冻结发包
│ → 模型引擎 write() 阻塞,暂停生成

├─ 情况 B:上游是已经读完的内存响应
│ 本地不读,压力传不回去,止步于本地,内存该占还是占

└─

所以「Eino 的背压能让大模型服务端停止生成」这句话,在直连流式的场景下成立,但不能当成普适结论。中间隔一层LLM API 中转站的话,背压能传到哪一层,取决于它边收边转还是收完再转。

两端都要 Close:写端 Close() 通知消费者 EOF,读端 Close() 通知写端「我不要了」,否则写端 Send 永远阻塞,goroutine 泄漏。

还有个 SetAutomaticClose(),内部用 runtime.SetFinalizer 在对象被 GC 时自动关闭。但它需要主动调用才生效,源码注释也写明 NOT concurrency safe。适合兜底,不适合放进关键路径。

Copy 是共享链表,不是复制

Copy(n) 把一个流扇出给 n 个消费者,但每个消费者读的不是独立副本。

1
2
3
4
5
6
7
// Each value comes from a hidden linked list of cpStreamElement.
elem.once.Do(func() {
t, err = p.sr.Recv() // 上游只读一次
elem.item = streamItem[T]{chunk: t, err: err}
elem.next = &cpStreamElement[T]{}
p.subStreamList[idx] = elem.next
})
1
2
3
4
5
6
7
8
9
10
11
12
13
14
┌─ Copy(3) 之后的读序
│ 上游 stream
│ │
│ ▼
│ parentStreamReader
│ │
│ ├─ subStreamList[0] ──→ ● ──→ ● ──→ ● 子 reader 0 游标
│ ├─ subStreamList[1] ──→ ● ──→ ● 子 reader 1 游标
│ └─ subStreamList[2] ──→ ● 子 reader 2 游标
│ ▲
│ 同一批链表节点,值不复制

│ 每个节点由 sync.Once 保护,上游 Recv() 只调一次
└─

代价是链表节点会保留到所有子 reader 都消费过为止。某个子 reader 慢,后面的节点就堆着不放。省的是复制开销,换的是内存占用随最慢消费者增长。

调用 Copy 之后原来的 sr 失效,只能用返回的副本。

7. 横切:Callback 读,Middleware 写

运行时有两套横切机制,分工不同。

  • Callback 的签名是 OnStart / OnEnd / OnError,只能观察,用来做 trace、指标、审计
  • Middleware 的签名是 func(Endpoint) Endpoint,能改输入、改输出、或者短路不调用,用来做重试、熔断、脱敏、缓存

重试、熔断这类失败处理,在《稳定性:从传统工程到大模型时代》里属于设计阶段就该定下来的事;落到 Eino 里,承载它们的就是 Middleware。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
┌─ 一次节点调用的剖面
│ 节点入参
│ │
│ ├── ◄── Middleware 在这一层:能改参数、能直接返回
│ │
│ ▼
│ ┌─ 节点执行体
│ │ OnStart ← Callback 只能看
│ │
│ │ 真正执行节点逻辑
│ │
│ │ OnEnd / OnError ← 只能看
│ └─
│ │
│ ▼
│ 节点出参
│ │
│ └── ◄── Middleware 也在这一层:能改结果
└─

一个容易踩的细节,源码注释写得很明确:

This middleware only applies to tools that implement the InvokableTool interface.

InvokableToolMiddleware 只作用于实现了 InvokableTool 的工具。工具如果只实现了 StreamableTool,这个中间件对它无效。挂中间件前得先确认工具实现了哪个接口。

ADK 层还有 AgentMiddleware / ChatModelAgentMiddleware,通过 WrapModelWrapTool 包住模型与工具。

观测与改造的分工

要观测就用 Callback。用 Middleware 做 trace 有个隐患:如果中间件会短路,比如缓存命中或熔断打开,trace 会记成一次成功调用,但底层根本没执行,指标上分不出来。

要改行为就用 Middleware。这时候 trace 作为副产品自然就有了。

8. 工具:五种接口变体与中断

工具不是一个接口,是五个。

1
2
3
4
5
BaseTool                    // 只要 Info(),能描述自己
InvokableTool // InvokableRun(ctx, argsJSON) (string, error)
StreamableTool // StreamableRun(ctx, argsJSON) (*StreamReader[string], error)
EnhancedInvokableTool // 带结构化的 ToolArgument / ToolResult
EnhancedStreamableTool

运行时按工具实际实现了哪些接口决定调用路径:调用方声明要流式,工具没实现 StreamableTool 就降级或报错;声明要结构化参数,没实现 EnhancedInvokableTool 就回落到 JSON 字符串。

工具可以主动中断

components/tool/interrupt.go 提供了一组中断原语:

1
2
3
4
5
tool.Interrupt(ctx, info)                            // 中断,带上下文信息
tool.StatefulInterrupt(ctx, info, state) // 中断并保存现场
tool.CompositeInterrupt(ctx, info, state, errs...) // 多个中断合并
tool.GetInterruptState[T](ctx) // 恢复时取回现场
tool.GetResumeContext[T](ctx) // 判断自己是不是恢复目标
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
┌─ 工具内中断的时间线
│ ToolsNode 执行到一半,需要人确认
│ │
│ ▼
│ tool.StatefulInterrupt(ctx, info, state)
│ │
│ ▼
│ 整图就地挂起,现场被 checkpoint 序列化
│ │
│ ⋯ 等人,几分钟到几天
│ │
│ ▼
│ 带 checkpoint 重新进入图
│ │
│ ▼
│ 跳过已完成节点,工具内 GetInterruptState 取回现场,继续跑
└─

这是 human-in-the-loop 的基础,跟 §9 的 checkpoint 是一条线。

MCP 在运行时里没有特殊地位

MCP 工具在 eino-ext/components/tool/mcp,不在核心仓库。从运行时视角看,它跟本地工具没有区别,就是一个 Tool 实现。

真正需要注意的是 schema 的时机:MCP 工具的 Info() 通常要连服务端拿 schema,比本地工具多一次网络往返,这次往返发生在初始化阶段而不是每次调用。

9. 请求终结:取消、错误、中断恢复

取消是 context 级联

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
┌─ 请求 ctx:Deadline 5s,TraceID abc123
│ │
│ ├── 派生 span ctx ──→ 协程 1:Model 节点
│ │
│ └── 派生 span ctx ──→ 协程 2:检索节点
│ │
│ │ ← 客户端断开 / 超时
│ ▼
│ ctx.Done() 触发
│ │
│ ▼
│ 级联响应
│ ├─ 中断 DB 连接
│ ├─ 关闭 PipeWriter
│ └─ 触发 OnError
│ │
│ ▼
│ 整棵协程树退出,状态随实例释放
└─

关键在于所有可能阻塞的操作都吃同一个 ctx。一处取消,整棵树跟着退。这也是为什么每个节点拿到的是派生 ctx 而不是原始 ctx。

错误有类型标签

1
2
internalErrorTypeNodeRun  = "NodeRunError"
internalErrorTypeGraphRun = "GraphRunError"

节点错误包一层打 NodeRunError,图调度层再包一层打 GraphRunError。调用方配合 ErrExceedMaxSteps 这类哨兵错误,可以用 errors.Is 分类,不用匹配错误字符串。

中断与恢复

实现分散在 compose/ 的三个文件里:interrupt.go 挂起,checkpoint.go 序列化现场(已完成节点、状态表、中断点),resume.go 带着 checkpoint 重新进图。定位是长流程和人工介入,一个需要人审批的流程可以停下来等几小时甚至几天,然后接着跑。

10. ADK 层:ReAct 循环与模型容错

compose 给的是通用编排能力。ADK 面向 Agent 场景,加了这些东西:

能力文件
ReAct 循环react.go
模型重试retry_chatmodel.go
模型故障转移failover_chatmodel.go
中断与取消interrupt.gocancel.go
轮次管理turn_loop.goturn_buffer.go
指令与配置instruction.goconfig.go

其中重试和故障转移是 §7 讲的 Middleware 模式的具体应用:

1
2
3
4
5
6
7
8
9
10
11
┌─ buildModelWrappers 的装饰顺序
│ 原始 BaseModel
│ │
│ ├── 包一层:callback 注入 typedCallbackInjectionModelWrapper
│ ├── 包一层:事件发送 NewEventSenderModelWrapper
│ ├── 包一层:重试 retry
│ └── 包一层:故障转移 failover

│ 最外层拿到的还是 model.BaseModel[M] 接口
│ 调用方看不出里面套了几层
└─

这是「不改节点代码插入横切逻辑」的现成例子,也解释了为什么 §7 要区分 Callback 和 Middleware:ADK 自身就是靠 Middleware 实现重试和故障转移的。

turn_loop.go 管 Agent 的轮次。ReAct 的「思考 → 调工具 → 再思考」就是个循环,ADK 把它抽成了可配置的循环,比手写图更贴合 Agent 的语义。

11. 五条设计取舍

抛开 Eino,这几条对我做 Go 服务有参考价值。

  1. 编译期校验,运行时只读。类型、字段名、拓扑有没有环,全在 Compile() 时发现,运行时不需要加锁。代价是灵活性。

  2. 并发安全的责任收进框架。GetStateProcessState 的改动值得学:能靠结构避免的错误,不要靠文档提醒。

  3. 背压复用 Go channel 语义,不自造机制。不发明新概念,学习成本和出错概率都低。

  4. 读和写分开两套机制。混在一起的框架最后往往变成一个万能 Hook,什么都往里塞。

  5. 广播用共享链表而不是复制。代价是最慢的消费者会拖住内存,是个明确的权衡。


Eino 运行时的大部分复杂度,来自它要在共享的只读拓扑上跑隔离的并发请求。

12. 参考

  • cloudwego/eino · 源码仓库,本文的 API 名称与实现细节均以 main 分支为准
  • Eino README · 三层模块划分,以及 ADK 的设计来源
  • Stream Processing Internals · StreamReader 五种内部类型与背压
  • eino v0.4.0 更新解析 · Eager 执行默认开启、GetState 移除的背景
  • Eino 发布记录 · 各版本 API 变更