Files
CamTalk/docs/08-Eino框架与编排设计.md
2026-06-21 19:31:17 +08:00

17 KiB
Raw Permalink Blame History

CamTalk Eino 框架与编排设计

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 的核心理由:

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

1.2 旧方案的问题

当前后端 AI 编排层(internal/orchestrator/pipeline.go)为手写 goroutine 管道存在以下问题:

  1. 编排逻辑硬编码STT→LLM→TTS 流程写死,扩展困难
  2. 并发控制粗糙:手动 goroutine 调度,缺乏结构化流式传递
  3. 无回调/AOP 机制:日志、指标、追踪散落各处
  4. 配置耦合模型名、TTS 参数等硬编码在结构体
  5. 错误处理不一致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() 条件路由
  • 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)    // 流式调用

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 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      # 单元测试

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 模式下会自动处理 TStreamReader[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_resultllm_chunk × N → llm_donetts_audio × N
  • 验证 interrupt 取消功能
  • 验证多并发请求隔离

6. 未来扩展路径

基于 Eino Graph 的重构完成后,可无缝扩展:

  1. ReAct AgentGraph 添加 Branch 节点,实现 LLM → Tool → LLM 循环
  2. 多模态理解:添加视觉分析 Lambda 节点(图像描述 → 上下文注入)
  3. Model RouterGraph 前置分支节点,按场景/成本路由不同 LLM
  4. Rate Limiter:通过 Callback 的 OnStart 实现令牌桶
  5. Checkpoint/Resume:利用 Eino 的 CheckpointStore 实现断点续传
  6. 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)