文章へ移動
ShemolEino学習ノート-1-ChatModel
Eino / LLM

Eino学習ノート-1-ChatModel

ChatModel は Eino フレームワークにおける対話 LLM の抽象で、OpenAI や Ollama など異なる大モデルサービスと話すための統一インターフェースを提供する。

このコンポーネントは次の場面で効いてくる:

  • 自然言語の対話
  • テキスト生成と補完
  • ツール呼び出しのパラメータ生成
  • マルチモーダルなやり取り(テキスト、画像、音声など)

コンポーネント定義

インターフェース定義

コード位置:eino/components/model/interface.go
go
type ChatModel interface {
    Generate(ctx context.Context, input []*schema.Message, opts ...Option) (*schema.Message, error)
    Stream(ctx context.Context, input []*schema.Message, opts ...Option) (*schema.StreamReader[*schema.Message], error)
    BindTools(tools []*schema.ToolInfo) error
}

Generate メソッド

  • 機能:モデルの完全な応答を生成する
  • パラメータ:
    • ctx:コンテキスト。リクエストレベルの情報を渡し、Callback Manager も渡す
    • input:入力メッセージのリスト
    • opts:モデルの振る舞いを設定する任意パラメータ
  • 戻り値:
    • *schema.Message:モデルが生成した応答メッセージ
    • error:生成中のエラー情報

Stream メソッド

  • 機能:ストリームでモデル応答を生成する
  • パラメータ:Generate と同じ
  • 戻り値:
    • *schema.StreamReader[*schema.Message]:モデル応答のストリームリーダー
    • error:生成中のエラー情報

BindTools メソッド

  • 機能:モデルに使えるツールをバインドする
  • パラメータ:
    • tools:ツール情報のリスト
  • 戻り値:
    • error:バインド中のエラー情報

コアの位置づけ:このインターフェースは対話モデルの中核抽象で、2つの呼び出しモードを支える:

  • Generate :同期で完全な応答(普通の対話向け)
  • Stream :ストリーム応答(長文生成/リアルタイム向け)

アーキテクチャの特徴

go
type ChatModel interface {
    // 同步生成(典型AI对话模式)
    Generate(ctx context.Context, input []*schema.Message, opts ...Option) (*schema.Message, error)

    // 流式处理(适合逐段输出场景)
    Stream(ctx context.Context, input []*schema.Message, opts ...Option) (
        *schema.StreamReader[*schema.Message], error)

    // 工具绑定机制(支持功能扩展)
    BindTools(tools []*schema.ToolInfo) error
}

設計のポイント :

  • マルチモデル :インターフェース抽象で異なる AI エンジン(OpenAI/MAAS)に対応
  • コンテキスト認識 :context.Context でタイムアウトやトレースなど
  • 拡張パラメータ : ...Option で実装ごとに設定を足せる
  • ツールのホットバインド : BindTools で実行時に機能を足す(Function Calling などを想定)

エンジニアリング実践 :

//go:generate で ChatModelMock のモック実装を自動生成している。つまり:

  • インターフェース優先の設計
  • 単体テストのサポートがしっかりしている
  • 依存性注入(環境ごとのテストがしやすい)

注意点 :

  • 並行安全 :コメントで BindTools と Generate は非アトミックと明記。同期が必要そう
  • メッセージプロトコル :<schema.Message> で定義された形式に依存(具体的なプロトコルと合わせて見る)
  • ストリームのライフサイクル : StreamReader は Close してリソースを解放する

Message 構造体

コード位置:eino/schema/message.go
go
type Message struct {   
    // Role 表示消息的角色(system/user/assistant/tool)
    Role RoleType
    // Content 是消息的文本内容
    Content string
    // MultiContent 是多模态内容,支持文本、图片、音频等
    MultiContent []ChatMessagePart
    // Name 是消息的发送者名称
    Name string
    // ToolCalls 是 assistant 消息中的工具调用信息
    ToolCalls []ToolCall
    // ToolCallID 是 tool 消息的工具调用 ID
    ToolCallID string
    // ResponseMeta 包含响应的元信息
    ResponseMeta *ResponseMeta
    // Extra 用于存储额外信息
    Extra map[string]any
}

Message 構造体はモデルとのやり取りの基本構造で、次を支える:

  • 複数ロール:system(システム)、user(ユーザー)、assistant(ai)、tool(ツール)
  • マルチモーダル:テキスト、画像、音声、動画、ファイル
  • ツール呼び出し:モデルが外部ツールや関数を呼べる
  • メタ情報:応答理由、token 使用量など

共通 Option

Model コンポーネントはモデルの振る舞いを設定する共通 Option を一式提供する:

コード位置:eino/components/model/option.go
go
type Options struct {
    // Temperature 控制输出的随机性
    Temperature *float32
    // MaxTokens 控制生成的最大 token 数量
    MaxTokens *int
    // Model 指定使用的模型名称
    Model *string
    // TopP 控制输出的多样性
    TopP *float32
    // Stop 指定停止生成的条件
    Stop []string
}

