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))既存の実装
- OpenAI ChatModel: OpenAI の GPT 系列 ChatModel - OpenAI
- Ollama ChatModel: Ollama のローカルモデル ChatModel - Ollama
- ARK ChatModel: ARK プラットフォームのモデルサービス ChatModel - ARK
自分で実装するときの参考
独自の ChatModel を実装するときは、次に注意:
- 共通 option を実装すること
- callback の仕組みを実装すること
- ストリーム出力では、出し終わったら 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
}