Day 17 / 共 20 天 · 第 4 周 横切与生态

流式深水区:StreamReader 内部

Day 03 学了流的"用法",今天潜入它的"实现"——Pipe 缓冲、背压、Copy 引用计数、多路合并、自动拼接。理解流式为什么能既快又安全。

📍 你在整门课的位置 · 第 4 周 进阶与生态(共 4 周 · 20 天)
D16 callbacks D17 流式深水 D18 类型/泛型 D19 eino-ext D20 构建收官
L01

回顾与今日目标

🤔 痛点:一个流只能读一次,可下游有好几个都要用它,怎么办? 流式最反直觉的地方:StreamReader 像一次性水管,Recv() 取走一块就没了、只能被读一次。但真实场景里,同一个模型输出流常常要"喂给下一个节点 + 同时给 callbacks 埋点 + 同时推给前端"——三个下游都想读同一条流。另一头,一个"只要完整值"的节点却收到了碎成 20 段的流。这两个矛盾(一个流多处消费碎片拼回完整)就是今天要解决的。

Day 03 记住的:StreamReader 像一次性水管,Recv() 取一块、取完 io.EOF、只能读一次、Copy 能分叉。今天回答"它内部怎么做到的",代码在 schema/stream.go

💡 本质:Copy 让一个流被多个下游消费,Concat 把碎片拼回完整——这是 4 范式自动转换的底层 Day 06 的"4 种流式范式能自动互转"不是魔法,靠的正是今天这两块地基:Copy(stream.go:261把一条流"一分多"给多个下游各读各的;Concat(message.go:1644ConcatMessagesagentic_message.go:901ConcatAgenticMessages把流里的碎片"多合一"拼回完整值。有了它们,框架才能在"流"和"非流"之间随意转换,组件作者只实现一两种范式即可。
为什么值得深入? 流式是 Eino 的一等公民(Day 01),几乎所有 Agent 交互都靠它。搞懂内部机制,你才能:写出不阻塞的流式代码、正确处理 Copy、排查"流读不到/读重了/卡死"的诡异 bug。这是从"会用"到"精通"的一课。
L02

Pipe 内部:一个带缓冲的管子

schema/stream.go:99Pipe(cap) 造出一对"写端 + 读端",中间是一个容量为 cap 的缓冲队列:

Writer
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),消费者从里取(Recv),两边各跑各的(不同 goroutine)。模型一边生成 token 一边 Send,你一边 Recv 一边往前端推——真正的边生成边消费。
L03

背压:缓冲满了写端就等

缓冲容量有限。如果生产者飞快、消费者很慢,缓冲会——这时 Send阻塞等待,直到消费者取走一些腾出空间。这叫 背压(back-pressure)

背压保护了什么? 想象没有背压:生产者疯狂生成,全堆在内存里,消费者来不及处理 → 内存爆炸。背压让"下游处理不过来时,自动让上游慢下来",内存占用被缓冲容量限制住,稳定可控。这是流式系统不 OOM 的关键。底层用 Go 的 channel 天然实现了背压(channel 满则发送阻塞)。
所以选 cap 是个权衡:太小 → 频繁阻塞、并发度低;太大 → 占内存、背压迟钝。Eino 内部按场景选合适的缓冲大小。
L04

Copy 的引用计数(真实代码)

Day 03 说 Copy(n)stream.go:261)把一个流分成 n 个独立流。它不是真的复制数据 n 份,而是用一个共享的底层 + 引用计数:

惰性链表 + sync.Once:Copy 出的 n 个 reader 共享同一条"块链表",每个 reader 有自己的读取位置指针。某个块被所有 reader 都读过了,才会被释放。用 sync.Once 保证每个块只从底层拉取一次(不会因为 n 个 reader 各拉一次而重复消费源)。所以 Copy 很轻——不复制数据,只维护"谁读到哪了"。
直觉 像 n 个人共看一部连载:小说只印一份(底层块链表),每人一个书签记自己看到第几章(读取位置)。所有人都看过的章节可以从桌上收走(释放)。没人重复印书,但每个人进度独立。这就是 Day 16 流式回调能"回调读副本、业务读原本"的底层原理。

👶 小白:Copy 说是"分成 n 个独立流",又说"不复制数据"。既然数据只有一份,两个下游各读各的,谁读走了另一个不就少一块吗?这不矛盾吗?

👨‍🏫 老师:不矛盾,关键在"读走"这个词。接着n 个人共看一本连载的比方——小说确实只印一份(底层那条块链表),但"读到第几章"是每个人自己的书签(各自的读取位置指针)。甲翻到第 3 章,不会把书页撕走,乙的书签还停在第 1 章,照样能从头翻。

