schema 核心数据类型
最底层的 L1。重点两个:Message(一个结构体同时当输入/输出/模板)和 StreamReader(Eino 最精妙的流式抽象)。贴真实源码。
schema 是"数据契约层"
schema/ 定义"数据长什么样"——它是全框架的通用货币。四个核心类型:
| 类型 | 文件 | 是什么 |
|---|---|---|
Message | message.go | 模型的输入输出统一结构(今天重点) |
StreamReader/Writer | stream.go | 流式抽象(今天重点,全框架最精妙) |
ToolInfo | tool.go | 工具的自描述(传给模型的工具定义) |
Document | document.go | RAG 文档(检索四组件间流通的货币) |
Message、检索返回 []Document、流式返回 StreamReader。统一了数据格式,各层才能无缝拼接。就像大家都用同一种货币交易,不用换汇。Message:一个结构体走天下(真实代码)
{role, content, tool_calls}、Claude 又是另一套字段名、工具结果的挂法也各不相同。你要是每接一个模型就写一套结构体和转换代码,编排层根本没法"一视同仁"地处理消息——换个模型,上层全得改。Message 既能表示你发的话、也能表示模型的回复、还能表示工具的结果——靠一个 Role 字段区分。各家模型实现只需在自己内部把私有格式翻译成这一个 Message,上层就统一了。生活类比:
Message 就像快递行业的统一快递单——不管发件的是顺丰、京东还是私人,单子上永远是「寄件人 / 收件人 / 内容」那几栏(对应 Role/Content/...)。分拣中心(编排层)只认这张标准单,不用管货是谁家的。Message(schema/message.go:498)同时承载"用户输入"和"模型输出",兼容纯文本和多模态。核心字段:
type Message struct {
Role RoleType // 谁说的:System/User/Assistant/Tool
Content string // 纯文本内容
// 多模态(新):用户输入 / 模型输出 分开
UserInputMultiContent []MessageInputPart
AssistantGenMultiContent []MessageOutputPart
// 仅 Assistant 消息:模型要发起的工具调用
ToolCalls []ToolCall
// 仅 Tool 消息:把工具执行结果关联回之前某次 tool_call
ToolCallID string
ToolName string
ResponseMeta *ResponseMeta // 元信息:FinishReason、Usage(token 数)
ReasoningContent string // 推理模型的思维链
Extra map[string]any // 各实现自定义扩展位
}
User Message,模型回的话是 Assistant Message,工具跑完的结果是 Tool Message。避免了"请求类型"和"响应类型"两套。[]*Message 再发给模型(就像 claude-code 教程里的消息回填)。用一个统一结构,消息列表里各种角色的消息能自然共存、追加。ToolCalls(模型说"我要调工具")和 ToolCallID(工具结果"回应哪次调用")靠 ID 配对——这就是 Agent 工具循环的数据基础。便捷构造器:SystemMessage() / UserMessage() / AssistantMessage() / ToolMessage()(message.go:1105 起)。
schema.UserMessage("北京天气?") 得到 Message{Role:"user", Content:"北京天气?"}模型要调工具 →
Message{Role:"assistant", Content:"", ToolCalls:[{ID:"call_1", Function:{Name:"get_weather", Arguments:"{\"city\":\"北京\"}"}}]}工具跑完 →
schema.ToolMessage("晴 25℃", "call_1") 得到 Message{Role:"tool", Content:"晴 25℃", ToolCallID:"call_1"}注意
ToolCallID:"call_1" 和上面 assistant 的 ID:"call_1" 对上了——这就是"工具结果回应哪次调用"的配对。四种 Role
RoleType 四种常量(message.go:111):
"你是谁、规矩"
你说的话
含 ToolCalls
含 ToolCallID
[System: 你是助手] [User: 北京天气?] [Assistant: 我要调 get_weather(北京), ToolCalls=...] [Tool: 晴 25℃, ToolCallID=...] [Assistant: 北京今天晴,25℃]——五条消息,四种角色都用上了。模型每次看到整个列表,就知道"聊到哪了、工具返回了啥"。Message 还是"模板"
有意思的是——Message 自己实现了 MessagesTemplate 接口(message.go:573),能当提示词模板用:
Format(ctx, vs map[string]any, formatType FormatType) ([]*Message, error)
支持三种模板语法(message.go:99):FString(Python 风格 {name})、GoTemplate、Jinja2。还有 MessagesPlaceholder(key)(message.go:595)——把"历史消息列表"原样插进模板某处。
[]*Message,能无缝喂给模型。MessagesPlaceholder 尤其重要——它让你在模板里留一个"历史对话插这里"的坑,多轮对话时把之前的消息填进去。这是 Day 04 ChatTemplate 组件的基础。Stream 是什么:边想边吐
大模型不是"憋完整段话再返回",而是一个 token 一个 token 地吐(你看 ChatGPT 打字机效果就是)。Eino 用 StreamReader[T] 抽象这种"流"。
StreamReader[T] 就是"一个能一块块读出 T 的管道"——T 通常是 *Message(每个 chunk 是一小片消息)。泛型 [T] 意思是“装什么类型都行”,这里装 Message 分片。Generate)像宴席大厨:所有菜全做完了才一起端上桌,你干坐着等半天。流式(Stream)像小炒摊:炒好一个菜就先端一个,你边吃边等下一道。StreamReader 就是"上菜口",Recv() 就是你去端下一道菜。好处一样:不用干等(体验好)、厨房也不用一次性摆满所有成品(省内存)。StreamReader / StreamWriter(真实代码)
核心是一对泛型类型(schema/stream.go):StreamWriter[T](生产端,往里写) 和 StreamReader[T](消费端,往外读)。用 Pipe 一次造出一对(stream.go:99):
sr, sw := schema.Pipe[*Message](cap) // 造一读一写,底层是带缓冲 channel
// 生产端(比如模型实现):
sw.Send(chunk, err) // 发一个分片(stream.go:126);返回 closed=true 说明读端已关,可停手
sw.Close() // 发完关闭
// 消费端(比如你的代码):
for {
chunk, err := sr.Recv() // 收一个分片(stream.go:195)
if err == io.EOF { break } // channel 关了 → EOF,读完了
// 用 chunk...
}
sr.Close() // ★ 必须调一次!即使已 EOF 或提前 break
channel(协程间通信的管道)。Recv 后一定要 Close? 因为流底层可能有一个 goroutine(协程)在往里写。如果你读一半就不读了、也不 Close,那个 goroutine 会永远卡在 Send——goroutine 泄漏(内存越占越多)。所以 Eino 的铁律:StreamReader 用完必须 Close 一次。SetAutomaticClose()(stream.go:279)还能借 GC 兜底自动关。channel 是 Go 里“一个协程写、另一个协程读”的线程安全队列——流的底层就是它。break 退出循环、不读了,为什么还要多此一举调 Close?👨🏫 老师:因为流的另一头往往有个后台协程在拼命往里写。你不读也不关,它就会一直卡在
Send 上下不来。👶 小白:卡着能咋样?程序又没崩。
👨🏫 老师:这就是「goroutine 泄漏」——每来一个请求就漏一个卡死的协程,它占的内存永远不释放。跑一天,成千上万个僵尸协程把内存吃光,服务 OOM 崩溃。
👶 小白:那我记不住怎么办?
👨🏫 老师:记铁律「拿到 StreamReader,就
defer sr.Close()」,跟开门要关门一样变成肌肉记忆。实在忘了,SetAutomaticClose() 还能借 GC 兜底。io.EOF(流读完了)就自动清理了、不用 Close。其实:EOF 只代表"没数据了",不代表"资源释放了"——底层协程和 channel 还在。无论读没读完、有没有提前 break,都要 Close 恰好一次。Copy:一份流喂给多个消费者(巧思)
StreamReader 是单消费者的(一个流只能被读一次)。但编排时经常要"同一份输出既喂给下一个节点、又喂给回调 handler"。解决办法是 Copy(n)(stream.go:261)——扇出成 n 个独立 reader,各自都能读到全部元素:
readers := sr.Copy(2) // 变成 2 个独立流,原 sr 作废
// readers[0] 给下一个节点,readers[1] 给回调 handler,互不影响
sync.Once,保证每个上游元素只从父流真正 Recv 一次,后到的子 reader 复用已缓存的元素(stream.go:837 的 peek)。所有子流都 Close 后才真正关父流。既支持多消费者,又不预先把整流塞进内存。这是 Eino 流式设计里最精妙的一处(Day 17 会更深入)。sync.Once = "保证某段代码只执行一次",用来确保每个元素只从父流取一次。流拼接成完整值:4 范式自动转换的地基
Day 01 讲的"Runnable 4 范式自动转换"(流↔非流),底层就靠这里——把一串消息分片拼成一条完整消息。ConcatMessages(message.go:1644):
// 把 [北] [京] [今] [天] [晴] 这些分片 → 拼成一条完整 Message{Content: "北京今天晴"}
full, _ := schema.ConcatMessages(chunks)
// ConcatMessageStream:直接把整条流拼成一条消息(message.go:1842)
full, _ := schema.ConcatMessageStream(stream)
ToolCalls 按 Index 分组合并(把碎片化的工具参数 JSON 拼回完整)、Usage 取最大值、多模态分片合并。关键机制:这些拼接函数在 init()(message.go:40)里通过 RegisterStreamChunkConcatFunc 注册——编排层遇到"需要把流降级成完整值"时(比如下一个节点只要非流输入),就自动按类型调用对应的拼接函数。这就是 Day 01 说的"组件只实现有意义的范式,框架补齐其余"的底层魔法(Day 06/17 深入)。Message 是"数据",StreamReader 是"流动的数据",Concat 是"把流动的数据凝固成完整数据"。这三者构成了 Eino 数据层的全部,也是上面所有编排/Agent 的基石。今日小结 + 动手
🧠 今天你应该能回答
- 为什么 Message 一个结构体同时当输入/输出/模板?
- 四种 Role?一次带工具对话的消息列表长什么样?
- Stream 为什么需要?StreamReader 用完为什么必须 Close?
- Copy 怎么做到"多消费者但不预存整流"?(惰性链表 + sync.Once)
- ConcatMessages 干什么?它怎么支撑"4 范式自动转换"?
✋ 动手:对着真实代码读一遍
# 1. Message 结构与 Role(L02/L03)
sed -n '498,532p' schema/message.go
sed -n '111,120p' schema/message.go
# 2. Stream 的 Pipe/Send/Recv/Copy(L06/L07)
sed -n '99,140p' schema/stream.go
sed -n '261,300p' schema/stream.go
# 3. 拼接函数的注册(L08,4 范式的地基)
sed -n '40,48p' schema/message.go
grep -n "func ConcatMessages\|func ConcatMessageStream" schema/message.go
components 组件接口。ChatModel/Tool/Retriever/ChatTemplate 各是什么接口、有哪些方法。它们是编排的"积木"。