Eino
Eino 是 CloudWeGo 的 Go 大模型应用开发框架。它把模型、提示模板、检索和工具等能力抽象为组件,再通过 Chain、Graph 等编排方式连接这些组件。本文整理组件接口和基础编排,示例使用 schema.Message,不展开 ADK 和 AgenticMessage 的用法。
本文以 Go 1.26.x 和 Eino v0.9.21 为基准,接口声明按该版本整理。eino-ext 的各组件是独立 Go 模块,版本号与 Eino 核心库不同。具体 API 可核对 Eino v0.9.21 文档。
1. 环境与依赖
创建独立的示例目录,在该目录初始化模块并安装依赖:
go mod init example.com/einodemo
go get github.com/cloudwego/eino@v0.9.21
go get github.com/cloudwego/eino-ext/components/model/ark@v0.1.71
go get github.com/cloudwego/eino-ext/components/embedding/ark@v0.1.2
go get github.com/cloudwego/eino-ext/components/indexer/redis@v0.0.0-20260924074145-3603a39473c3
go get github.com/cloudwego/eino-ext/components/retriever/redis@v0.0.0-20260924074145-3603a39473c3
go get github.com/cloudwego/eino-ext/components/document/transformer/splitter/markdown@v0.0.0-20260924074145-3603a39473c3
go get github.com/cloudwego/eino-ext/components/tool/httprequest@v0.0.0-20260924074145-3603a39473c3
go get github.com/cloudwego/eino-ext/callbacks/cozeloop@v0.3.1
go get github.com/coze-dev/cozeloop-go@v0.1.23
go get github.com/redis/go-redis/v9@v9.23.0其中 Ark ChatModel 使用 v0.1.71,Ark Embedding 使用 v0.1.2,CozeLoop 回调扩展使用 v0.3.1。Redis 索引、检索、Markdown 分割和 HTTP 工具模块使用同一个 Go 伪版本号,它们尚未提供可用于本例的语义版本标签。
每个包含 package main 的示例都是独立程序。例如,将生成示例保存为 cmd/generate/main.go,然后在模块根目录执行 go run ./cmd/generate。不同示例不能全部放进同一个包,否则 main 等名称会重复。接口和结构体摘录用于说明框架 API,其中的 Option 等类型属于相应组件包,不需要在业务代码里重新定义。
模型相关示例从进程环境读取配置,不会自动加载 .env。运行前设置相应变量:
| 变量 | 含义 |
|---|---|
ARK_API_KEY | Ark 服务的 API Key |
MODEL | 对话模型或推理接入点的 ID |
EMBEDDER | 嵌入模型或推理接入点的 ID |
ARK_BASE_URL | 可选,所用地域的 API 基础地址,留空使用 SDK 默认的北京地域地址 |
EMBEDDER_API_TYPE | 可选,text_api 对应 /embeddings,multi_modal_api 对应 /embeddings/multimodal,默认 text_api |
REDIS_ADDR | Redis 的 host:port,例如 127.0.0.1:6379 |
REDIS_PASSWORD | 可选,Redis 密码 |
MODEL 和 EMBEDDER 是这里约定的环境变量名,不是 Eino 自动识别的配置。应填写账号中可用的模型或接入点 ID,并选择该模型支持的接口类型。模板、文档分割、自定义工具和不含模型的编排示例无需 Ark 凭据。
2. ChatModel 组件
2.1 接口
ChatModel 抽象模型调用。BaseChatModel 是 BaseModel[*schema.Message] 的类型别名,包含 Generate 和 Stream。WithTools 属于扩展接口 ToolCallingChatModel,并非每个基础模型实现都支持它。
type messageType interface {
*schema.Message | *schema.AgenticMessage
}
type BaseModel[M messageType] interface {
Generate(ctx context.Context, input []M, opts ...Option) (M, error)
Stream(ctx context.Context, input []M, opts ...Option) (*schema.StreamReader[M], error)
}
type BaseChatModel = BaseModel[*schema.Message]
type ToolCallingChatModel interface {
BaseChatModel
WithTools(tools []*schema.ToolInfo) (ToolCallingChatModel, error)
}Generate 返回一条完整响应消息。Stream 返回流读取器,后续响应或错误从 Recv 获取。两者都接收消息列表、上下文和模型选项。上下文用于取消、超时和传递回调信息,选项的支持情况还取决于模型实现。
接口可以容纳多模态消息,但模型和服务适配器是否支持图片、音频等输入输出,需要分别确认。
2.2 schema.Message
type Message struct {
// Role 表示消息的角色(system/user/assistant/tool)
Role RoleType `json:"role"`
// Content 是消息的文本内容
Content string `json:"content"`
// if MultiContent is not empty, use this instead of Content
// if MultiContent is empty, use Content
// MultiContent 是多模态内容,支持文本、图片、音频等
// Deprecated: 用户输入用 UserInputMultiContent,模型输出用 AssistantGenMultiContent
MultiContent []ChatMessagePart `json:"multi_content,omitempty"`
// UserInputMultiContent 用来存储用户输入的多模态数据,支持文本、图片、音频、视频、文件
// 使用此字段时限制模型角色为 User
UserInputMultiContent []MessageInputPart `json:"user_input_multi_content,omitempty"`
// AssistantGenMultiContent 用来承接模型输出的多模态数据,支持文本、图片、音频、视频
// 使用此字段时限制模型角色为 Assistant
AssistantGenMultiContent []MessageOutputPart `json:"assistant_output_multi_content,omitempty"`
// Name 是消息的发送者名称
Name string `json:"name,omitempty"`
// ToolCalls 是 assistant 消息中的工具调用信息
ToolCalls []ToolCall `json:"tool_calls,omitempty"`
// ToolCallID 是 tool 消息的工具调用 ID
ToolCallID string `json:"tool_call_id,omitempty"`
// only for ToolMessage
ToolName string `json:"tool_name,omitempty"`
// ResponseMeta 包含响应的元信息
ResponseMeta *ResponseMeta `json:"response_meta,omitempty"`
// ReasoningContent 是模型服务返回的推理相关内容
ReasoningContent string `json:"reasoning_content,omitempty"`
// Extra 用于存储额外信息
Extra map[string]any `json:"extra,omitempty"`
}Role 区分系统指令、用户输入、模型回复和工具结果。用户问题应放在 user 消息中,不应因经过模板处理就改成 system 消息。ToolCalls 描述模型请求的工具调用,工具结果通过 ToolCallID 与对应调用关联。
MultiContent 已弃用。用户的多模态输入使用 UserInputMultiContent,模型的多模态输出使用 AssistantGenMultiContent。ReasoningContent 只在服务实际返回相应内容时有值,不能据此认为可以获得模型完整的内部推理过程。
2.3 完整生成
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino-ext/components/model/ark"
"github.com/cloudwego/eino/schema"
"log"
"os"
"time"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
apiKey, modelID := os.Getenv("ARK_API_KEY"), os.Getenv("MODEL")
if apiKey == "" || modelID == "" {
return fmt.Errorf("请设置 ARK_API_KEY 和 MODEL")
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
chatModel, err := ark.NewChatModel(ctx, &ark.ChatModelConfig{
APIKey: apiKey,
Model: modelID,
BaseURL: os.Getenv("ARK_BASE_URL"),
})
if err != nil {
return err
}
input := []*schema.Message{
schema.SystemMessage("你是一名技术助理,请用中文简要回答。"),
schema.UserMessage("什么是软件定义网络?"),
}
response, err := chatModel.Generate(ctx, input)
if err != nil {
return err
}
fmt.Println(response.Content)
return nil
}示例使用 30 秒调用超时,模型响应的具体内容不固定。返回错误时应先处理错误,再读取响应。
2.4 流式生成
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino-ext/components/model/ark"
"github.com/cloudwego/eino/schema"
"io"
"log"
"os"
"time"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
apiKey, modelID := os.Getenv("ARK_API_KEY"), os.Getenv("MODEL")
if apiKey == "" || modelID == "" {
return fmt.Errorf("请设置 ARK_API_KEY 和 MODEL")
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
chatModel, err := ark.NewChatModel(ctx, &ark.ChatModelConfig{
APIKey: apiKey,
Model: modelID,
BaseURL: os.Getenv("ARK_BASE_URL"),
})
if err != nil {
return err
}
input := []*schema.Message{
schema.SystemMessage("你是一名技术助理,请用中文简要回答。"),
schema.UserMessage("什么是软件定义网络?"),
}
reader, err := chatModel.Stream(ctx, input)
if err != nil {
return err
}
defer reader.Close()
for {
chunk, err := reader.Recv()
if err == io.EOF {
fmt.Println()
return nil
}
if err != nil {
return err
}
fmt.Print(chunk.Content)
}
}流块通常是增量内容,不是截至当前时刻的完整答案,也不保证按词或句子划分。只打印 Content 适合展示文本,工具调用参数、推理内容和多模态结果还可能位于其他字段中。
Recv 返回 io.EOF 表示流正常结束,其他错误应返回给调用方。无论是否读到末尾,都要关闭读取器。需要合并消息时使用 schema.ConcatMessages,不能只拼接文本就丢掉其他字段。一个读取器只有一个消费进度,多处消费应先复制流,并分别关闭各自的读取器。
2.5 绑定工具
下面的辅助函数接收支持工具调用的模型和工具实例,取得工具描述后,调用 WithTools 派生一个绑定工具的新模型。它复用第 9 节创建的工具,需要导入 model、tool 和 schema 包。
func generateWithTool(ctx context.Context, base model.ToolCallingChatModel, noteTool tool.InvokableTool) (*schema.Message, error) {
info, err := noteTool.Info(ctx)
if err != nil {
return nil, err
}
bound, err := base.WithTools([]*schema.ToolInfo{info})
if err != nil {
return nil, err
}
return bound.Generate(ctx, []*schema.Message{
schema.UserMessage("请查找 P4 笔记的链接。"),
})
}WithTools 不修改原实例。绑定工具后,模型可能返回 ToolCalls,也可能直接回答。模型生成调用名称和参数,并不执行 Go 函数,执行过程由应用代码或 ToolsNode 完成。
3. ChatTemplate 组件
3.1 接口与格式化
type ChatTemplate interface {
Format(ctx context.Context, vs map[string]any, opts ...Option) ([]*schema.Message, error)
}prompt.FromMessages 把实现 schema.MessagesTemplate 的对象组合为模板。*schema.Message 实现了相应格式化方法,消息构造函数返回的指针可以直接使用。Format 接收变量映射,输出 []*schema.Message,它本身不会调用模型。
SystemMessage、UserMessage、AssistantMessage 和 ToolMessage 用于构造不同角色的消息。AssistantMessage 还接收工具调用列表,ToolMessage 需要对应的工具调用 ID。
3.2 示例
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino/components/prompt"
"github.com/cloudwego/eino/schema"
"log"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
ctx := context.Background()
template := prompt.FromMessages(schema.FString,
schema.SystemMessage("你是一名{role},请用中文回答。"),
schema.MessagesPlaceholder("history", true),
schema.UserMessage("请解释{topic}。"),
)
params := map[string]any{
"role": "网络技术助理",
"topic": "软件定义网络",
"history": []*schema.Message{
schema.UserMessage("我了解基本的 TCP/IP 协议。"),
schema.AssistantMessage("接下来可以学习控制平面和数据平面的分工。", nil),
},
}
messages, err := template.Format(ctx, params)
if err != nil {
return err
}
for _, message := range messages {
fmt.Printf("%s: %s\n", message.Role, message.Content)
}
return nil
}这里使用 schema.FString,{role} 和 {topic} 从变量映射取值。需要输出字面量花括号时,将对应的花括号写两次。缺少必需的格式化变量会返回错误。
MessagesPlaceholder("history", true) 把历史消息列表插入模板,第二个参数表示该变量可以缺省。传入历史时,值应为 []*schema.Message。模板不会自动保存历史,需要由应用准备并传入。
4. RAG
RAG(Retrieval-Augmented Generation,检索增强生成)在生成前检索外部资料,再把相关片段作为上下文交给模型。它可以补充训练知识之外的信息,也有助于减少缺乏依据的回答,但不能保证消除幻觉。结果仍取决于资料质量、检索覆盖和模型是否正确使用资料。
典型流程分为两个阶段:
- 建立索引。收集资料,用 Document Transformer 清理或切分文档。采用向量检索时,用 Embedding 生成向量,再由 Indexer 写入后端。
- 检索与生成。Retriever 根据用户问题取回文档,应用整理来源和上下文,通过 ChatTemplate 或消息列表交给 ChatModel。
RAG 不要求必须使用向量数据库,也可以采用关键词检索、稀疏检索或混合检索。向量只是其中一种方式。检索结果是否及时更新,取决于资料和索引的更新流程,并不天然具有实时性。
下面的 Redis 示例演示索引和检索两个步骤,还没有把检索结果传给模型生成最终答案。另有 RAG_Learning 练习项目,使用时应核对其自身依赖。
5. Embedding 组件
5.1 接口
type Embedder interface {
EmbedStrings(ctx context.Context, texts []string, opts ...Option) ([][]float64, error) // invoke
}EmbedStrings 将文本列表转换为向量列表,[][]float64 的外层对应输入文本,内层是各维的值。向量维度、适用语言和检索效果由所选模型决定。
5.2 示例
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino-ext/components/embedding/ark"
"log"
"os"
"time"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
apiKey, modelID := os.Getenv("ARK_API_KEY"), os.Getenv("EMBEDDER")
if apiKey == "" || modelID == "" {
return fmt.Errorf("请设置 ARK_API_KEY 和 EMBEDDER")
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
apiType := ark.APITypeText
if value := os.Getenv("EMBEDDER_API_TYPE"); value != "" {
apiType = ark.APIType(value)
if apiType != ark.APITypeText && apiType != ark.APITypeMultiModal {
return fmt.Errorf("EMBEDDER_API_TYPE 应为 text_api 或 multi_modal_api")
}
}
embedder, err := ark.NewEmbedder(ctx, &ark.EmbeddingConfig{
APIKey: apiKey,
Model: modelID,
BaseURL: os.Getenv("ARK_BASE_URL"),
APIType: &apiType,
})
if err != nil {
return err
}
texts := []string{"软件定义网络", "可编程数据平面", "向量检索"}
vectors, err := embedder.EmbedStrings(ctx, texts)
if err != nil {
return err
}
if len(vectors) != len(texts) {
return fmt.Errorf("向量数量与输入文本数量不一致")
}
for i, vector := range vectors {
fmt.Printf("文本 %d 的向量维度:%d\n", i+1, len(vector))
}
return nil
}文本索引和查询应使用同一个嵌入模型,并保持接口类型、维度及其他相关配置一致。两个模型即使输出维度相同,向量空间也不一定兼容。更换模型后通常需要重新生成文档向量。
语义接近的文本在合适的向量空间中通常较接近,但具体排序还取决于余弦距离、内积或欧氏距离等度量,不能只凭向量维度判断检索质量。
6. Indexer 组件
6.1 接口
type Indexer interface {
// Store stores the documents.
Store(ctx context.Context, docs []*schema.Document, opts ...Option) (ids []string, err error) // invoke
}Store 接收文档并返回写入后的 ID。Indexer 负责与索引后端交互,并不局限于向量数据库。是否生成向量、如何映射字段以及写入后的检索可见性,由具体实现决定。
6.2 Redis 示例的条件
本例需要支持 FT.CREATE 和 FT.SEARCH 的 Redis Search/Query Engine。可使用包含这些功能的 Redis 8.2.x 发行包,并确认服务已加载查询引擎。只有普通数据命令可用的 Redis 服务不足以运行这个示例。
Redis 扩展把嵌入结果转换为 FLOAT32 二进制向量,因此索引的向量字段也使用 FLOAT32。以下配置必须对应:
| 项目 | 本例配置 |
|---|---|
| 索引名称 | eino_notes |
| Hash 键前缀 | eino:note: |
| 文本字段 | content |
| 向量字段 | vector_content |
| 向量维度 | 从所选模型的实际响应中读取 |
| 距离度量 | COSINE |
6.3 写入文档
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino-ext/components/embedding/ark"
ri "github.com/cloudwego/eino-ext/components/indexer/redis"
"github.com/cloudwego/eino/schema"
goredis "github.com/redis/go-redis/v9"
"log"
"os"
"time"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
apiKey, modelID := os.Getenv("ARK_API_KEY"), os.Getenv("EMBEDDER")
if apiKey == "" || modelID == "" {
return fmt.Errorf("请设置 ARK_API_KEY 和 EMBEDDER")
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
apiType := ark.APITypeText
if value := os.Getenv("EMBEDDER_API_TYPE"); value != "" {
apiType = ark.APIType(value)
if apiType != ark.APITypeText && apiType != ark.APITypeMultiModal {
return fmt.Errorf("EMBEDDER_API_TYPE 应为 text_api 或 multi_modal_api")
}
}
embedder, err := ark.NewEmbedder(ctx, &ark.EmbeddingConfig{
APIKey: apiKey,
Model: modelID,
BaseURL: os.Getenv("ARK_BASE_URL"),
APIType: &apiType,
})
if err != nil {
return err
}
address := os.Getenv("REDIS_ADDR")
if address == "" {
return fmt.Errorf("请设置 REDIS_ADDR")
}
client := goredis.NewClient(&goredis.Options{
Addr: address,
Password: os.Getenv("REDIS_PASSWORD"),
Protocol: 2,
})
defer client.Close()
if err := client.Ping(ctx).Err(); err != nil {
return err
}
// 按所选嵌入模型的实际输出维度创建新索引。
probe, err := embedder.EmbedStrings(ctx, []string{"维度检查"})
if err != nil {
return err
}
if len(probe) != 1 || len(probe[0]) == 0 {
return fmt.Errorf("嵌入模型未返回有效向量")
}
dimension := len(probe[0])
err = client.FTCreate(ctx, "eino_notes", &goredis.FTCreateOptions{
OnHash: true,
Prefix: []any{"eino:note:"},
},
&goredis.FieldSchema{FieldName: "content", FieldType: goredis.SearchFieldTypeText},
&goredis.FieldSchema{
FieldName: "vector_content",
FieldType: goredis.SearchFieldTypeVector,
VectorArgs: &goredis.FTVectorArgs{
FlatOptions: &goredis.FTFlatOptions{
Type: "FLOAT32", Dim: dimension, DistanceMetric: "COSINE",
},
},
},
).Err()
if err != nil {
return fmt.Errorf("创建索引失败: %w", err)
}
indexer, err := ri.NewIndexer(ctx, &ri.IndexerConfig{
Client: client,
KeyPrefix: "eino:note:",
Embedding: embedder,
})
if err != nil {
return err
}
ids, err := indexer.Store(ctx, []*schema.Document{
{ID: "p4", Content: "P4 用于描述可编程数据平面的包处理逻辑。"},
{ID: "sdn", Content: "SDN 将网络控制逻辑与数据转发功能分离。"},
{ID: "int", Content: "INT 在数据包中携带网络遥测信息。"},
})
if err != nil {
return err
}
fmt.Printf("索引维度:%d,文档 ID:%v\n", dimension, ids)
return nil
}程序先生成一个探测向量,用其长度创建索引,再写入三条文档。索引已经存在时,创建步骤会返回错误。复用已有索引前,应核对字段、维度、数据类型和距离度量,不能仅凭索引同名就判断配置相同。
这里的 Store 返回文档 ID,例如 p4,实际 Hash 键为 eino:note:p4。同一前缀下重复写入相同 ID 会更新对应 Hash 中的字段。示例使用 Protocol: 2 让查询结果按 RESP2 解析,不需要同时启用 UnstableResp3。
7. Retriever 组件
7.1 接口与选项
type Retriever interface {
Retrieve(ctx context.Context, query string, opts ...Option) ([]*schema.Document, error)
}Retrieve 根据查询返回 []*schema.Document。它可以使用向量、关键词或其他检索方式,公共选项如下:
type Options struct {
Index *string
SubIndex *string
TopK *int
// 阈值的含义和支持情况由检索器实现决定。
ScoreThreshold *float64
Embedding embedding.Embedder
// 后端特有的过滤或查询参数,不限定为某一种检索器。
DSLInfo map[string]any
}TopK 控制返回数量。Index、SubIndex、DSLInfo 等选项的含义和支持情况由后端实现决定。ScoreThreshold 也不是所有检索器通用的“相似度下限”,不能把某个后端的数值解释直接套到另一个后端。
7.2 Redis 检索
先运行上一节的索引程序,再运行检索程序:
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino-ext/components/embedding/ark"
rr "github.com/cloudwego/eino-ext/components/retriever/redis"
goredis "github.com/redis/go-redis/v9"
"log"
"os"
"time"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
apiKey, modelID := os.Getenv("ARK_API_KEY"), os.Getenv("EMBEDDER")
if apiKey == "" || modelID == "" {
return fmt.Errorf("请设置 ARK_API_KEY 和 EMBEDDER")
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
apiType := ark.APITypeText
if value := os.Getenv("EMBEDDER_API_TYPE"); value != "" {
apiType = ark.APIType(value)
if apiType != ark.APITypeText && apiType != ark.APITypeMultiModal {
return fmt.Errorf("EMBEDDER_API_TYPE 应为 text_api 或 multi_modal_api")
}
}
embedder, err := ark.NewEmbedder(ctx, &ark.EmbeddingConfig{
APIKey: apiKey,
Model: modelID,
BaseURL: os.Getenv("ARK_BASE_URL"),
APIType: &apiType,
})
if err != nil {
return err
}
address := os.Getenv("REDIS_ADDR")
if address == "" {
return fmt.Errorf("请设置 REDIS_ADDR")
}
client := goredis.NewClient(&goredis.Options{
Addr: address,
Password: os.Getenv("REDIS_PASSWORD"),
Protocol: 2,
})
defer client.Close()
if err := client.Ping(ctx).Err(); err != nil {
return err
}
retriever, err := rr.NewRetriever(ctx, &rr.RetrieverConfig{
Client: client,
Index: "eino_notes",
VectorField: "vector_content",
ReturnFields: []string{"content"},
TopK: 2,
Embedding: embedder,
})
if err != nil {
return err
}
docs, err := retriever.Retrieve(ctx, "什么是可编程数据平面?")
if err != nil {
return err
}
for _, doc := range docs {
fmt.Printf("%s: %s\n", doc.ID, doc.Content)
}
return nil
}示例用 KNN 查询返回最多两条文档,并只取回 content 字段。该扩展返回的文档 ID 是实际 Redis 键,包含 eino:note: 前缀。检索顺序取决于模型的向量响应和距离,正文不固定具体命中的文档。
本文固定的 Redis 扩展使用 RetrieverConfig.DistanceThreshold 配置向量范围查询,距离越小表示越接近。它不会把公共 WithScoreThreshold 选项作为这个配置的替代。需要限制距离时,应配置 DistanceThreshold,而不是按“分数必须大于 0.5”理解阈值。
8. Document Transformer 组件
8.1 接口
type Transformer interface {
Transform(ctx context.Context, src []*schema.Document, opts ...TransformerOption) ([]*schema.Document, error)
}Transformer 接收并返回文档列表,用于切分、清理、过滤或其他转换。它不等同于嵌入模型,也不保证每个实现都会生成向量。
8.2 按 Markdown 标题切分
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino-ext/components/document/transformer/splitter/markdown"
"github.com/cloudwego/eino/schema"
"log"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
ctx := context.Background()
splitter, err := markdown.NewHeaderSplitter(ctx, &markdown.HeaderConfig{
Headers: map[string]string{"#": "h1", "##": "h2", "###": "h3"},
TrimHeaders: true,
})
if err != nil {
return err
}
docs := []*schema.Document{{
ID: "network-note",
Content: "# 网络笔记\n\n介绍网络体系结构。\n\n## SDN\n\n控制平面与数据平面分离。\n\n### P4\n\n描述数据平面的包处理逻辑。\n\n## INT\n\n记录转发路径上的遥测信息。",
}}
results, err := splitter.Transform(ctx, docs)
if err != nil {
return err
}
for _, result := range results {
fmt.Println(result.Content)
fmt.Printf("标题:%v / %v / %v\n", result.MetaData["h1"], result.MetaData["h2"], result.MetaData["h3"])
}
return nil
}Headers 把标题级别映射为元数据字段。TrimHeaders: true 从输出正文中移除匹配的标题行,标题仍可保存在 h1、h2、h3 元数据中。
标题切分并不保证每块都小于模型的 token 上限,较长章节还需要进一步按长度切分。分块后写入索引时,也需要为每个块分配合适的文档 ID。
9. Tool 与 ToolsNode
9.1 工具接口
Tool 描述一项可执行能力。ToolsNode 是编排中的执行节点,它根据 assistant 消息里的工具名称和参数调用已注册的工具,再生成工具结果消息,两者职责不同。
// BaseTool 基础工具接口,提供工具信息
type BaseTool interface {
Info(ctx context.Context) (*schema.ToolInfo, error)
}
// InvokableTool 支持同步调用的工具接口
type InvokableTool interface {
BaseTool
// InvokableRun call function with arguments in JSON format
InvokableRun(ctx context.Context, argumentsInJSON string, opts ...Option) (string, error)
}
// StreamableTool 支持流式输出的工具接口
type StreamableTool interface {
BaseTool
StreamableRun(ctx context.Context, argumentsInJSON string, opts ...Option) (*schema.StreamReader[string], error)
}BaseTool.Info 提供名称、描述和参数定义。InvokableTool.InvokableRun 接收 JSON 参数字符串并返回完整结果,StreamableTool.StreamableRun 则返回字符串流。这里展示的是文本结果接口,该版本还提供支持结构化多模态结果的增强接口。
9.2 工具描述
type ToolInfo struct {
Name string
Desc string
Extra map[string]any
// 用 NewParamsOneOfByParams 或 NewParamsOneOfByJSONSchema 描述参数。
*ParamsOneOf
}参数可通过 NewParamsOneOfByParams 或 NewParamsOneOfByJSONSchema 描述。参数声明用于指导模型生成调用,不应代替工具自身的输入校验。
9.3 HTTP GET 工具
package main
import (
"context"
"encoding/json"
"fmt"
req "github.com/cloudwego/eino-ext/components/tool/httprequest/get"
"log"
"net/http"
"time"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
getTool, err := req.NewTool(ctx, &req.Config{
Headers: map[string]string{"User-Agent": "EinoNotesExample"},
HttpClient: &http.Client{Timeout: 10 * time.Second},
})
if err != nil {
return err
}
arguments, err := json.Marshal(&req.GetRequest{URL: "https://zhh2001.github.io/sitemap.xml"})
if err != nil {
return err
}
result, err := getTool.InvokableRun(ctx, string(arguments))
if err != nil {
return err
}
fmt.Println(result)
return nil
}示例直接调用工具获取本站的 sitemap.xml,不经过模型。工具使用 HTTP 客户端超时,调用本身也带上下文超时。
9.4 创建本地工具
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino/components/tool"
"github.com/cloudwego/eino/components/tool/utils"
"github.com/cloudwego/eino/schema"
"log"
"strings"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
noteTool := CreateTool()
result, err := noteTool.InvokableRun(context.Background(), `{"name":"P4"}`)
if err != nil {
return err
}
fmt.Println(result)
return nil
}
type InputParams struct {
Name string `json:"name" jsonschema:"description=技术名称"`
}
func GetNote(ctx context.Context, params *InputParams) (string, error) {
if err := ctx.Err(); err != nil {
return "", err
}
if params == nil || strings.TrimSpace(params.Name) == "" {
return "", fmt.Errorf("技术名称不能为空")
}
notes := map[string]string{
"p4": "https://zhh2001.github.io/sdn/p4",
"int": "https://zhh2001.github.io/sdn/int",
"mininet": "https://zhh2001.github.io/sdn/mininet",
"iperf": "https://zhh2001.github.io/sdn/iperf",
}
url, ok := notes[strings.ToLower(strings.TrimSpace(params.Name))]
if !ok {
return "", fmt.Errorf("未找到对应笔记")
}
return url, nil
}
func CreateTool() tool.InvokableTool {
return utils.NewTool(&schema.ToolInfo{
Name: "get_note",
Desc: "根据技术名称获取学习笔记链接",
ParamsOneOf: schema.NewParamsOneOfByParams(map[string]*schema.ParameterInfo{
"name": {Type: schema.String, Desc: "技术名称", Required: true},
}),
}, GetNote)
}utils.NewTool 根据显式的 ToolInfo 封装 Go 函数,负责 JSON 参数解码和结果转换。GetNote 检查空名称和未知名称,正常运行时输出 P4 笔记链接。也可以用 utils.InferTool 从函数参数类型推导参数定义。
9.5 通过 ToolsNode 执行
下面复用上一节的 InputParams、GetNote 和 CreateTool。将辅助函数与这些定义放在同一个包,补充 compose 导入,在已有 main 中调用 executeToolCall 并处理返回值:
func executeToolCall(ctx context.Context) ([]*schema.Message, error) {
node, err := compose.NewToolNode(ctx, &compose.ToolsNodeConfig{
Tools: []tool.BaseTool{CreateTool()},
})
if err != nil {
return nil, err
}
message := schema.AssistantMessage("", []schema.ToolCall{{
ID: "call_1",
Type: "function",
Function: schema.FunctionCall{Name: "get_note", Arguments: `{"name":"P4"}`},
}})
return node.Invoke(ctx, message)
}这里手工构造一条工具调用消息,因此无需模型服务。返回的工具消息通过 ToolCallID 关联到 call_1。实际对话中,还要把 assistant 的工具调用消息和 tool 结果追加到消息历史,再交给模型继续回答。ToolsNode 本身不会完成这一整轮对话。
10. 编排
10.1 Chain
Chain 适合顺序执行的流程。编译时检查节点连接和类型是否匹配,得到 Runnable 后再调用 Invoke 或 Stream。本例由模板生成消息列表,再调用模型:

