Eino学习笔记
Eino学习笔记 —— 入门篇
本笔记通过ai辅助生成,但对于代码阅读顺序应当没有问题(均使用eino官方的例子进行,个人认为官方的教程感觉有点太难懂了)
另附,可能是更新速度快,部分代码的方法已经被标记为废弃,但在示例中仍然未修改为最新的方法,建议自行对应
Eino 框架学习笔记
📂 源码仓库:github.com/cloudwego/eino-examples↗
📖 系列文档:入门笔记 | Agentic 进阶 | 附录一 Flow | 附录二 组件 | 附录三 工具
速览:两条调用路径 & 核心速记
在深入细节之前,先在心里放一张总地图。Eino 只有两条调用路径:
sequenceDiagram
actor U as 用户代码
participant CM as ChatModel
participant A as Agent
participant R as Runner
participant I as Iterator
Note over U,I: === 轻量路径:直接用 ChatModel(入门笔记 §2) ===
U->>CM: Generate(messages) / Stream(messages)
CM-->>U: *schema.Message / StreamReader
Note over U,I: === 完整路径:Agent + Runner(入门笔记 §3~§5) ===
U->>CM: ① NewChatModel(config)
CM-->>U: model
U->>A: ② NewChatModelAgent(model, instruction, tools)
A-->>U: agent
U->>R: ③ NewRunner(agent, streaming, checkpoint)
R-->>U: runner
U->>R: ④ Run(input) / Query(input, checkpointID)
R-->>I: AsyncIterator[*AgentEvent]
loop 消费事件
I->>I: Next() → (event, ok)
end
Note over U,I: === Agentic 路径:Responses API(Agentic 进阶) ===
U->>CM: ① New(agenticModel)
CM-->>U: model
U->>A: ② NewTypedChatModelAgent[*AgenticMessage](model, ...)
A-->>U: agent
U->>R: ③ NewTypedRunner[*AgenticMessage](agent, ...)
R-->>U: runner
U->>R: ④ Run(AgenticMessage输入)
R-->>I: Iterator
loop 消费事件
I->>I: TypedGetMessage(event) → msg.String()
end📇 全篇速记卡片(看完这篇你只需要记住这些)
创建链:NewChatModel → NewChatModelAgent → NewRunner → Run (模) (代) (跑) (消息)
两种迭代:Agent 用 Next() → (event, ok) —— !ok 结束 Stream 用 Recv() → (chunk, err) —— io.EOF 结束
Agent 管能力:{Name, Instruction, ToolsConfig, Model}Runner 管执行:{Agent, EnableStreaming, CheckPointStore}
三种 Tool 创建:90% 用 InferTool(函数+struct tag) 需要约束用 NewTool(schema+函数) 有状态用结构体实现接口
中断恢复三要素:CheckPointStore + CheckPointID + Resume
核心包(按使用频率): github.com/cloudwego/eino/adk ← Agent, Runner, Message github.com/cloudwego/eino/schema ← UserMessage, SystemMessage github.com/cloudwego/eino/components/tool/utils ← InferTool, NewTool github.com/cloudwego/eino-ext/components/model/openai ← ChatModel⚡ 关键概念:Iterator(事件迭代器)是什么?
在所有 Eino 示例中,runner.Run() 和 runner.Query() 都会立即返回一个 AsyncIterator[*AgentEvent](简称 Iterator),然后你通过 Next() 逐个消费事件。这是 Eino 最核心的数据消费模式。
为什么要有 Iterator?
Agent 执行是异步的:模型调用可能耗时几秒到几十秒,期间可能经历”思考→调Tool→再思考→输出文本”多个阶段。如果 Run() 等全部完成才返回,你的程序就卡住了。
Iterator 解决了这个问题——立即返回一个”事件管道”,Agent 在后台生成事件、推入管道,你在前台逐个取出处理。类比:
Runner.Run() 返回 Iterator = 你拿到一个"对讲机"Agent 在后台干活,每隔一会通过对讲机说句话你在前台: Next() → 收到一条消息 → Next() → 收到下一条 → ...Next() 返回 (event, false) = 对讲机没信号了 = Agent 干完了Iterator 的本质
type AsyncIterator[T any] struct { // 内部有 channel,连接"生产者"(Agent goroutine)和"消费者"(你的代码)}
func (iter *AsyncIterator[T]) Next() (item T, ok bool) { // 阻塞等待下一个事件 // ok=true: 拿到了新事件 // ok=false: 管道关闭,Agent 执行完毕}与 Go channel 的对比
| Go channel | AsyncIterator | |
|---|---|---|
| 创建 | make(chan T) | Runner.Run() / Runner.Query() 返回 |
| 发送 | ch <- item | Agent 内部通过 gen.Send(event) |
| 接收 | item := <-ch | event, ok := iter.Next() |
| 关闭 | close(ch) | Agent 内部通过 gen.Close() |
| 关闭检测 | item, ok := <-ch | event, ok := iter.Next() — 同样是 ok=false |
🧠 一句话:Iterator 就是包装了 channel 的”事件流”——Agent 在后台往里面写事件,你在前台用
Next()逐个读。读完(ok=false)Agent 就结束了。
一、Hello World:认识四步组装线
📁 代码:
adk/helloworld/helloworld.go
为什么这样设计?
Eino 把一个 AI 对话应用拆成四个独立零件,像组装流水线一样逐个创建、最后拼起来运行。这样拆的好处是:每个零件可以单独替换(比如换个模型、换个 Agent 类型),互不影响。
四步组装线(记忆口诀:“模→代→跑→消息”)
第1步 创建 ChatModel → 第2步 创建 Agent → 第3步 创建 Runner → 第4步 发送消息并运行("嘴巴",负责调用LLM) ("大脑",封装业务逻辑) ("引擎",管理执行过程) (用户说的话)示范代码(带注释)
package main
import ( "context" "fmt" "log" "os"
"github.com/cloudwego/eino/adk" "github.com/cloudwego/eino/schema" "github.com/cloudwego/eino-ext/components/model/openai")
func main() { // 【固定写法】所有 Eino 调用的"通行证",贯穿整个调用链 ctx := context.Background()
// ─── 第1步:创建 ChatModel("嘴巴")─── // ChatModel = 与 LLM 通信的组件,屏蔽不同厂商差异 model, err := openai.NewChatModel(ctx, &openai.ChatModelConfig{ APIKey: os.Getenv("OPENAI_API_KEY"), Model: os.Getenv("OPENAI_MODEL"), BaseURL: os.Getenv("OPENAI_BASE_URL"), }) if err != nil { log.Fatal(err) }
// ─── 第2步:创建 Agent("大脑")─── agent, err := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{ Name: "hello_agent", Description: "A friendly greeting assistant", Instruction: "You are a friendly assistant. Please respond warmly.", Model: model, // ← 把第1步的 model 注入进来 }) if err != nil { log.Fatal(err) }
// ─── 第3步:创建 Runner("引擎")─── runner := adk.NewRunner(ctx, adk.RunnerConfig{ Agent: agent, EnableStreaming: true, })
// ─── 第4步:组装消息并运行 ─── input := []adk.Message{ schema.UserMessage("Hello, please introduce yourself."), }
// runner.Run() 返回事件迭代器,逐个消费事件 events := runner.Run(ctx, input) for { event, ok := events.Next() if !ok { break } if event.Err != nil { log.Printf("错误: %v", event.Err) break } if msg, err := event.Output.MessageOutput.GetMessage(); err == nil { fmt.Printf("Agent: %s\n", msg.Content) } }}💡 关键理解
| 概念 | 一句话 | 类比 |
|---|---|---|
context.Background() | 所有调用的”通行证”,固定写法 | 身份证,每次办事都要出示 |
ChatModel | 负责与 LLM 通信的组件 | 数据库驱动(屏蔽 MySQL/PG 差异) |
Agent | 封装”怎么用模型”的业务逻辑 | 业务逻辑层 |
Runner | 管理 Agent 的执行生命周期 | 汽车引擎(你只转方向盘,不管内部) |
events.Next() | 消费下一个事件,没事件了返回 false | 自助餐流水线,一盘盘拿 |
🧠 记忆技巧
把创建过程记成一句话:“用某个模型(ChatModel),造一个智能体(Agent),交给引擎(Runner)去跑(Run)”
NewChatModel → NewChatModelAgent → NewRunner → Run (模) (代) (跑) (消息)常用 import 路径:
"github.com/cloudwego/eino-ext/components/model/openai" // OpenAI 兼容模型,更多代码支持不再列出,直接参考官方文档"github.com/cloudwego/eino/adk" // Agent, Runner, Message"github.com/cloudwego/eino/schema" // UserMessage, SystemMessage 等💡 类型别名说明:
adk.Message就是*schema.Message(Go type alias)。代码里写[]adk.Message{ schema.UserMessage(...) }和[]*schema.Message{ schema.UserMessage(...) }是等价的。文档里两种写法都会出现——见到时知道是同一个东西即可。
📇 本节速记卡片
NewChatModel → NewChatModelAgent → NewRunner → Runevent, ok := iter.Next(); if !ok { break }msg, err := event.Output.MessageOutput.GetMessage()二、直接用 ChatModel(不用 Agent)
📁 代码:
quickstart/chat/main.go、components/prompt/chat_prompt/chat_prompt.go
什么时候不用 Agent?
如果只需要一次性的问答(不涉及多轮对话、工具调用、中断恢复),直接用 ChatModel 更轻量。
两种输出方式
| 方式 | 方法 | 效果 | 适用场景 |
|---|---|---|---|
| 一次性 | model.Generate(ctx, messages) | 等全部生成完再返回 *schema.Message | 批量处理、非实时场景 |
| 流式 | model.Stream(ctx, messages) | 返回 *schema.StreamReader,逐块接收 | 聊天、实时显示 |
示范代码
2.1 用模板构造消息
import ( "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/schema")
// 创建模板:用 {变量名} 做占位符func createTemplate() prompt.ChatTemplate { return prompt.FromMessages(schema.FString, schema.SystemMessage("你是一个{role}。用{style}的语气回答。"), schema.MessagesPlaceholder("chat_history", true), schema.UserMessage("问题: {question}"), )}
// 用模板生成实际消息func createMessages() []*schema.Message { // 正如上面所说,这里使用adk.Message也是一样的 template := createTemplate() messages, err := template.Format(context.Background(), map[string]any{ "role": "程序员鼓励师", "style": "积极、温暖且专业", "question": "我的代码一直报错,该怎么办?", "chat_history": []*schema.Message{ schema.UserMessage("你好"), schema.AssistantMessage("嘿!加油!", nil), }, }) // ... 处理 err return messages}💡
components/prompt/chat_prompt/chat_prompt.go展示了一个更完整的模板用法。
2.2 一次性生成(Generate)
result, err := model.Generate(ctx, messages)if err != nil { log.Fatal(err)}fmt.Println(result.Content)2.3 流式生成(Stream)—— 重点掌握
streamReader, err := model.Stream(ctx, messages)if err != nil { log.Fatal(err)}defer streamReader.Close() // ⚠️ 用完记得关闭!
for { chunk, err := streamReader.Recv() if err == io.EOF { break } if err != nil { log.Fatal(err) } fmt.Print(chunk.Content)}💡 StreamReader 是什么?
model.Stream()返回的*schema.StreamReader[*schema.Message]和 Iterator 类似——模型在后台逐块生成,你通过Recv()逐块消费。但有关键区别:
- Iterator(Agent 用):
Next()返回(event, bool),!ok表示 Agent 执行完毕- StreamReader(ChatModel 用):
Recv()返回(chunk, error),io.EOF表示流结束defer Close()是必须的:StreamReader 底层持有 HTTP 连接,不 Close 会导致连接泄漏 / goroutine 泄漏
⚠️ 本节最容易搞混的地方
Eino 里有两种迭代模式,结束判断条件不同:
// ─── 模式A:Agent Runner 事件迭代 ───// Next() 返回 (event, bool),!ok 表示结束for { event, ok := iter.Next() if !ok { break }}
// ─── 模式B:ChatModel Stream 流式迭代 ───// Recv() 返回 (chunk, error),io.EOF 表示结束defer streamReader.Close()for { chunk, err := streamReader.Recv() if err == io.EOF { break }}记忆诀窍:Agent 用 Next + ok,Stream 用 Recv + EOF。
常用 import 路径(本节新增):
"github.com/cloudwego/eino/components/prompt" // ChatTemplate, FromMessages📇 本节速记卡片
一次性:model.Generate(ctx, messages) → (*Message, error)流式: model.Stream(ctx, messages) → (*StreamReader, error) defer close → Recv() 循环 → io.EOF 结束
模板: prompt.FromMessages(schema.FString, ...messages) template.Format(ctx, map[string]any{...})三、ChatModelAgent 基础:Tool + Agent + Runner
📁 代码:
adk/intro/chatmodel/chatmodel.go📖 理论:
quickstart/chatwitheino/docs/ch02_chatmodel_agent_runner_console
🆕 本节新包速览(先记名字,后面会细讲):
"github.com/cloudwego/eino/components/tool" // BaseTool, InvokableTool"github.com/cloudwego/eino/components/tool/utils" // InferTool, NewTool"github.com/cloudwego/eino/compose" // ToolsNodeConfig(Tool 配置容器)
compose包会在第七节深入讲解——这里你只需要知道compose.ToolsNodeConfig是注册 Tool 时用的”配置壳”即可。
3.1 为什么需要 Agent 而不是直接用 ChatModel?
| 维度 | ChatModel(组件) | ChatModelAgent(智能体) |
|---|---|---|
| 定位 | 单个能力单元 | 完整的 AI 应用 |
| 输出 | Generate() / Stream() 直接返回消息 | Run() 返回事件流,可包含工具调用、中断等 |
| 多轮对话 | 需要自己管理 history | Agent 内部处理 |
| 工具调用 | ❌ 不支持 | ✅ 通过 ToolsConfig 配置 |
| 中断恢复 | ❌ 不支持 | ✅ 通过 CheckPointStore |
| 适用场景 | 简单的一次性问答 | 复杂智能体应用 |
3.2 创建 Tool 的三种方式
📁 补充代码:
quickstart/todoagent/main.go
Eino 提供三种创建 Tool 的方式:
| 方式 | 做法 | 适合场景 |
|---|---|---|
| 方式一:InferTool | 写一个普通函数,自动从 struct tag 推断 schema | 最常用,90% 场景 |
| 方式二:NewTool | 手动写 schema.ToolInfo,传给 utils.NewTool | 需要精确控制参数约束 |
| 方式三:结构体实现接口 | struct 实现 Info() + InvokableRun() | Tool 逻辑复杂,需要带状态 |
方式一:InferTool(自动推断,最常用)⭐
import ( "github.com/cloudwego/eino/components/tool" "github.com/cloudwego/eino/components/tool/utils")
// ① 定义输入结构体(用 json tag + jsonschema tag 描述参数)type BookSearchInput struct { Genre string `json:"genre" jsonschema_description:"书籍类型"` MaxPages int `json:"max_pages" jsonschema_description:"最大页数"`}
type BookSearchOutput struct { Books []string `json:"books"`}
// ② 写业务函数 + utils.InferTool 自动包装func NewBookSearchTool() tool.InvokableTool { t, _ := utils.InferTool( "search_book", "Search books by preferences", func(ctx context.Context, input *BookSearchInput) (*BookSearchOutput, error) { return &BookSearchOutput{Books: []string{"《三体》"}}, nil }, ) return t}方式二:NewTool(手动定义 schema)
当需要更精确控制参数描述、添加枚举约束等场景时使用:
func getAddTodoTool() tool.InvokableTool { info := &schema.ToolInfo{ Name: "add_todo", Desc: "Add a todo item", ParamsOneOf: schema.NewParamsOneOfByParams(map[string]*schema.ParameterInfo{ "content": { Desc: "The content of the todo item", Type: schema.String, Required: true, }, "deadline": { Desc: "The deadline of the todo item, in unix timestamp", Type: schema.Integer, }, }), } return utils.NewTool(info, AddTodoFunc)}
type TodoAddParams struct { Content string `json:"content"` Deadline *int64 `json:"deadline,omitempty"`}
func AddTodoFunc(_ context.Context, params *TodoAddParams) (string, error) { return `{"msg": "add todo success"}`, nil}方式三:结构体实现接口(最灵活)
Tool 需要带状态、或者逻辑很复杂时,直接实现接口:
type ListTodoTool struct{}
func (lt *ListTodoTool) Info(_ context.Context) (*schema.ToolInfo, error) { return &schema.ToolInfo{ Name: "list_todo", Desc: "List all todo items", ParamsOneOf: schema.NewParamsOneOfByParams(map[string]*schema.ParameterInfo{ "finished": { Desc: "filter todo items if finished", Type: schema.Boolean, }, }), }, nil}
func (lt *ListTodoTool) InvokableRun(_ context.Context, argumentsInJSON string, _ ...tool.Option) (string, error) { return `{"todos": [...]}`, nil}🧠 怎么选:90% 场景用
InferTool就够了;需要手动约束参数类型/枚举值时用NewTool;Tool 本身带状态(比如有数据库连接)时用结构体方式。三种方式注册到 Agent 的方式完全一样——都丢进ToolsConfig.Tools数组里。
3.3 组装完整应用(Tool + 流式)
ℹ️ 本节聚焦:只展示 Tool + Agent + Runner + Query 的主流程。中断恢复(Resume)会单独在第五节讲。
package main
import ( "context" "fmt" "log" "os"
"github.com/cloudwego/eino/adk" "github.com/cloudwego/eino/components/tool" "github.com/cloudwego/eino/compose" "github.com/cloudwego/eino/schema" "github.com/cloudwego/eino-ext/components/model/openai")
func main() { ctx := context.Background()
// ─── 创建模型 ─── model, err := openai.NewChatModel(ctx, &openai.ChatModelConfig{ APIKey: os.Getenv("OPENAI_API_KEY"), Model: os.Getenv("OPENAI_MODEL"), BaseURL: os.Getenv("OPENAI_BASE_URL"), }) if err != nil { log.Fatal(err) }
// ─── 创建 Agent(配置 Tool)─── agent, err := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{ Name: "BookRecommender", Description: "An agent that recommends books", Instruction: "You are a book expert. Use search_book tool to find books.", Model: model, ToolsConfig: adk.ToolsConfig{ ToolsNodeConfig: compose.ToolsNodeConfig{ Tools: []tool.BaseTool{ NewBookSearchTool(), }, }, }, }) if err != nil { log.Fatal(err) }
// ─── 创建 Runner(配置流式)─── runner := adk.NewRunner(ctx, adk.RunnerConfig{ Agent: agent, EnableStreaming: true, })
// ─── 运行 ─── iter := runner.Query(ctx, "推荐一本科幻小说") for { event, ok := iter.Next() if !ok { break } if event.Err != nil { log.Fatal(event.Err) } fmt.Printf("Event: %+v\n", event) }}3.4 配置项归属速查
记住一句话:Agent 管”有什么能力”,Runner 管”怎么执行”。
| 配置 | 在哪里设置 | 作用 |
|---|---|---|
Instruction | ChatModelAgentConfig | 系统提示词,定义 Agent 的”人设” |
ToolsConfig | ChatModelAgentConfig | 注册 Agent 可用的工具 |
EnableStreaming: true | RunnerConfig | 流式输出(打字机效果) |
CheckPointStore | RunnerConfig | 保存执行状态,支持中断恢复(第五节) |
3.5 底层视角:Agent 的本质是一个接口
📁 代码:
adk/intro/custom/myagent.go
adk.NewChatModelAgent 内部实现了 adk.Agent 接口:
type Agent interface { Name(ctx) string Description(ctx) string Run(ctx, input, ...options) *AsyncIterator[*AgentEvent]}myagent.go 展示了最简实现——三个关键点:
NewAsyncIteratorPair= 建一条事件流水线,gen管发送,iter管接收- 开 goroutine 是因为
Run()要立即返回iter,事件在后台异步生成 gen.Close()必须调用,否则消费端的Next()会永远阻塞
常用 import 路径(本节新增):
"github.com/cloudwego/eino/components/tool" // BaseTool, InvokableTool, Option"github.com/cloudwego/eino/components/tool/utils" // InferTool, NewTool"github.com/cloudwego/eino/compose" // ToolsNodeConfig📇 本节速记卡片
Agent.{Name, Instruction, ToolsConfig, Model}Runner.{Agent, EnableStreaming, CheckPointStore}
Tool 三种创建:InferTool(90%)/ NewTool / 结构体接口注册方式都一样:ToolsConfig.Tools = []tool.BaseTool{ toolfuction...}
附注:如果忘记某特定的struct需要填写什么可以直接看源码接口定义四、Middleware:给 Agent 加拦截器
📁 代码:
quickstart/chatwitheino/helpers/middleware.go- SafeToolMiddleware,自定义一个新的middlewarequickstart/chatwitheino/helpers/retry.go- 错误自动重试adk/middlewares/dynamictool/toolsearch/— 动态工具检索中间件adk/middlewares/skill/main.go— 技能加载中间件📖 理论:
quickstart/chatwitheino/docs/ch05_middleware
4.1 为什么需要 Middleware?
两个常见问题:
问题一:Tool 报错导致整个对话崩溃
[tool call] read_file(file_path: "nonexistent.txt")Error: open nonexistent.txt: no such file or directory// 💥 对话直接中断问题二:模型 API 限流导致失败
Error: rate limit exceeded (429)// 💥 对话中断你希望的行为是:Tool 报错时把错误交给 LLM 自愈;模型限流时自动重试。这就是 Middleware 要解决的问题——Agent 的拦截器,在调用前后插入自定义逻辑。
4.2 Middleware 的本质:洋葱模型
请求 → A.Wrap → B.Wrap → C.Wrap → 实际执行 → C返回 → B返回 → A返回 → 响应 ↑ ↑ 最外层先拦截 最内层先接触实际执行结果4.3 核心场景一:SafeToolMiddleware(自定义一个新的middleware:Tool 错误转字符串)
// 定义一个结构体来承载中间件的逻辑。内嵌一个基础中间件(Base Middleware),以复用其默认行为。type safeToolMiddleware struct { *adk.BaseChatModelAgentMiddleware}
// 拦截同步工具调用,BaseChatModelAgentMiddleware中定义了方法的签名func (m *safeToolMiddleware) WrapInvokableToolCall( _ context.Context, endpoint adk.InvokableToolCallEndpoint, _ *adk.ToolContext,) (adk.InvokableToolCallEndpoint, error) { return func(ctx context.Context, args string, opts ...tool.Option) (string, error) { // 拦截原始 endpoint的结果和err result, err := endpoint(ctx, args, opts...) if err != nil { // ⚠️ 中断错误必须继续传播,不能吞掉 if _, ok := compose.IsInterruptRerunError(err); ok { return "", err } // 普通错误 → 转成字符串,让 LLM 看到 return fmt.Sprintf("[tool error] %v", err), nil } // 没有错误就直接返回原始的结果 return result, nil }, nil}
// 还可以拦截其他的调用如WrapStreamableToolCall 拦截流式工具调用等等,可以自行查看官方示例💡 关键区分:
compose.IsInterruptRerunError是中断恢复(第五节)抛出的特殊错误,它必须继续向上传播,不能转成字符串。
4.4 核心场景二:ModelRetryConfig(模型调用自动重试)
原来的IsRetryAble被标记弃用,现在应该使用ShouldRetry进行替代
部分中间件可以直接在agent定义时注明,如ModelRetryConfig
ModelRetryConfig: &adk.ModelRetryConfig{ MaxRetries: 5, ShouldRetry: func(ctx context.Context, retryCtx *adk.TypedRetryContext[M]) *adk.TypedRetryDecision[M] { // 获取错误信息 err := retryCtx.Error // 判断是否需要重试 if err != nil && (strings.Contains(err.Error(), "429") || strings.Contains(err.Error(), "Too Many Requests")) { // 返回重试决策 return &adk.TypedRetryDecision[M]{ Retry: true, // 在这里你还可以修改重试时的输入或选项,等等 } } // 不需要重试,返回 nil 或 Retry: false return nil },},4.5 核心场景三:动态工具检索(ToolSearch Middleware)
eino还预先定义了部分middleware,如agentsmd、dynamictool/toolsearch等,可以直接调用,下面介绍一些常见的
📁
adk/middlewares/dynamictool/toolsearch/
痛点:Tool 太多会撑爆上下文窗口。ToolSearch 先把所有 Tool 放进”工具库”,调用时搜索过滤:
toolSearchMiddleware, _ := toolsearch.New(ctx, &toolsearch.Config{ DynamicTools: allDynamicTools,})
agent, _ := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{ Name: "tool_search_agent", Model: chatModel, Handlers: []adk.ChatModelAgentMiddleware{ toolSearchMiddleware, // ← 不是 ToolsConfig! },})🧠 与第三节 ToolsConfig 的区别:第三节把 Tool 写在
ToolsConfig.Tools里——全部发给 LLM。ToolSearch 先用搜索过滤,只把相关的发给 LLM。
4.6 核心场景四:Skill Middleware(动态加载技能)
📁
adk/middlewares/skill/main.go
创建 Skill 后端
skillBackend, err := skill.NewBackendFromFilesystem(ctx, &skill.BackendFromFilesystemConfig{ Backend: be, BaseDir: skillsDir, // 指定扫描的根目录})创建并注入 Skill 中间件
sm, err := skill.NewMiddleware(ctx, &skill.Config{ Backend: skillBackend,})之后在运行时从文件系统动态加载”技能”(预定义的提示词+工具组合)
// 文件系统中间件 → Skill 中间件(有序!)agent, _ := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{ Name: "LogAnalysisAgent", Model: cm, Handlers: []adk.ChatModelAgentMiddleware{ fsm, // 文件读写能力(洋葱模型外层,最先wrap,最后返回) sm, // Skill 加载能力(洋葱模型内层,之后wrap,先返回) },})4.7 Middleware 全景速查
| Middleware 类型 | 解决什么痛点 | 来源 | 配置位置 |
|---|---|---|---|
| SafeToolMiddleware | Tool 报错中断对话 → 转成字符串让 LLM 自愈 | 自己实现 | Handlers |
| ModelRetryConfig | 模型 API 限流 → 自动重试 | adk.ModelRetryConfig(内置) | ChatModelAgentConfig |
| ToolSearch | Tool 太多撑爆上下文 → 动态检索注入 | toolsearch.New() | Handlers |
| Skill | 运行时动态加载技能 → 从文件系统读 | skill.NewMiddleware() | Handlers |
| filesystem | Agent 需要读写文件系统 | filesystem.New() | Handlers |
🧠 模型:安检通道
用户请求 → [安检A:查证件] → [安检B:过X光] → [安检C:金属探测] → 登机用户收到 ← [安检A:盖章] ← [安检B:贴标签] ← [安检C:称重] ← 行李出来📇 本节速记卡片
SafeToolMiddleware = 错误转字符串(除了 InterruptRerunError)ModelRetryConfig = 限流自动重试(内置配置项)ToolSearch = 工具池动态搜索(替代 ToolsConfig)Skill = 文件系统加载技能
通用模式:嵌入 BaseChatModelAgentMiddleware → 覆写一个钩子 → 放入 Handlers五、中断与恢复(Interrupt & Resume)
📁 代码:
adk/intro/chatmodel/— Tool 中断adk/cancel/graceful-exit/— 外部取消adk/human-in-the-loop/— 人机协作全套示例📖 Memory 理论:
quickstart/chatwitheino/docs/ch03_memory_session_jsonl📖 Tool 理论:
quickstart/chatwitheino/docs/ch04_tool_backend_filesystem
5.1 中断机制总览
这是 Eino ADK 的核心高级特性。整个机制基于一个简单的状态机:
stateDiagram-v2
[*] --> Running: runner.Query/Run(checkpointID)
Running --> Interrupted: compose.Interrupt()
Running --> Cancelled: Ctrl-C → cancelFn()
Running --> Completed: 正常结束
Running --> HITL_Interrupted: compose.Interrupt()
Interrupted --> Resuming: runner.Resume(checkpointID, WithToolOptions)
Cancelled --> ResumingNew: 新runner.Resume(checkpointID)
HITL_Interrupted --> ResumingHITL: runner.ResumeWithParams(checkpointID, params)
Resuming --> Running
ResumingNew --> Running
ResumingHITL --> Running
Completed --> [*]三种触发方式对比如下:
| 方式一:Tool 中断(本质就是方式三) | 方式二:外部取消 | 方式三:Human-in-the-Loop | |
|---|---|---|---|
| 触发源 | Tool 调用 `compose.Interrupt“ | 系统信号 Ctrl-C | Tool 调用 compose.Interrupt |
| Runner | 同一个 runner | 新建 runner | 同一个 runner |
| Resume API | runner.Resume + WithToolOptions | 新 runner.Resume | runner.ResumeWithParams |
| 典型场景 | 信息不足,等用户补充 | 长时间运行,用户想中止 | 敏感操作审批、参数审核 |
共同点:都依赖
CheckPointStore+CheckPointID。
🆕 本节新包速览:
"github.com/cloudwego/eino/compose" // CheckPointStore(中断状态存储接口)// CheckPointStore 有多种实现,如:// store.NewInMemoryStore() ← 内存实现(开发用)// 第五节只关注中断API,CheckPointStore 的底层原理见第七节
5.2 方式一:Tool 内部中断(主动暂停等输入)
用户输入 "推荐一本书" │ ▼ runner := adk.NewRunner(ctx, adk.RunnerConfig{ Agent: a, CheckPointStore: store.NewInMemoryStore(), // runner设置CheckPointStore以储存后续的CheckPointID }) │ ▼ runner.Query(ctx, "推荐一本书", WithCheckPointID("1")) 传入CheckPointID,后续恢复通过这个 │ ▼ Agent 分析:信息不够 → Tool 返回 Interrupt("请问你喜欢什么类型的书?") 注:示例代码中NewInterruptAndRerunErr被标记启用 │ ▼ 事件流中收到 Interrupted 事件 → 暂停,等待用户输入 │ ▼ runner.Resume(ctx, "1", WithToolOptions(WithNewInput("科幻"))) 通过CheckPointID进行恢复 │ ▼ Agent 继续执行 → 用 "科幻" 重新调用 Tool → 返回结果注:既然这个已经过时+本质和方法三相同,可以直接看方法三是如何中断的
5.3 方式二:外部取消(Ctrl-C 优雅退出)
📁
adk/cancel/graceful-exit/main.go
设计意图:用户按 Ctrl-C 时不直接杀进程白跑——而是安全暂停、自动保存进度、之后从断点恢复。
用户按 Ctrl-C │ ▼ cancelFn(CancelAfterChatModel, WithRecursive, 30s超时) │ ├─→ CancelAfterChatModel: 等当前 ChatModel 调用完成再停 ├─→ WithRecursive: 递归传播到子 Agent └─→ 超时兜底: 30s 内没到安全点 → 强制 CancelImmediate │ ▼ CheckPointStore 自动保存状态 → 新 Runner.Resume 恢复取消的创建
// 创建取消选项与外部触发函数cancelOpt, cancelFn := adk.WithCancel()
// 启动Agent运行,绑定断点IDiter := runner.Run(ctx, input, cancelOpt, adk.WithCheckPointID(checkpointID))
// 注册系统信号监听:SIGINT(Ctrl-C)、SIGTERMsigCh := make(chan os.Signal, 1)signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
go func() { sig := <-sigCh fmt.Printf("\n收到操作系统信号: %v,发起Agent优雅取消\n", sig)
// 发起优雅取消:等待ChatModel安全点、递归传播子Agent、30s超时兜底 handle, contributed := cancelFn( adk.WithAgentCancelMode(adk.CancelAfterChatModel), adk.WithRecursive(), adk.WithAgentCancelTimeout(30*time.Second), ) fmt.Printf("取消请求是否成功提交: %v\n", contributed)
// 阻塞等待取消流程完成(断点持久化完成后返回) if waitErr := handle.Wait(); waitErr != nil { fmt.Printf("优雅取消异常: %v\n", waitErr) } else { fmt.Println("优雅取消完成,断点已保存") }}()取消的恢复
// 使用同一个CheckPointStore,新建Runner恢复resumeRunner := adk.NewRunner(ctx, adk.RunnerConfig{ Agent: agent, EnableStreaming: true, CheckPointStore: cpStore,})resumeIter, err := resumeRunner.Resume(ctx, checkpointID)if err != nil { log.Fatalf("断点恢复失败: %v", err)}drainEvents(resumeIter)5.4 方式三:Human-in-the-Loop 人机协作(ResumeWithParams)
📁
adk/human-in-the-loop/1_approval/~8_supervisor-plan-execute/
通过 interruptID 精确匹配中断点,注入任意类型的数据:
// ① 启动iter := runner.Query(ctx, query, adk.WithCheckPointID("1"))
// ② 检测中断,提取 interruptIDinterruptCtx := lastEvent.Action.Interrupted.InterruptContexts[0]interruptID := interruptCtx.ID
// ③ 用户决策后,通过 ResumeWithParams 传回iter, _ = runner.ResumeWithParams(ctx, "1", &adk.ResumeParams{ Targets: map[string]any{ interruptID: modifiedInfo, },})中断方:
func FollowUp(ctx context.Context, input *FollowUpToolInput) (string, error) { // 首次进入(触发中断) wasInterrupted, _, storedState := tool.GetInterruptState[*FollowUpState](ctx) if !wasInterrupted { // 准备要展示给用户的信息(比如问题列表) info := &FollowUpInfo{Questions: input.Questions} // 准备要保存的状态(为了恢复时知道我们在干嘛) state := &FollowUpState{Questions: input.Questions}
// 触发中断!把 info 发给前端展示,把 state 存起来,然后函数立刻返回 return "", tool.StatefulInterrupt(ctx, info, state) }
// 第二次进入 isResumeTarget, hasData, resumeData := tool.GetResumeContext[*FollowUpInfo](ctx) if !isResumeTarget { // 如果系统唤醒了这个函数,但发现它不是当前该恢复的目标 // 那就把之前保存的状态再拿出来,继续中断等待 info := &FollowUpInfo{Questions: storedState.Questions} return "", tool.StatefulInterrupt(ctx, info, storedState) } // 被唤醒,但还没拿到答案(防御性重入) if !hasData || resumeData.UserAnswer == "" { return "", fmt.Errorf("tool resumed without a user answer") }
// 成功恢复,拿到答案(正常退出) return resumeData.UserAnswer, nil}四种交互模式(区别仅在于 info 类型和用户操作):
| 模式 | 中断时携带的信息 | 用户操作 | 示例目录 |
|---|---|---|---|
| 审批 | ApprovalResult | Y/N 批准或拒绝 | 1_approval, 5_supervisor |
| 审核编辑 | ReviewEditInfo | 修改参数 / 批准 / 拒绝 | 2_review-and-edit |
| 反馈循环 | FeedbackInfo | 输入反馈意见 | 3_feedback-loop |
| 追问 | FollowUpInfo | 逐个回答问题 | 4_follow-up, 7_deep-agents |
💡 ResumeWithParams vs 旧 Resume:
ResumeWithParams是通用 API,Resume+WithToolOptions是其特例。
📇 本节速记卡片
中断恢复三要素:CheckPointStore(存状态)+ CheckPointID(标识)+ Resume(恢复)
方式一:compose.Interrupt → runner.Resume(id, WithToolOptions)方式二:Ctrl-C → cancelFn → 新 runner.Resume(id)方式三:compose.Interrupt → runner.ResumeWithParams(id, {Targets: {interruptID: data}})六、进阶主题
6.1 Callbacks 与 Tracing(可观测性)
📁
adk/common/trace/coze_loop.go
traceCloseFn, startSpanFn := trace.AppendCozeLoopCallbackIfConfigured(ctx)defer traceCloseFn(ctx)
ctx, endSpanFn := startSpanFn(ctx, "MyTask", "user query")// ... runner.Query(ctx, query) ...endSpanFn(ctx, lastMessage)6.2 Workflow Agent:编排多个 Agent 的协作方式
📁
adk/intro/workflow/loop/,parallel/,sequential/
三种内置工作流 Agent:
| 类型 | 模式 | 适用场景 |
|---|---|---|
| LoopAgent | 主Agent生成 → 评审审查 → 不满意重改 | 需要反复打磨的任务 |
| ParallelAgent | 多个子Agent同时接收相同输入,各自执行 | 同时从多数据源收集信息 |
| SequentialAgent | 子Agent逐个执行,前一个输出→后一个输入 | 分步任务流水线 |
// LoopAgenta, _ := adk.NewLoopAgent(ctx, &adk.LoopAgentConfig{ SubAgents: []adk.Agent{mainAgent, critiqueAgent}, MaxIterations: 5,})
// ParallelAgenta, _ := adk.NewParallelAgent(ctx, &adk.ParallelAgentConfig{ SubAgents: []adk.Agent{stockAgent, newsAgent, socialAgent},})
// SequentialAgenta, _ := adk.NewSequentialAgent(ctx, &adk.SequentialAgentConfig{ SubAgents: []adk.Agent{planAgent, writerAgent},})6.3 Multi-Agent:Supervisor 模式
📁
adk/multiagent/supervisor/,layered-supervisor/,plan-execute-replan/
Supervisor:一个”主管”Agent 分析用户请求,决定派给哪个”专家”子 Agent。
sv, _ := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{ Name: "supervisor", Instruction: "You are a supervisor managing a research_agent and math_agent.", Model: m,})
supervisorAgent, _ := supervisor.New(ctx, &supervisor.Config{ Supervisor: sv, SubAgents: []adk.Agent{searchAgent, mathAgent},})Layered Supervisor:子 Agent 本身也可以是 Supervisor,形成树状调度。
Plan-Execute-Replan:Planner 制定计划 → Executor 逐步执行 → Replanner 检查并修改计划。
entryAgent, _ := planexecute.New(ctx, &planexecute.Config{ Planner: planAgent, Executor: executeAgent, Replanner: replanAgent, MaxIterations: 20,})🧠 三种多 Agent 模式对比:
- Workflow Agent(6.2):编排是固定的(顺序/并行/循环)
- Supervisor(6.3):编排是动态的,由 LLM 决定派给谁
- Plan-Execute-Replan:适合需要多步规划且计划可能中途调整的复杂任务
6.4 Compose 编排 API 总览
按抽象层级从高到低:
ADK (Agent+Runner) > Chain > Graph > Workflow 高层,开箱即用 底层,精细控制每个节点和数据流| API | 特点 | 适合场景 |
|---|---|---|
| Workflow | 字段映射 + 自动依赖推导 | 纯函数数据流水线 |
| Graph | 手动 AddEdge 建图,支持 State | LLM 调用流程、有状态分支 |
| Chain | 链式 Append 节点 | ChatModel → ToolsNode 线性流程 |
| Batch | 批量并行处理 N 个输入 | 文档批量审核等 |
Graph 基础用法
g := compose.NewGraph[map[string]any, *schema.Message]()_ = g.AddChatTemplateNode("prompt", pt)_ = g.AddChatModelNode("model", chatModel)_ = g.AddEdge(compose.START, "prompt")_ = g.AddEdge("prompt", "model")_ = g.AddEdge("model", compose.END)runner, _ := g.Compile(ctx)result, _ := runner.Invoke(ctx, input)Chain 基础用法
chain := compose.NewChain[map[string]any, string]()chain. AppendLambda(preprocess). AppendBranch(compose.NewChainBranch(cond).AddLambda("b1", fn1).AddLambda("b2", fn2)). AppendLambda(postprocess)runner, _ := chain.Compile(ctx)📇 本节速记卡片
Workflow Agent: LoopAgent / ParallelAgent / SequentialAgentSupervisor: supervisor.New({Supervisor, SubAgents})Compose 编排: Graph(AddEdge) / Chain(Append) / Workflow(MapFields)七、完整架构分层图
graph TD
Q["你要做什么?"] --> Q1{"简单一次性问答?"}
Q1 -->|是| CM["ChatModel.Generate/Stream
§2"]
Q1 -->|否| Q2{"标准 Agent
问答+工具+记忆?"}
Q2 -->|是| ADK["ChatModelAgent + Runner
§3"]
Q2 -->|否| Q3{"自定义 ReAct 循环?"}
Q3 -->|是| REACT["react.NewAgent
appendix1 §1"]
Q3 -->|否| Q4{"多 Agent 协作?"}
Q4 -->|调度模式| SUP["ADK Supervisor
§6.3
或 Host MultiAgent
appendix1 §6"]
Q4 -->|流水线| PE["Plan-Execute
appendix1 §7"]
Q4 -->|复杂状态图| CUSTOM["自定义 Graph+State
appendix1 §9"]
Q4 -->|否| Q5{"精细控制每个节点?"}
Q5 -->|是| GRAPH["Compose Graph/Chain
§6.4"]
Q5 -->|否| CHECK["回到 ADK 路径"]📇 全篇速记卡片(完整版)
┌─ 创建链 ─────────────────────────────────────────────┐│ NewChatModel → NewChatModelAgent → NewRunner → Run ││ (模) (代) (跑) (消息) │├─ 两种迭代 ───────────────────────────────────────────┤│ Agent: for { event, ok := iter.Next(); if !ok {break}}││ Stream: for { chunk, err := sr.Recv(); io.EOF → break}│├─ 职责分离 ───────────────────────────────────────────┤│ Agent 管能力:{Name, Instruction, ToolsConfig, Model} ││ Runner管执行:{Agent, EnableStreaming, CheckPointStore}│├─ Tool 创建 ──────────────────────────────────────────┤│ InferTool(90%)/ NewTool / 结构体接口 │├─ Middleware ────────────────────────────────────────┤│ SafeTool / ModelRetry / ToolSearch / Skill │├─ 中断恢复 ──────────────────────────────────────────┤│ CheckPointStore + CheckPointID + Resume ││ 三种方式:Tool中断 / Ctrl-C取消 / Human-in-the-Loop │├─ 编排层次 ──────────────────────────────────────────┤│ ADK(Agent+Runner) > Chain > Graph > Workflow │├─ Agentic 路径([Agentic 进阶](./agentic))──────────────────────────┤│ 所有类型加 Typed[*AgenticMessage] 即可 │└──────────────────────────────────────────────────────┘常用 import 路径(全篇汇总):
// 模型"github.com/cloudwego/eino-ext/components/model/openai" // OpenAI 兼容 ChatModel
// ADK 核心"github.com/cloudwego/eino/adk" // Agent, Runner, Message, Middleware"github.com/cloudwego/eino/schema" // UserMessage, SystemMessage, ToolInfo
// Tool"github.com/cloudwego/eino/components/tool" // BaseTool, InvokableTool"github.com/cloudwego/eino/components/tool/utils" // InferTool, NewTool
// 编排"github.com/cloudwego/eino/compose" // Graph, Chain, Workflow, ToolsNodeConfig
// Middleware 扩展"github.com/cloudwego/eino/adk/middlewares/dynamictool/toolsearch" // ToolSearch"github.com/cloudwego/eino/adk/middlewares/skill" // Skill"github.com/cloudwego/eino/adk/middlewares/filesystem" // 文件系统
// Prompt"github.com/cloudwego/eino/components/prompt" // ChatTemplate, FromMessages⚠️ 常见坑(Common Gotchas)
开发 Eino 应用时,下面几个坑最容易踩到——每个都用一个”坏代码 vs 好代码”的对比来展示:
坑 1:忘记 defer streamReader.Close() → goroutine 泄漏
// ❌ 坏:忘记 Close,HTTP 连接永不释放,goroutine 泄漏sr, _ := model.Stream(ctx, messages)for { chunk, err := sr.Recv(); ... }
// ✅ 好:defer Close 保证无论怎么退出都会关闭sr, _ := model.Stream(ctx, messages)defer sr.Close()for { chunk, err := sr.Recv(); ... }坑 2:自定义 Agent 忘记 gen.Close() → Next() 永远阻塞
// ❌ 坏:忘记 Close,消费端 Next() 永远等不到 falsefunc (m *MyAgent) Run(...) *AsyncIterator[*AgentEvent] { iter, gen := adk.NewAsyncIteratorPair[*AgentEvent]() go func() { gen.Send(&AgentEvent{...}) // 忘记 gen.Close()!!! }() return iter}
// ✅ 好:defer + recover 确保 Close 一定被调用func (m *MyAgent) Run(...) *AsyncIterator[*AgentEvent] { iter, gen := adk.NewAsyncIteratorPair[*AgentEvent]() go func() { defer func() { if e := recover() { gen.Send(&AgentEvent{Err: ...}) } gen.Close() // ← 必须调用! }() gen.Send(&AgentEvent{...}) }() return iter}坑 3:SafeToolMiddleware 误吞中断错误
// ❌ 坏:所有 error 都转成字符串,中断错误也被吞了func wrap(endpoint) endpoint { return func(ctx, args) (string, error) { result, err := endpoint(ctx, args) if err != nil { return fmt.Sprintf("[error] %v", err), nil } return result, nil }}
// ✅ 好:区分中断错误和普通错误func wrap(endpoint) endpoint { return func(ctx, args) (string, error) { result, err := endpoint(ctx, args) if err != nil { if _, ok := compose.IsInterruptRerunError(err); ok { return "", err // ← 中断错误必须传播! } return fmt.Sprintf("[tool error] %v", err), nil } return result, nil }}坑 4:Resume 时 CheckPointID 不一致
// ❌ 坏:Query 和 Resume 用了不同的 ID——Runner 找不到之前的保存点iter := runner.Query(ctx, "帮我推荐", adk.WithCheckPointID("session-1"))// ... 中断 ...iter, _ = runner.Resume(ctx, "session-2") // ← ID 不一致!应该是 "session-1"
// ✅ 好:同一个 IDiter := runner.Query(ctx, "帮我推荐", adk.WithCheckPointID("session-1"))// ... 中断 ...iter, _ = runner.Resume(ctx, "session-1") // ← 一致!坑 5:Resume 时用了新的 Runner 但没传同一个 CheckPointStore
// ❌ 坏:新 Runner 没用同一个 CheckPointStore——状态丢失runner1 := adk.NewRunner(ctx, adk.RunnerConfig{Agent: agent}) // 默认空 Storeiter := runner1.Query(ctx, "query", adk.WithCheckPointID("1"))// ... Ctrl-C 取消 ...runner2 := adk.NewRunner(ctx, adk.RunnerConfig{Agent: agent}) // 另一个空 Storeiter, _ = runner2.Resume(ctx, "1") // ← 找不到 checkpoint!
// ✅ 好:共享同一个 CheckPointStorestore := store.NewInMemoryStore()runner1 := adk.NewRunner(ctx, adk.RunnerConfig{Agent: agent, CheckPointStore: store})runner1.Query(...)// ...runner2 := adk.NewRunner(ctx, adk.RunnerConfig{Agent: agent, CheckPointStore: store})runner2.Resume(ctx, "1") // ← 找得到!坑 6:服务端工具(web_search)期待 function_tool_result
// ❌ 坏:等不到 web_search 的 function_tool_resultfor _, block := range msg.ContentBlocks { if block.Type == schema.ContentBlockTypeFunctionToolResult { // web_search 永远不会产生这个! }}
// ✅ 好:理解 server_tool_call 在厂商侧执行,没有对应的 function_tool_resultfor _, block := range msg.ContentBlocks { if block.Type == schema.ContentBlockTypeServerToolCall { // 模型已经消费了搜索结果——你只会看到后续的 reasoning/text block fmt.Println("服务端搜索完成:", block.ServerToolCall.Name) }}🧠 总结:上面六个坑覆盖了 90% 的 Eino 开发问题。记住三条原则:
- 打开要关(StreamReader → defer Close,AsyncIterator → gen.Close)
- 中断要透传(InterruptRerunError 不能被 Middleware 吞掉)
- 恢复要对齐(CheckPointStore + CheckPointID 跨 Runner 必须一致)
📂 源码仓库:github.com/cloudwego/eino-examples↗
📖 继续阅读:
- 入门笔记 — Eino ADK 从 Hello World 到 Compose 编排 ← 你在这里
- Agentic 进阶 — Responses API / AgenticMessage / Typed 泛型
- 附录一:Flow 流程模块 — ReAct Agent / Multi-Agent / 状态图
- 附录二:Components 组件模块 — A/B 路由 / HTTP 日志 / 检索增强 / 文档解析
- 附录三:Lambda 与调试工具 — Lambda 写法 / Devops / Mermaid 可视化