流式深水区:StreamReader 内部
Day 03 学了流的"用法",今天潜入它的"实现"——Pipe 缓冲、背压、Copy 引用计数、多路合并、自动拼接。理解流式为什么能既快又安全。
回顾与今日目标
StreamReader 像一次性水管,Recv() 取走一块就没了、只能被读一次。但真实场景里,同一个模型输出流常常要"喂给下一个节点 + 同时给 callbacks 埋点 + 同时推给前端"——三个下游都想读同一条流。另一头,一个"只要完整值"的节点却收到了碎成 20 段的流。这两个矛盾(一个流多处消费 和 碎片拼回完整)就是今天要解决的。Day 03 记住的:StreamReader 像一次性水管,Recv() 取一块、取完 io.EOF、只能读一次、Copy 能分叉。今天回答"它内部怎么做到的",代码在 schema/stream.go。
stream.go:261)把一条流"一分多"给多个下游各读各的;Concat(message.go:1644 的 ConcatMessages、agentic_message.go:901 的 ConcatAgenticMessages)把流里的碎片"多合一"拼回完整值。有了它们,框架才能在"流"和"非流"之间随意转换,组件作者只实现一两种范式即可。Pipe 内部:一个带缓冲的管子
schema/stream.go:99 的 Pipe(cap) 造出一对"写端 + 读端",中间是一个容量为 cap 的缓冲队列:
Send()→ 块块块→ Reader
Recv()
sr, sw := schema.Pipe[*Message](cap) // 造一对读写端,缓冲容量 cap
go func() { // 生产者(比如模型 SDK)在另一个 goroutine
sw.Send(chunk1, nil); sw.Send(chunk2, nil)
sw.Close() // 生产完,关闭写端 → 读端收到 EOF
}()
for { chunk, err := sr.Recv(); if err==io.EOF {break}; 用(chunk) } // 消费者
背压:缓冲满了写端就等
缓冲容量有限。如果生产者飞快、消费者很慢,缓冲会满——这时 Send 会阻塞等待,直到消费者取走一些腾出空间。这叫 背压(back-pressure)。
cap 是个权衡:太小 → 频繁阻塞、并发度低;太大 → 占内存、背压迟钝。Eino 内部按场景选合适的缓冲大小。Copy 的引用计数(真实代码)
Day 03 说 Copy(n)(stream.go:261)把一个流分成 n 个独立流。它不是真的复制数据 n 份,而是用一个共享的底层 + 引用计数:
sync.Once 保证每个块只从底层拉取一次(不会因为 n 个 reader 各拉一次而重复消费源)。所以 Copy 很轻——不复制数据,只维护"谁读到哪了"。👶 小白:Copy 说是"分成 n 个独立流",又说"不复制数据"。既然数据只有一份,两个下游各读各的,谁读走了另一个不就少一块吗?这不矛盾吗?
👨🏫 老师:不矛盾,关键在"读走"这个词。接着n 个人共看一本连载的比方——小说确实只印一份(底层那条块链表),但"读到第几章"是每个人自己的书签(各自的读取位置指针)。甲翻到第 3 章,不会把书页撕走,乙的书签还停在第 1 章,照样能从头翻。
👨🏫 所以"独立"独立在各自的进度指针,"不复制"不复制的是那份底层块数据。某一块只有当所有 reader 的书签都越过它,才会被释放(引用计数归零);而 sync.Once 保证每块只从源头真正拉取一次,不会因为 n 个人各拉一次而把上游源重复消费。这就是为什么 Copy 很轻——它只加了几个书签,没多印一本书。
多路合并 merge
Day 08 的并行分支会产生多个流,需要合并成一个。stream.go 的 merge 把 n 个 StreamReader 合成一个——从任意一个有数据的流里取,谁先来先出。
reflect.Select 或多路监听,从 n 个源流里"谁准备好了就读谁",汇成一个输出流。全部源流 EOF 了,合并流才 EOF。这让"并行调 3 个检索器、结果流式合并返回"成为可能——三个流的块交错着流出来,而不用等某一个全部完成。channel(Day 09)在有多个上游时就用它来合并流。自动拼接 concat
Day 06 讲过:一个只会 Invoke(要完整值)的节点,遇到上游给的流,框架自动把流收集拼接成完整值再喂它。拼接逻辑就是 Day 03 的 ConcatMessages(message.go:1644)等 concat 函数。
// stream.go 里:把流全收下来,再按类型拼接成一个完整值
func concatStreamReader[T](sr *StreamReader[T]) (T, error) {
var chunks []T
for { c, err := sr.Recv(); if err==io.EOF {break}; chunks=append(chunks,c) }
return concat(chunks) // Message → ConcatMessages 拼成完整消息
}
schema 的 init(Day 03 提过)。get_weather({"city":"北京","days":3}) 时,流式返回常把这段 JSON 拆成 {"ci → ty":"北 → 京","da → ys":3} 好几块。单看任何一块都不是合法 JSON,没法解析。框架在节点边界调 ConcatMessages(schema/message.go:1644)把这些块按 tool_calls 的 index 归并、字符串拼接,还原出完整参数再交给工具。AgenticMessage(Claude 那种带 content block 的消息)走的是 ConcatAgenticMessages(schema/agentic_message.go:901),内部再靠 concatChunksOfSameContentBlock(agentic_message.go:1137)把同一个内容块的碎片拼齐。所以你在业务里拿到的永远是完整参数,分块细节被 concat 藏起来了。关闭与资源释放
流用完必须关。读端 Close() 释放资源;写端 Close() 通知读端 EOF。忘了关会泄漏 goroutine 和内存。
defer sr.Close()。即使你没读完就想提前退出(比如用户取消),Close 也会正确清理,让上游生产者的 goroutine 收到信号停止。Eino 内部各节点都严格 defer Close,你自己消费顶层流时也务必如此。这是流式代码不泄漏的铁律。今日小结 + 动手
🧠 今天你应该能回答
- Pipe 内部是什么?(带缓冲的写端+读端,goroutine 解耦)
- 背压是什么?保护了什么?
- Copy 为什么很轻?(共享底层 + 引用计数 + 读位置)
- merge 和 concat 各干什么?
- 为什么必须 Close?不 Close 会怎样?
✋ 动手
sed -n '99,130p' schema/stream.go # Pipe
sed -n '195,320p' schema/stream.go | head -60 # Recv / Copy
grep -n 'concat\|merge\|Close' schema/stream.go | head
Runnable[I,O]、Graph[I,O])保证编译期类型检查。今天讲泛型怎么用、类型擦除在哪发生、编译期检查如何前移错误。