Option は次のように設定できる:

go
// 设置温度
WithTemperature(temperature float32) Option

// 设置最大 token 数
WithMaxTokens(maxTokens int) Option

// 设置模型名称
WithModel(name string) Option

// 设置 top_p 值
WithTopP(topP float32) Option

// 设置停止词
WithStop(stop []string) Option

使い方

単体で使う

go
import (
    "context"
    "fmt"
    "io"

    "github.com/cloudwego/eino-ext/components/model/openai"
    "github.com/cloudwego/eino/components/model"
    "github.com/cloudwego/eino/schema"
)

// 初始化模型 (以openai为例)
cm, err := openai.NewChatModel(ctx, &openai.ChatModelConfig{
    // 配置参数
})

// 准备输入消息
messages := []*schema.Message{
    {
       Role:    schema.System,
       Content: "你是一个有帮助的助手。",
    },
    {
       Role:    schema.User,
       Content: "你好!",
    },
}

// 生成响应
response, err := cm.Generate(ctx, messages, model.WithTemperature(0.8))

// 响应处理
fmt.Print(response.Content)

// 流式生成
streamResult, err := cm.Stream(ctx, messages)

defer streamResult.Close()

for {
    chunk, err := streamResult.Recv()
    if err == io.EOF {
       break
    }
    if err != nil {
       // 错误处理
    }
    // 响应片段处理
    fmt.Print(chunk.Content)
}

オーケストレーションで使う

go
import (
    "github.com/cloudwego/eino/schema"
    "github.com/cloudwego/eino/compose"
)

/*** 初始化ChatModel
* cm, err := xxx
*/

// 在 Chain 中使用
c := compose.NewChain[[]*schema.Message, *schema.Message]()
c.AppendChatModel(cm)


// 在 Graph 中使用
g := compose.NewGraph[[]*schema.Message, *schema.Message]()
g.AddChatModelNode("model_node", cm)

Option と Callback の使い方

Option の使用例

go
import "github.com/cloudwego/eino/components/model"

// 使用 Option
response, err := cm.Generate(ctx, messages,
    model.WithTemperature(0.7),
    model.WithMaxTokens(2000),
    model.WithModel("gpt-4"),
)

Callback の使用例

go
import (
    "context"
    "fmt"

    "github.com/cloudwego/eino/callbacks"
    "github.com/cloudwego/eino/components/model"
    "github.com/cloudwego/eino/compose"
    "github.com/cloudwego/eino/schema"
    callbacksHelper "github.com/cloudwego/eino/utils/callbacks"
)

// 创建 callback handler
handler := &callbacksHelper.ModelCallbackHandler{
    OnStart: func(ctx context.Context, info *callbacks.RunInfo, input *model.CallbackInput) context.Context {
       fmt.Printf("开始生成,输入消息数量: %d\n", len(input.Messages))
       return ctx
    },
    OnEnd: func(ctx context.Context, info *callbacks.RunInfo, output *model.CallbackOutput) context.Context {
       fmt.Printf("生成完成,Token 使用情况: %+v\n", output.TokenUsage)
       return ctx
    },
    OnEndWithStreamOutput: func(ctx context.Context, info *callbacks.RunInfo, output *schema.StreamReader[*model.CallbackOutput]) context.Context {
       fmt.Println("开始接收流式输出")
       defer output.Close()
       return ctx
    },
}

// 使用 callback handler
helper := callbacksHelper.NewHandlerHelper().
    ChatModel(handler).
    Handler()

/*** compose a chain
* chain := NewChain
* chain.appendxxx().
*       appendxxx().
*       ...
*/

// 在运行时使用
runnable, err := chain.Compile()
if err != nil {
    return err
}
result, err := runnable.Invoke(ctx, messages, compose.WithCallbacks(helper))

既存の実装

  1. OpenAI ChatModel: OpenAI の GPT 系列 ChatModel - OpenAI
  2. Ollama ChatModel: Ollama のローカルモデル ChatModel - Ollama
  3. ARK ChatModel: ARK プラットフォームのモデルサービス ChatModel - ARK

自分で実装するときの参考

独自の ChatModel を実装するときは、次に注意:

  1. 共通 option を実装すること
  2. callback の仕組みを実装すること
  3. ストリーム出力では、出し終わったら writer を close すること

Option の仕組み

独自 ChatModel が共通 Option 以外も欲しいなら、コンポーネント抽象のヘルパーで独自 Option を作れる。たとえば:

go
import (
    "time"

    "github.com/cloudwego/eino/components/model"
)

// 定义 Option 结构体
type MyChatModelOptions struct {
    Options    *model.Options
    RetryCount int
    Timeout    time.Duration
}

// 定义 Option 函数
func WithRetryCount(count int) model.Option {
    return model.WrapImplSpecificOptFn(func(o *MyChatModelOptions) {
       o.RetryCount = count
    })
}