👨‍🏫 所以"独立"独立在各自的进度指针,"不复制"不复制的是那份底层块数据。某一块只有当所有 reader 的书签都越过它,才会被释放(引用计数归零);而 sync.Once 保证每块只从源头真正拉取一次,不会因为 n 个人各拉一次而把上游源重复消费。这就是为什么 Copy 很轻——它只加了几个书签,没多印一本书。

① Copy(n):一条流「一分多」——多个下游各读各的 源 StreamReader 副本1 → 下一个节点 副本2 → callbacks 埋点 副本3 → 推前端 共享底层块链表 + 各自读取位置,不复制数据 ② Concat:碎片「多合一」——拼回完整值喂给非流节点 你好 ,世 界! …(20 段碎片) ConcatMessages 完整 Message「你好,世界!」 tool_call 参数被分块吐出时,也在这里拼成完整 JSON
上:Copy 把一条流分给多个下游(不复制数据);下:Concat 把碎片拼回一个完整值。
L05

多路合并 merge

Day 08 的并行分支会产生多个流,需要合并成一个。stream.go 的 merge 把 n 个 StreamReader 合成一个——从任意一个有数据的流里取,谁先来先出。

实现上用 Go 的 reflect.Select 或多路监听,从 n 个源流里"谁准备好了就读谁",汇成一个输出流。全部源流 EOF 了,合并流才 EOF。这让"并行调 3 个检索器、结果流式合并返回"成为可能——三个流的块交错着流出来,而不用等某一个全部完成。channel(Day 09)在有多个上游时就用它来合并流。
L06

自动拼接 concat

Day 06 讲过:一个只会 Invoke(要完整值)的节点,遇到上游给的流,框架自动把流收集拼接成完整值再喂它。拼接逻辑就是 Day 03 的 ConcatMessagesmessage.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 拼成完整消息
}
读法:把流里所有块收进 slice,再调对应类型的拼接函数合成一个完整值。Message 的拼接会把内容片段接起来、合并 tool_calls 片段(Day 03 讲的模型分块吐 tool_call 参数,在这里拼完整)。类型→拼接函数的注册表在 schema 的 init(Day 03 提过)。
📝 例子:模型分块吐 tool_call,靠 concat 拼回完整参数 模型要调 get_weather({"city":"北京","days":3}) 时,流式返回常把这段 JSON 拆成 {"city":"北京","days":3} 好几块。单看任何一块都不是合法 JSON,没法解析。框架在节点边界调 ConcatMessagesschema/message.go:1644)把这些块按 tool_calls 的 index 归并、字符串拼接,还原出完整参数再交给工具。AgenticMessage(Claude 那种带 content block 的消息)走的是 ConcatAgenticMessagesschema/agentic_message.go:901),内部再靠 concatChunksOfSameContentBlockagentic_message.go:1137)把同一个内容块的碎片拼齐。所以你在业务里拿到的永远是完整参数,分块细节被 concat 藏起来了。
这解释了 Day 06 的"一个非流节点冻结全链" 因为 concat 必须等流全部到齐才能拼成完整值。所以链路里只要有一个"只要完整值"的节点,流到它这里就被"蓄水成池",前面的流式优势断了。设计流水线时,尽量让流能一路贯通到最后(前端)。
L07

关闭与资源释放

流用完必须关。读端 Close() 释放资源;写端 Close() 通知读端 EOF。忘了关会泄漏 goroutine 和内存。

最佳实践:拿到 StreamReader 后立刻 defer sr.Close()。即使你没读完就想提前退出(比如用户取消),Close 也会正确清理,让上游生产者的 goroutine 收到信号停止。Eino 内部各节点都严格 defer Close,你自己消费顶层流时也务必如此。这是流式代码不泄漏的铁律。
为什么不 Close 会泄漏? 因为生产者可能是个在后台 goroutine 里跑的循环(不停 Send)。如果你不读也不 Close,它可能永远卡在"缓冲满了等人取"的 Send 上——goroutine 永远不退出,内存永远不释放。Close 就是给它一个"别发了,收工"的信号。
L08

今日小结 + 动手

🧠 今天你应该能回答

  • 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
明天预告 · Day 18类型安全与 Go 泛型——Eino 大量用泛型(Runnable[I,O]Graph[I,O])保证编译期类型检查。今天讲泛型怎么用、类型擦除在哪发生、编译期检查如何前移错误。
← Day 16 callbacks Day 18 · 类型安全 →