Files
CamTalk/docs/11-Eino框架技术文档.md
cfy666 9576884619 docs: 新增 Eino 框架技术文档和重构实施记录
- docs/11-Eino框架技术文档.md: 框架简介、技术选型对比、核心概念(Lambda/Graph/ChatModel/StreamReader/Callback/State)、CamTalk Graph 设计、目录结构、注意事项
- docs/12-Eino重构实施记录.md: 重构背景、架构变更、四阶段实施详情、代码统计、遗留事项

Co-Authored-By: Claude <noreply@anthropic.com>
2026-06-19 22:09:07 +08:00

8.7 KiB
Raw Blame History

CamTalk Eino 框架技术文档

创建日期2026-06-19 状态:已实施

1. 框架简介

CloudWeGo Eino 是字节跳动 CloudWeGo 团队开源的 AI 应用开发框架提供基于图Graph的编排能力、组件抽象和流式处理支持。

CamTalk 使用 Eino 替代原有的手写 goroutine 管道,实现 STT → LLM → TTS 的声明式编排。

2. 技术选型

2.1 为什么选 Eino

维度 手写 goroutine旧方案 Eino Graph新方案
编排方式 手动 go func() + sync.WaitGroup 声明式 DAG类型安全
流式处理 自定义 chan 传递 StreamReader + Pipe,自动转换
错误处理 各节点独立处理,不一致 Graph 级别统一错误传播
回调/AOP 日志散落各处 callbacks.Handler 统一注入
配置灵活性 Pipeline 创建时固定 每请求 Option 动态注入
可测试性 需启动 goroutine Graph.Invoke() 直接测试
扩展性 修改 Pipeline 代码 添加节点 + 边,无侵入
并发安全 手动 sync State 自动加锁

2.2 Eino vs 其他编排框架

框架 特点 CamTalk 适用性
Eino Go 原生、类型安全、流式原生 完美匹配
LangChain Go 生态丰富但较重 过度抽象
自研编排 完全可控 维护成本高

选择 Eino 的核心理由

  1. Go 原生,泛型支持,编译时类型检查
  2. 原生流式处理(StreamReader),适合 LLM token 级推送
  3. Graph 支持分支、并行、循环,满足当前和未来需求
  4. Callback 机制实现 AOP日志、指标、消息推送
  5. eino-ext 提供 OpenAI ChatModel 实现,直接对接 DashScope

2.3 核心依赖版本

github.com/cloudwego/eino v0.9.9
github.com/cloudwego/eino-ext/components/model/openai v0.1.13

3. Eino 核心概念

3.1 Lambda

Lambda 是 Graph 中的可组合函数单元,支持四种模式:

模式 函数签名 构造方法 说明
Invoke I → O compose.InvokableLambda() 同步调用
Stream I → StreamReader[O] compose.StreamableLambda() 流式输出
Collect StreamReader[I] → O compose.CollectableLambda() 流式输入
Transform StreamReader[I] → StreamReader[O] compose.TransformableLambda() 双向流式

返回类型:所有 Lambda 构造函数返回 *compose.Lambda

3.2 Graph

Graph 是有向无环图DAG编排器支持

  • 节点Lambda、ChatModel、ToolsNode 等
  • g.AddEdge(from, to) 定义数据流向
  • 分支g.AddBranch() 条件路由
  • Statecompose.WithGenLocalState() 跨节点共享状态
g := compose.NewGraph[PipelineInput, PipelineOutput]()
g.AddLambdaNode("stt", sttLambda)
g.AddChatModelNode("llm", chatModel)
g.AddEdge(compose.START, "stt")
g.AddEdge("stt", "llm")
g.AddEdge("llm", compose.END)

runnable, err := g.Compile(ctx)
output, err := runnable.Invoke(ctx, input)   // 同步调用
stream, err := runnable.Stream(ctx, input)    // 流式调用

3.3 ChatModel

ChatModel 是 LLM 组件抽象,接口定义:

type BaseChatModel interface {
    Generate(ctx, []*schema.Message, ...Option) (*schema.Message, error)
    Stream(ctx, []*schema.Message, ...Option) (*schema.StreamReader[*schema.Message], error)
}

CamTalk 使用 eino-ext/components/model/openai 实现,通过 BaseURL 对接 DashScope

chatModel, _ := openai.NewChatModel(ctx, &openai.ChatModelConfig{
    APIKey:  cfg.AI.LLM.APIKey,
    Model:   cfg.AI.LLM.Model,
    BaseURL: cfg.AI.LLM.Endpoint, // "https://dashscope.aliyuncs.com/compatible-mode/v1"
})

3.4 StreamReader