func WithTimeout(timeout time.Duration) model.Option {
    return model.WrapImplSpecificOptFn(func(o *MyChatModelOptions) {
       o.Timeout = timeout
    })
}

Callback の処理

ChatModel の実装は適切なタイミングでコールバックを飛ばす。以下の構造は ChatModel コンポーネントが定義する:

go
import (
    "github.com/cloudwego/eino/schema"
)

// 定义回调输入输出
type CallbackInput struct {
    Messages    []*schema.Message
    Model       string
    Temperature *float32
    MaxTokens   *int
    Extra       map[string]any
}

type CallbackOutput struct {
    Message    *schema.Message
    TokenUsage *schema.TokenUsage
    Extra      map[string]any
}

完全な実装例

go
import (
    "context"
    "errors"
    "net/http"
    "time"

    "github.com/cloudwego/eino/callbacks"
    "github.com/cloudwego/eino/components/model"
    "github.com/cloudwego/eino/schema"
)

type MyChatModel struct {
    client     *http.Client
    apiKey     string
    baseURL    string
    model      string
    timeout    time.Duration
    retryCount int
}

type MyChatModelConfig struct {
    APIKey string
}

func NewMyChatModel(config *MyChatModelConfig) (*MyChatModel, error) {
    if config.APIKey == "" {
       return nil, errors.New("api key is required")
    }

    return &MyChatModel{
       client: &http.Client{},
       apiKey: config.APIKey,
    }, nil
}

func (m *MyChatModel) Generate(ctx context.Context, messages []*schema.Message, opts ...model.Option) (*schema.Message, error) {
    // 1. 处理选项
    options := &MyChatModelOptions{
       Options: &model.Options{
          Model: &m.model,
       },
       RetryCount: m.retryCount,
       Timeout:    m.timeout,
    }
    options.Options = model.GetCommonOptions(options.Options, opts...)
    options = model.GetImplSpecificOptions(options, opts...)

    // 2. 开始生成前的回调
    ctx = callbacks.OnStart(ctx, &model.CallbackInput{
       Messages: messages,
       Config: &model.Config{
          Model: *options.Options.Model,
       },
    })

    // 3. 执行生成逻辑
    response, err := m.doGenerate(ctx, messages, options)

    // 4. 处理错误和完成回调
    if err != nil {
       ctx = callbacks.OnError(ctx, err)
       return nil, err
    }

    ctx = callbacks.OnEnd(ctx, &model.CallbackOutput{
       Message: response,
    })

    return response, nil
}

func (m *MyChatModel) Stream(ctx context.Context, messages []*schema.Message, opts ...model.Option) (*schema.StreamReader[*schema.Message], error) {
    // 1. 处理选项
    options := &MyChatModelOptions{
       Options: &model.Options{
          Model: &m.model,
       },
       RetryCount: m.retryCount,
       Timeout:    m.timeout,
    }
    options.Options = model.GetCommonOptions(options.Options, opts...)
    options = model.GetImplSpecificOptions(options, opts...)

    // 2. 开始流式生成前的回调
    ctx = callbacks.OnStart(ctx, &model.CallbackInput{
       Messages: messages,
       Config: &model.Config{
          Model: *options.Options.Model,
       },
    })

    // 3. 创建流式响应
    // Pipe产生一个StreamReader和一个StreamWrite,向StreamWrite中写入可以从StreamReader中读到,二者并发安全。
    // 实现中异步向StreamWrite中写入生成内容,返回StreamReader作为返回值
    // ***StreamReader是一个数据流,仅可读一次,组件自行实现Callback时,既需要通过OnEndWithCallbackOutput向callback传递数据流,也需要向返回一个数据流,需要对数据流进行一次拷贝
    // 考虑到此种情形总是需要拷贝数据流,OnEndWithCallbackOutput函数会在内部拷贝并返回一个未被读取的流
    // 以下代码演示了一种流处理方式,处理方式不唯一
    sr, sw := schema.Pipe[*model.CallbackOutput](1)

    // 4. 启动异步生成
    go func() {
       defer sw.Close()

       // 流式写入
       m.doStream(ctx, messages, options, sw)
    }()

    // 5. 完成回调
    _, nsr := callbacks.OnEndWithStreamOutput(ctx, sr)

    return schema.StreamReaderWithConvert(nsr, func(t *model.CallbackOutput) (*schema.Message, error) {
       return t.Message, nil
    }), nil
}

func (m *MyChatModel) BindTools(tools []*schema.ToolInfo) error {
    // 实现工具绑定逻辑
    return nil
}

func (m *MyChatModel) doGenerate(ctx context.Context, messages []*schema.Message, opts *MyChatModelOptions) (*schema.Message, error) {
    // 实现生成逻辑
    return nil, nil
}

func (m *MyChatModel) doStream(ctx context.Context, messages []*schema.Message, opts *MyChatModelOptions, sr *schema.StreamWriter[*model.CallbackOutput]) {
    // 流式生成文本写入sr中
    return
}

参考資料