package main
import (
"context"
"fmt"
"github.com/cloudwego/eino-ext/components/model/ark"
"github.com/cloudwego/eino/components/prompt"
"github.com/cloudwego/eino/compose"
"github.com/cloudwego/eino/schema"
"log"
"os"
"time"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
apiKey, modelID := os.Getenv("ARK_API_KEY"), os.Getenv("MODEL")
if apiKey == "" || modelID == "" {
return fmt.Errorf("请设置 ARK_API_KEY 和 MODEL")
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
chatModel, err := ark.NewChatModel(ctx, &ark.ChatModelConfig{
APIKey: apiKey,
Model: modelID,
BaseURL: os.Getenv("ARK_BASE_URL"),
})
if err != nil {
return err
}
template := prompt.FromMessages(schema.FString,
schema.SystemMessage("你是一名{role},请用中文回答。"),
schema.MessagesPlaceholder("history", true),
schema.UserMessage("请解释{topic}。"),
)
chain := compose.NewChain[map[string]any, *schema.Message]()
chain.AppendChatTemplate(template).AppendChatModel(chatModel)
runnable, err := chain.Compile(ctx)
if err != nil {
return err
}
output, err := runnable.Invoke(ctx, map[string]any{
"role": "网络技术助理", "topic": "软件定义网络",
})
if err != nil {
return err
}
fmt.Println(output.Content)
return nil
}输入是 map[string]any,模板输出是 []*schema.Message,模型输出是 *schema.Message。模板允许缺省历史,所以此处只传入角色和主题。
10.2 Graph 与分支
Graph 通过节点、边和分支描述执行路径,可以支持循环,也可以配置为 DAG。AddBranch 的条件函数根据上游输出选择目标,可能到达的目标需要在分支定义中声明。
本例的执行路径为:
START → classify → cat → END
→ dog → END
→ other → ENDpackage main
import (
"context"
"fmt"
"github.com/cloudwego/eino/compose"
"log"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
ctx := context.Background()
graph := compose.NewGraph[string, string]()
classify := compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
switch input {
case "1":
return "cat", nil
case "2":
return "dog", nil
default:
return "other", nil
}
})
if err := graph.AddLambdaNode("classify", classify); err != nil {
return err
}
for _, name := range []string{"cat", "dog", "other"} {
lambda := compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
return map[string]string{"cat": "喵喵喵", "dog": "汪汪汪", "other": "你好"}[input], nil
})
if err := graph.AddLambdaNode(name, lambda); err != nil {
return err
}
if err := graph.AddEdge(name, compose.END); err != nil {
return err
}
}
branch := compose.NewGraphBranch(func(ctx context.Context, input string) (string, error) {
return input, nil
}, map[string]bool{"cat": true, "dog": true, "other": true})
if err := graph.AddBranch("classify", branch); err != nil {
return err
}
if err := graph.AddEdge(compose.START, "classify"); err != nil {
return err
}
runnable, err := graph.Compile(ctx)
if err != nil {
return err
}
for _, input := range []string{"1", "2", "3"} {
output, err := runnable.Invoke(ctx, input)
if err != nil {
return err
}
fmt.Println(output)
}
return nil
}输入 1、2 和其他值分别选择三个分支,示例依次输出 喵喵喵、汪汪汪 和 你好。分支由普通 Go 逻辑决定,不需要模型判断。
10.3 Graph 中的模型节点
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino-ext/components/model/ark"
"github.com/cloudwego/eino/compose"
"github.com/cloudwego/eino/schema"
"log"
"os"
"time"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
apiKey, modelID := os.Getenv("ARK_API_KEY"), os.Getenv("MODEL")
if apiKey == "" || modelID == "" {
return fmt.Errorf("请设置 ARK_API_KEY 和 MODEL")
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
chatModel, err := ark.NewChatModel(ctx, &ark.ChatModelConfig{
APIKey: apiKey,
Model: modelID,
BaseURL: os.Getenv("ARK_BASE_URL"),
})
if err != nil {
return err
}
graph := compose.NewGraph[map[string]string, *schema.Message]()
choose := compose.InvokableLambda(func(ctx context.Context, input map[string]string) (map[string]string, error) {
if input["style"] != "brief" && input["style"] != "detail" {
return nil, fmt.Errorf("style 应为 brief 或 detail")
}
return input, nil
})
if err := graph.AddLambdaNode("choose", choose); err != nil {
return err
}
for _, style := range []string{"brief", "detail"} {
prepare := compose.InvokableLambda(func(ctx context.Context, input map[string]string) ([]*schema.Message, error) {
instruction := "请简要解释技术概念。"
if input["style"] == "detail" {
instruction = "请分步骤解释技术概念,并给出一个例子。"
}
return []*schema.Message{schema.SystemMessage(instruction), schema.UserMessage(input["content"])}, nil
})
if err := graph.AddLambdaNode(style, prepare); err != nil {
return err
}
}
if err := graph.AddChatModelNode("model", chatModel); err != nil {
return err
}
branch := compose.NewGraphBranch(func(ctx context.Context, input map[string]string) (string, error) {
return input["style"], nil
}, map[string]bool{"brief": true, "detail": true})
if err := graph.AddBranch("choose", branch); err != nil {
return err
}
for _, edge := range [][2]string{{compose.START, "choose"}, {"brief", "model"}, {"detail", "model"}, {"model", compose.END}} {
if err := graph.AddEdge(edge[0], edge[1]); err != nil {
return err
}
}
runnable, err := graph.Compile(ctx)
if err != nil {
return err
}
output, err := runnable.Invoke(ctx, map[string]string{"style": "brief", "content": "什么是 P4?"})
if err != nil {
return err
}
fmt.Println(output.Content)
return nil
}输入包含 style 和 content。分支选择简要或详细的系统提示,两条路径都把用户问题保留为 user 消息,再交给同一个模型节点。节点之间的输入输出类型必须匹配,添加节点、边和编译时的错误都需要处理。
10.4 调用内状态
WithGenLocalState 为每次图调用创建状态。它在本次调用的节点间共享,不会自动变成跨请求的会话历史或持久化数据。状态生成函数应创建新对象,重复返回同一个共享指针会破坏请求隔离。
StatePreHandler 和 StatePostHandler 在节点前后读取或修改状态,也可以改变节点的输入输出。流式场景有对应的 StreamStatePreHandler 和 StreamStatePostHandler。在流上使用非流式处理器可能触发聚合,不能据此保证增量输出不受影响。
节点内部访问状态可以使用 compose.ProcessState,公开签名如下:
// 公开 API 签名,省略框架内部的状态查找和加锁实现。
func ProcessState[S any](ctx context.Context, handler func(context.Context, S) error) error框架在状态处理器和 ProcessState 的回调执行期间加锁。不要把状态指针带到回调外直接修改,也不要在持有同一状态锁的回调里再次调用 ProcessState,否则可能产生数据竞争或死锁。流式处理器返回后,另起 goroutine 的访问也不自动受这把锁保护。
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino/compose"
"log"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
ctx := context.Background()
graph := compose.NewGraph[string, string](compose.WithGenLocalState(func(ctx context.Context) *State {
return &State{}
}))
count := compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
err := compose.ProcessState(ctx, func(ctx context.Context, state *State) error {
state.Count++
return nil
})
return input, err
})
format := compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
return input, nil
})
pre := func(ctx context.Context, input string, state *State) (string, error) {
return fmt.Sprintf("%s: count=%d", input, state.Count), nil
}
if err := graph.AddLambdaNode("count", count); err != nil {
return err
}
if err := graph.AddLambdaNode("format", format, compose.WithStatePreHandler(pre)); err != nil {
return err
}
for _, edge := range [][2]string{{compose.START, "count"}, {"count", "format"}, {"format", compose.END}} {
if err := graph.AddEdge(edge[0], edge[1]); err != nil {
return err
}
}
runnable, err := graph.Compile(ctx)
if err != nil {
return err
}
for i := 0; i < 2; i++ {
output, err := runnable.Invoke(ctx, "request")
if err != nil {
return err
}
fmt.Println(output)
}
return nil
}
type State struct{ Count int }程序连续调用同一个 Runnable 两次,每次都输出 request: count=1。count 节点通过 ProcessState 更新状态,format 节点的前处理器读取它,展示的是调用内共享和调用间隔离。
10.5 回调
回调用于记录组件和编排执行过程。RunInfo 描述触发回调的实体,Name 是展示名称,Type 是实现类型,Component 是组件类别:
type RunInfo struct {
// Name is the graph node name for display purposes, not unique.
// Passed from compose.WithNodeName().
Name string
Type string
Component components.Component
}回调输入输出按组件区分,下面的类型定义不代表它们具有统一的结构:
type CallbackInput any
type CallbackOutput anypackage main
import (
"context"
"fmt"
"github.com/cloudwego/eino/callbacks"
"github.com/cloudwego/eino/compose"
"log"
"strings"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
ctx := context.Background()
chain := compose.NewChain[string, string]()
chain.AppendLambda(compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
return strings.ToUpper(input), nil
}), compose.WithNodeName("upper"))
runnable, err := chain.Compile(ctx)
if err != nil {
return err
}
output, err := runnable.Invoke(ctx, "eino", compose.WithCallbacks(genCallback()))
if err != nil {
return err
}
fmt.Println(output)
return nil
}
func genCallback() callbacks.Handler {
return callbacks.NewHandlerBuilder().
OnStartFn(func(ctx context.Context, info *callbacks.RunInfo, input callbacks.CallbackInput) context.Context {
if info != nil {
fmt.Printf("start: %s %s\n", info.Component, info.Name)
}
return ctx
}).
OnEndFn(func(ctx context.Context, info *callbacks.RunInfo, output callbacks.CallbackOutput) context.Context {
if info != nil {
fmt.Printf("end: %s %s\n", info.Component, info.Name)
}
return ctx
}).
OnErrorFn(func(ctx context.Context, info *callbacks.RunInfo, err error) context.Context {
fmt.Printf("error: %v\n", err)
return ctx
}).Build()
}示例记录开始、结束和错误,并通过 compose.WithCallbacks 应用于本次调用。需要读取特定组件的数据时,应使用该组件的转换函数,例如 model.ConvCallbackInput,并检查转换结果。
回调不应修改共享的输入输出对象。多个处理器之间没有可依赖的执行顺序,并发节点的日志也可能交错。流式回调拿到的是流副本,使用后同样需要关闭,不能把副本遗留到调用结束以后。
10.6 嵌套图
AddGraphNode 可以把图或 Chain 直接加入外层图,外层 Compile 会处理内部编排,不要求先把内层图编译后再包装为 Lambda。嵌套前后仍需保持输入输出类型匹配。
package main
import (
"context"
"fmt"
"github.com/cloudwego/eino/compose"
"log"
"strings"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
ctx := context.Background()
inner := compose.NewGraph[string, string]()
upper := compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
return strings.ToUpper(strings.TrimSpace(input)), nil
})
if err := inner.AddLambdaNode("upper", upper); err != nil {
return err
}
for _, edge := range [][2]string{{compose.START, "upper"}, {"upper", compose.END}} {
if err := inner.AddEdge(edge[0], edge[1]); err != nil {
return err
}
}
outer := compose.NewGraph[string, string]()
if err := outer.AddGraphNode("inner", inner); err != nil {
return err
}
decorate := compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
return "result: " + input, nil
})
if err := outer.AddLambdaNode("decorate", decorate); err != nil {
return err
}
for _, edge := range [][2]string{{compose.START, "inner"}, {"inner", "decorate"}, {"decorate", compose.END}} {
if err := outer.AddEdge(edge[0], edge[1]); err != nil {
return err
}
}
runnable, err := outer.Compile(ctx)
if err != nil {
return err
}
output, err := runnable.Invoke(ctx, " eino ")
if err != nil {
return err
}
fmt.Println(output)
return nil
}内层图去除首尾空白并转换为大写,外层图追加前缀,最终输出 result: EINO。
11. CozeLoop
CozeLoop 回调扩展把 Eino 的执行信息转换为 trace。运行示例前,还需要设置 COZELOOP_WORKSPACE_ID 和 COZELOOP_API_TOKEN,对应账号中的工作空间与访问凭据。
package main
import (
"context"
"fmt"
ccb "github.com/cloudwego/eino-ext/callbacks/cozeloop"
"github.com/cloudwego/eino/callbacks"
"github.com/cloudwego/eino/compose"
"github.com/coze-dev/cozeloop-go"
"log"
"os"
"strings"
"time"
)
func main() {
if err := run(); err != nil {
log.Fatal(err)
}
}
func run() error {
if os.Getenv("COZELOOP_WORKSPACE_ID") == "" || os.Getenv("COZELOOP_API_TOKEN") == "" {
return fmt.Errorf("请设置 COZELOOP_WORKSPACE_ID 和 COZELOOP_API_TOKEN")
}
client, err := cozeloop.NewClient()
if err != nil {
return err
}
defer func() {
closeCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
client.Close(closeCtx)
}()
// 进程初始化时注册一次,不要在每次请求时重复追加。
callbacks.AppendGlobalHandlers(ccb.NewLoopHandler(client))
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
chain := compose.NewChain[string, string]()
chain.AppendLambda(compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
return strings.ToUpper(input), nil
}), compose.WithNodeName("upper"))
runnable, err := chain.Compile(ctx)
if err != nil {
return err
}
output, err := runnable.Invoke(ctx, "eino")
if err != nil {
return err
}
fmt.Println(output)
return nil
}AppendGlobalHandlers 影响整个进程,应在初始化阶段注册一次,不要在每次请求时重复追加,也不要与正在执行的调用并发修改全局回调列表。只想给某次调用添加追踪时,可以使用 compose.WithCallbacks(handler)。
客户端关闭时需要留出上报时间。示例另建关闭上下文,防止使用已经取消的业务上下文。程序得到 EINO 只表示本地编排成功,trace 是否出现在平台上还取决于凭据、网络和服务端接收情况。