## 主要变更 ### 文档重构(减少 1199 行,-23%) - 01-架构设计.md: 503→369 行 (-27%),删除 DDL/配置示例,精简鉴权/存储描述 - 02-接口文档.md: 1313→570 行 (-57%),删除 Go 接口/Orchestrator 实现/配置管理 - 07-成本控制.md: 65→59 行 (-9%),代码块替换为文件引用 ### 文档编号规范化 - 08-功能创意.md → 删除(内容整合到 README.md "功能扩展方向") - 10-Eino框架与编排设计.md → 08-Eino框架与编排设计.md - 情景切换.md → 09-情景切换.md - 12-鉴权体系.md → 10-鉴权体系.md - 13-令牌桶限流.md → 11-令牌桶限流.md ### 交叉引用更新 - 01-架构设计.md: 更新对鉴权体系/令牌桶限流的引用为新编号 - README.md: 更新文档索引表、推荐阅读顺序、新增功能扩展方向 ### 删除过时文档 - 09-技术名词解释.md(内容已整合到 03-技术选型.md) - 10-Eino重构方案.md(历史记录,已完成) - 11-Eino框架技术文档.md(已合并到 08) - 情景切换功能完整文档.md(已规范化为 09) ## 重构原则 - 架构文档聚焦系统结构,移除实现细节 - 接口文档保留纯契约,删除内部实现 - 编号连续(01-11),语义清晰 - 通过交叉引用连接相关文档,避免重复
17 KiB
CamTalk Eino 框架与编排设计
创建日期:2026-06-19
状态:已实施
合并自:10-Eino重构方案.md+11-Eino框架技术文档.md
1. 概述
1.1 为什么选择 Eino
CloudWeGo Eino 是字节跳动 CloudWeGo 团队开源的 AI 应用开发框架,提供基于图(Graph)的编排能力、组件抽象和流式处理支持。
CamTalk 使用 Eino 替代原有的手写 goroutine 管道,实现 STT → LLM → TTS 的声明式编排。
技术选型对比:
| 维度 | 手写 goroutine(旧方案) | Eino Graph(新方案) |
|---|---|---|
| 编排方式 | 手动 go func() + sync.WaitGroup |
声明式 DAG,类型安全 |
| 流式处理 | 自定义 chan 传递 |
StreamReader + Pipe,自动转换 |
| 错误处理 | 各节点独立处理,不一致 | Graph 级别统一错误传播 |
| 回调/AOP | 日志散落各处 | callbacks.Handler 统一注入 |
| 配置灵活性 | Pipeline 创建时固定 | 每请求 Option 动态注入 |
| 可测试性 | 需启动 goroutine | Graph.Invoke() 直接测试 |
| 扩展性 | 修改 Pipeline 代码 | 添加节点 + 边,无侵入 |
| 并发安全 | 手动 sync |
State 自动加锁 |
选择 Eino 的核心理由:
- Go 原生,泛型支持,编译时类型检查
- 原生流式处理(
StreamReader),适合 LLM token 级推送 - Graph 支持分支、并行、循环,满足当前和未来需求
- Callback 机制实现 AOP(日志、指标、消息推送)
- eino-ext 提供 OpenAI ChatModel 实现,直接对接 DashScope
1.2 旧方案的问题
当前后端 AI 编排层(internal/orchestrator/pipeline.go)为手写 goroutine 管道存在以下问题:
- 编排逻辑硬编码:STT→LLM→TTS 流程写死,扩展困难
- 并发控制粗糙:手动 goroutine 调度,缺乏结构化流式传递
- 无回调/AOP 机制:日志、指标、追踪散落各处
- 配置耦合:模型名、TTS 参数等硬编码在结构体
- 错误处理不一致:TTS 错误静默吞掉,STT/LLM 错误通过 Sender 发送
1.3 核心依赖版本
github.com/cloudwego/eino v0.9.9
github.com/cloudwego/eino-ext/components/model/openai v0.1.13
2. Eino 核心概念
2.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。
2.2 Graph
Graph 是有向无环图(DAG)编排器,支持:
- 节点:Lambda、ChatModel、ToolsNode 等
- 边:
g.AddEdge(from, to)定义数据流向 - 分支:
g.AddBranch()条件路由 - State:
compose.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) // 流式调用
2.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"
})
2.4 StreamReader
schema.StreamReader[T] 是 Eino 的流式数据抽象:
sr.Recv()读取一帧,io.EOF表示流结束schema.Pipe[T](bufSize)创建StreamReader+StreamWriter对- 框架自动处理
T ↔ StreamReader[T]的转换(装箱/concat)
2.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
2.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。
3. CamTalk Graph 设计
3.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 |
3.2 数据类型定义
// Graph 统一输入
type PipelineInput struct {
AudioData []byte // base64 解码后的音频(可选)
ImageData []byte // base64 解码后的图像(可选)
Text string // 直接文本输入(可选,跳过 STT)
SessionID string
RequestID string
Language string // zh / en
Scenario string // free_chat, interviewer, etc.
}
// Graph 统一输出
type PipelineOutput struct {
TranscribedText string // STT 结果
FullResponse string // LLM 完整回复
}
// Pipeline State(跨节点共享)
type PipelineState struct {
FullResponse strings.Builder
TranscribedText string
TokenUsage *TokenUsage
}
3.3 流式模式
Graph 使用 Stream 模式调用:
- 内部所有节点以 Transform 模式运行
- ChatModel 的
Stream()方法实现真正的 token 级流式 - 适配器消费
StreamReader[PipelineOutput]触发整条链路
3.4 消息推送机制
| 消息 | 推送方式 | 时机 |
|---|---|---|
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 传递。
3.5 多模态支持
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,
},
},
}
4. 实现要点
4.1 目录结构
backend/internal/eino/
├── types.go # PipelineInput/Output、STTOutput、TokenUsage
├── state.go # PipelineState(跨节点状态)
├── callback.go # Callback handler(LLM token 推送)
├── graph.go # Graph 构建与编译
├── adapter.go # EinoOrchestrator(Orchestrator 接口适配器)
├── 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 # 单元测试
4.2 关键节点实现
STT Lambda(可选跳过)
func sttLambda(sttSvc stt.Service) func(ctx context.Context, input PipelineInput) (STTOutput, error) {
return func(ctx context.Context, input PipelineInput) (STTOutput, error) {
// 文本模式:跳过 STT
if input.Text != "" {
return STTOutput{Text: input.Text, Language: input.Language}, nil
}
// 调用 STT 服务
result, err := sttSvc.Recognize(ctx, input.AudioData, stt.Options{
Language: input.Language,
})
if err != nil {
return STTOutput{}, fmt.Errorf("STT error: %w", err)
}
return STTOutput{Text: result.Text, Language: result.Language}, nil
}
}
Splitter Transform Lambda(句子切分)
func splitterLambda() func(ctx, *schema.StreamReader[*schema.Message]) (*schema.StreamReader[[]string], error) {
return func(ctx context.Context, stream *schema.StreamReader[*schema.Message]) (*schema.StreamReader[[]string], error) {
sr, sw := schema.Pipe[[]string](8)
go func() {
defer sw.Close()
var buffer []rune
for {
chunk, err := stream.Recv()
if err != nil {
if err == io.EOF {
if len(buffer) > 0 {
sw.Send([]string{string(buffer)}, nil)
}
return
}
sw.Send(nil, err)
return
}
for _, r := range chunk.Content {
buffer = append(buffer, r)
if isSentenceDelimiter(r) {
sw.Send([]string{string(buffer)}, nil)
buffer = buffer[:0]
}
}
}
}()
return sr, nil
}
}
TTS Lambda(并行合成)
func ttsLambda(ttsSvc tts.Service, sender orchestrator.Sender) func(ctx, []string) (struct{}, error) {
return func(ctx context.Context, sentences []string) (struct{}, error) {
for _, sentence := range sentences {
if sentence == "" {
continue
}
// 调用 TTS 服务
audioData, err := ttsSvc.Synthesize(ctx, sentence, tts.Options{})
if err != nil {
// TTS 失败不中断流程,仅记录日志
log.Warn("TTS synthesis failed", zap.Error(err))
continue
}
// 推送音频到客户端
sender.SendTTSAudio(orchestrator.TTSAudioPayload{
Audio: audioData,
Format: "mp3",
})
}
return struct{}{}, nil
}
}
4.3 Callback 集成
// ModelCallbackHandler 用于 LLM token 推送
type ModelCallbackHandler struct {
sender orchestrator.Sender
}
func (h *ModelCallbackHandler) OnEndWithStreamOutput(
ctx context.Context,
info *callbacks.RunInfo,
output *schema.StreamReader[*schema.Message],
) context.Context {
// 逐 token 推送到客户端
for {
msg, err := output.Recv()
if err == io.EOF {
break
}
if err != nil {
return ctx
}
h.sender.SendLLMChunk(orchestrator.LLMChunkPayload{
Content: msg.Content,
})
}
return ctx
}
4.4 按请求动态配置
// 运行时 Option:每请求可变
func WithModelName(name string) compose.Option {
return compose.WithChatModelOption(model.WithModel(name))
}
func WithTemperature(temp float32) compose.Option {
return compose.WithChatModelOption(model.WithTemperature(temp))
}
// WebSocket Handler 中的调用
func (c *Client) handleQuery(req QueryRequest) {
opts := []compose.Option{}
if req.Model != "" {
opts = append(opts, WithModelName(req.Model))
}
output, err := c.pipeline.Stream(ctx, PipelineInput{...}, opts...)
}
4.5 注意事项
值类型 vs 指针类型
Graph 泛型参数必须使用值类型(PipelineInput/PipelineOutput),所有 Lambda 的输入输出也使用值类型。框架在 Transform 模式下会自动处理 T 和 StreamReader[T] 的转换。
Callback 运行时传入
Callback 通过 Stream() 的 option 传入,不在 Compile() 时注册:
streamReader, err := runnable.Stream(ctx, input, compose.WithCallbacks(handler))
eino-ext 与 DashScope 兼容性
eino-ext OpenAI ChatModel 通过 BaseURL 对接 DashScope 兼容接口。需注意:
- 多模态图片使用
Base64Data+MIMEType格式 Timeout控制单次请求超时- 流式输出通过
Stream()方法获取StreamReader[*schema.Message]
框架自动类型转换
Eino 框架在编排场景中自动处理以下转换:
- T → StreamReader[T]:将完整值装箱为单帧流(非流式 → 假流式)
- StreamReader[T] → T:将流 concat 为完整值(流式 → 非流式)
这使得不同流式模式的节点可以无缝连接。
5. 测试策略
5.1 单元测试
func TestPipelineGraph_WithTextInput(t *testing.T) {
mockLLM := &mockChatModel{responses: []string{"你好!"}}
mockSender := &mockSender{}
graph, err := NewPipelineGraph(ctx, &GraphOption{
ChatModel: mockLLM,
Sender: mockSender,
})
require.NoError(t, err)
output, err := graph.Invoke(ctx, PipelineInput{
Text: "你好",
SessionID: "test-session",
})
require.NoError(t, err)
assert.Equal(t, "你好!", output.FullResponse)
assert.True(t, mockSender.LLMDoneSent)
}
func TestPipelineGraph_WithAudioInput(t *testing.T) {
mockSTT := &mockSTT{text: "你好"}
mockLLM := &mockChatModel{responses: []string{"你好!"}}
mockTTS := &mockTTS{audio: []byte("fake-audio")}
mockSender := &mockSender{}
graph, _ := NewPipelineGraph(ctx, &GraphOption{
ChatModel: mockLLM,
STTService: mockSTT,
TTSService: mockTTS,
Sender: mockSender,
})
output, err := graph.Invoke(ctx, PipelineInput{
AudioData: []byte("fake-audio-data"),
SessionID: "test-session",
})
require.NoError(t, err)
assert.True(t, mockSender.TTSAudioSent)
}
5.2 集成测试
- 启动真实 OpenAI API 调用(使用测试 key)
- 验证 WebSocket 消息序列:
stt_result→llm_chunk× N →llm_done→tts_audio× N - 验证 interrupt 取消功能
- 验证多并发请求隔离
6. 未来扩展路径
基于 Eino Graph 的重构完成后,可无缝扩展:
- ReAct Agent:Graph 添加 Branch 节点,实现 LLM → Tool → LLM 循环
- 多模态理解:添加视觉分析 Lambda 节点(图像描述 → 上下文注入)
- Model Router:Graph 前置分支节点,按场景/成本路由不同 LLM
- Rate Limiter:通过 Callback 的 OnStart 实现令牌桶
- Checkpoint/Resume:利用 Eino 的 CheckpointStore 实现断点续传
- Multi-Agent:利用 ADK 的 Supervisor/SequentialAgent 编排复杂对话流程
附录:关键 Eino API 参考
// 构建 Graph
g := compose.NewGraph[I, O](opts...)
g.AddChatModelNode(key, chatModel)
g.AddLambdaNode(key, lambda, opts...)
g.AddEdge(from, to)
g.AddBranch(from, branchFunc, mapping)
// 编译
runnable, err := g.Compile(ctx, opts...)
// 执行四种模式
output, err := runnable.Invoke(ctx, input, opts...)
stream, err := runnable.Stream(ctx, input, opts...)
output, err := runnable.Collect(ctx, inputStream, opts...)
stream, err := runnable.Transform(ctx, inputStream, opts...)
// Lambda 四种构造器
lambda := compose.InvokableLambda(fn) // I → O
lambda := compose.StreamableLambda(fn) // I → StreamReader[O]
lambda := compose.CollectableLambda(fn) // StreamReader[I] → O
lambda := compose.TransformableLambda(fn) // StreamReader[I] → StreamReader[O]
// Stream 操作
sr, sw := schema.Pipe[T](bufSize)
sw.Send(chunk, err)
chunk, err := sr.Recv()
sw.Close()
// Option
compose.WithCallbacks(handler)
compose.WithCallbacks(handler).DesignateNode("node_key")
compose.WithChatModelOption(model.WithTemperature(0.7))
compose.WithGenLocalState(genFunc)