schema.StreamReader[T] 是 Eino 的流式数据抽象:

  • sr.Recv() 读取一帧,io.EOF 表示流结束
  • schema.Pipe[T](bufSize) 创建 StreamReader + StreamWriter
  • 框架自动处理 T ↔ StreamReader[T] 的转换(装箱/concat

3.5 Callback

Callback 是 Eino 的 AOP 机制,支持节点生命周期钩子:

type Handler interface {
    OnStart(ctx, *RunInfo, CallbackInput) context.Context
    OnEnd(ctx, *RunInfo, CallbackOutput) context.Context
    OnError(ctx, *RunInfo, error) context.Context
    OnStartWithStreamInput(ctx, *RunInfo, *StreamReader[CallbackInput]) context.Context
    OnEndWithStreamOutput(ctx, *RunInfo, *StreamReader[CallbackOutput]) context.Context
}

CamTalk 使用 utils/callbacks.NewHandlerHelper() 构建 typed handler

  • ModelCallbackHandler.OnEndWithStreamOutput:逐 token 推送 llm_chunk

3.6 State

Graph 全局状态,通过 WithGenLocalState 注册:

type PipelineState struct {
    FullResponse    strings.Builder
    TranscribedText string
    TokenUsage      *TokenUsage
}

g := compose.NewGraph[I, O](compose.WithGenLocalState(func(ctx context.Context) *PipelineState {
    return &PipelineState{}
}))

节点通过 compose.ProcessState 读写 State。

4. CamTalk Graph 设计

4.1 拓扑

START → STT → History → ChatModel → Splitter → TTS → Done → END
节点 类型 输入 → 输出 职责
STT InvokableLambda PipelineInput → STTOutput 语音识别,写入 State
History InvokableLambda STTOutput → []*schema.Message 组装提示词和历史
ChatModel ChatModel原生 []*schema.Message → StreamReader[*Message] LLM 流式推理
Splitter TransformableLambda StreamReader[string] → StreamReader[[]string] 句子切分
TTS InvokableLambda []string → struct{} 语音合成,推送音频
Done InvokableLambda struct{} → PipelineOutput 发送 llm_done

4.2 流式模式

Graph 使用 Stream 模式调用:

  • 内部所有节点以 Transform 模式运行
  • ChatModel 的 Stream() 方法实现真正的 token 级流式
  • 适配器消费 StreamReader[PipelineOutput] 触发整条链路

4.3 消息推送机制

消息 推送方式 时机
stt_result Lambda 内部直接调用 Sender STT 完成后
llm_chunk Callback OnEndWithStreamOutput ChatModel 逐 token
tts_audio Lambda 内部直接调用 Sender TTS 逐句合成
llm_done Lambda 内部直接调用 Sender Done 节点执行时

Context 注入Sender、RequestID、SessionID、PipelineState 通过 context.WithValue 传递。

4.4 多模态支持

History 节点将图片构建为 schema.Message.UserInputMultiContent

systemMsg.UserInputMultiContent = []schema.MessageInputPart{
    {
        Type: schema.ChatMessagePartTypeImageURL,
        Image: &schema.MessageInputImage{
            MessagePartCommon: schema.MessagePartCommon{
                Base64Data: &base64Str,
                MIMEType:   "image/jpeg",
            },
            Detail: schema.ImageURLDetailAuto,
        },
    },
}

5. 目录结构

backend/internal/eino/
├── types.go           # PipelineInput/Output、STTOutput、TokenUsage
├── state.go           # PipelineState跨节点状态
├── callback.go        # Callback handlerLLM token 推送)
├── graph.go           # Graph 构建与编译
├── adapter.go         # EinoOrchestratorOrchestrator 接口适配器)
├── nodes_stt.go       # STT Lambda
├── nodes_history.go   # 历史组装 Lambda
├── nodes_splitter.go  # 句子分割 Transform Lambda
├── nodes_tts.go       # TTS Lambda
├── nodes_done.go      # Done Lambda
└── graph_test.go      # 单元测试

6. 注意事项

6.1 值类型 vs 指针类型

Graph 泛型参数必须使用值类型(PipelineInput/PipelineOutput),所有 Lambda 的输入输出也使用值类型。框架在 Transform 模式下会自动处理 TStreamReader[T] 的转换。

6.2 Callback 运行时传入

Callback 通过 Stream() 的 option 传入,不在 Compile() 时注册:

streamReader, err := runnable.Stream(ctx, input, compose.WithCallbacks(handler))

6.3 eino-ext 与 DashScope 兼容性

eino-ext OpenAI ChatModel 通过 BaseURL 对接 DashScope 兼容接口。需注意:

  • 多模态图片使用 Base64Data + MIMEType 格式
  • Timeout 控制单次请求超时
  • 流式输出通过 Stream() 方法获取 StreamReader[*schema.Message]

6.4 框架自动类型转换

Eino 框架在编排场景中自动处理以下转换:

  • T → StreamReader[T]:将完整值装箱为单帧流(非流式 → 假流式)
  • StreamReader[T] → T:将流 concat 为完整值(流式 → 非流式)

这使得不同流式模式的节点可以无缝连接。