流式传输引擎:Eino StreamReader 源码拆解(第61篇-E47)

流式传输引擎:Eino StreamReader 源码拆解(第61篇-E47) 系列「企业级 AI Agent 实现拆解」E47 篇Part 10 生产工程篇第五章。上一篇 讲了 RAG 流水线。这篇往下看底层LLM 的流式输出在 Eino 内部是怎么流动的——StreamReader 如何实现 fan-out、fan-in、类型转换以及 Graph 层如何包装它。读完这篇你会知道Pipe[T]怎么创建一对 StreamReader/StreamWriterStreamReader 的 5 种内部类型是什么各自解决什么问题Copy(n)fan-out 的懒读链表实现cpStreamElementsync.OnceMergeStreamReadersfan-in 和MergeNamedStreamReaders的区别StreamReaderWithConvert类型转换 ErrNoValue过滤 WithOnEOFSetAutomaticCloseGC finalizer 兜底但不是 Close 的替代Graph 层的streamReaderPacker是什么写流式代码的 4 条实用规则为什么流式传输在 Agent 里是个难题LLM 流式输出SSE/Server-Sent Events是用户体验的关键——用户看到第一个 Token 的时间比完整回答更重要。但在 Agent 内部一个流式输出往往需要同时被多个消费者读取主流程继续传给下游节点Callback 系统同时统计 Token客户端 SSE 输出同时推给浏览器这就是流式 fan-out一个流分叉给多个读者的需求。反过来MergeStreamReaders解决 fan-in——多个并行节点的流要合并给下游。Go 的 channel 天然读一次即消费不支持多读。Eino 在schema/stream.go里实现了一套完整的流式抽象来解决这些问题。5 种内部类型StreamReader[T]是个接口外壳内部有 5 种具体实现// schema/stream.go 里的 5 种 reader 类型typestream[T any]// 基础类型底层是 channeltypearrayReader[T any]// 数组转流把 []T 包成流非流式组件兼容用typemultiStream[T any]// fan-in合并多个 StreamReadertypewithConvert[T,O any]// 类型转换T → O可过滤typecpStreamReader[T any]// Copy 产生的子 readerfan-out外部代码只看到*StreamReader[T]具体类型由工厂函数决定。Pipe[T]创建一对读写端// 创建一个容量为 cap 的流reader,writer:schema.Pipe[string](32)// 写端通常在 goroutine 里gofunc(){writer.Send(hello,nil)// 发送数据writer.Send(world,nil)writer.Close()// 关闭写端}()// 读端for{chunk,err:reader.Recv()iferrio.EOF{break}iferr!nil{/* 处理错误 */break}fmt.Print(chunk)}reader.Close()// 读完必须关闭Pipe[T]内部创建的stream[T]结构typestream[T any]struct{itemschanstreamItem[T]// 带缓冲的数据通道容量 cap 参数closedchanstruct{}// 关闭信号}typestreamItem[T any]struct{val T errerror// 数据和错误复用同一个 chan}Send(val, err)就是往items chan里放streamItem。Recv()是取出来。Close()关闭closed chan并消费掉items chan里剩余的数据防止写端 goroutine 阻塞。关键约定StreamReader 是 read-once 的——和 Go channel 一样Recv 一次数据就消费掉了不能回放。如果需要多个消费者必须用Copy。Copy(n)fan-out 的懒读实现// Copy 把一个 reader 分裂成 n 个独立的 readerreaders:reader.Copy(3)// readers[0], readers[1], readers[2]// 原 reader 调用 Copy 后不可再用Copy内部用懒读链表实现不是直接复制 channel// 每个 item 是一个链表节点typecpStreamElement[T any]struct{val T errerrornext*cpStreamElement[T]// 指向下一个节点once sync.Once// 保证只从原 stream 读一次}每个子 readercpStreamReader持有一个指向当前位置的链表节点指针。多个子 reader 共享同一个链表但各自维护自己的当前节点。当子 reader 调用Recv()时查看当前节点是否已有数据once已执行没有 → 用sync.Once从原 stream 读一次写入节点并设置next指针已有 → 直接读取移动到next这样最慢的子 reader 决定原 stream 的读取速度最快的 reader 可以缓存在链表节点里等着不会丢数据。// cpStreamReader.Recv() 的核心逻辑func(c*cpStreamReader[T])Recv()(T,error){cur:c.cur// 当前链表节点// sync.Once 保证只从原 stream 读一次cur.once.Do(func(){cur.val,cur.errc.parent.src.Recv()// 从原 stream 读cur.nextcpStreamElement[T]{}// 为下一个 item 创建节点})c.curcur.next// 移动到下一个节点returncur.val,cur.err}Copy后的子 reader 独立Close()。当所有子 reader 都关闭了parentStreamReader才会关闭原 stream——如果有子 reader 没关闭原 stream 就一直泄漏。MergeStreamReadersfan-in// 把多个 reader 合并成一个merged:schema.MergeStreamReaders([]*StreamReader[string]{r1,r2,r3})// merged.Recv() 非确定性地从 r1/r2/r3 中读取谁有数据先返回谁for{chunk,err:merged.Recv()iferrio.EOF{break}fmt.Println(chunk)}内部实现启动 N 个 goroutine每个负责从一个 reader 读数据全部往同一个output chan里写。merged.Recv()就是从output chan取。顺序不保证——多个 LLM 并行调用的输出顺序是 race 决定的用于拿到就处理的场景。MergeNamedStreamReaders知道数据来自哪里typeNamedStreamReader[T any]struct{NamestringReader*StreamReader[T]}merged:schema.MergeNamedStreamReaders([]NamedStreamReader[string]{{Name:model-a,Reader:r1},{Name:model-b,Reader:r2},})// 每当某个 reader 读完会发出一个 SourceEOF 错误chunk,err:merged.Recv()iferrors.Is(err,schema.SourceEOF){// 某个 reader 结束了但 merged 还没结束// err.(*SourceEOFError).Source 是哪个 Name}SourceEOF让你知道 r1 先结束了而 r2 还在继续——适合需要知道每个子流各自结束的场景比如并行工具调用的结果流聚合。StreamReaderWithConvert类型转换 过滤// 把 *StreamReader[string] 转成 *StreamReader[int]取字符串长度intReader:schema.StreamReaderWithConvert(stringReader,func(sstring)(int,error){ifs{return0,schema.ErrNoValue// 过滤掉空字符串}returnlen(s),nil},)ErrNoValue是过滤哨兵——转换函数返回它时Recv()不返回这个值自动跳过继续读下一个。适合从流里过滤无关的 Token 或格式化标记。其他可选参数schema.StreamReaderWithConvert(reader,fn,schema.WithOnEOF(func(){/* 流结束时执行一次 */}),schema.WithErrWrapper(func(errerror)error{returnfmt.Errorf(convert stage: %w,err)}),)WithOnEOF用于在流结束时做清理比如关闭资源WithErrWrapper包装错误增加上下文不破坏原始错误链。ErrNoValue / ErrRecvAfterClosed / SourceEOFEino 流式系统的三个哨兵错误错误含义谁产生io.EOF流正常结束stream[T].Recv()写端 Close 后schema.ErrNoValue过滤信号跳过这个 item转换函数返回schema.ErrRecvAfterClosed对已关闭的 reader 继续调用 Recv内部 panic 保护通常是 bugschema.SourceEOF某个子流结束named merge 用MergeNamedStreamReaders用errors.Is检查不要直接 io.EOF——因为这些错误可能被WithErrWrapper包装过。SetAutomaticCloseGC finalizer 兜底schema.SetAutomaticClose(reader)这行代码给reader注册了一个runtime.SetFinalizer——当 reader 被 GC 回收时自动调用Close()。但这不是 Close 的替代GC 时机不确定finalizer 可能延迟很久才触发底层 goroutine如 MergeStreamReaders 启动的那些在这段时间里一直跑着消耗资源有些 reader 类型的 Close 需要通知上游finalizer 时序无法保证SetAutomaticClose是兜底安全网防止某些代码路径忘记 Close 导致永久泄漏。它不等于你可以不写defer reader.Close()。正确模式始终是reader:schema.Pipe[string](32)deferreader.Close()// 保证 Close哪怕中间 paniccompose 层的 streamReaderPackerGraph 层compose/stream_reader.go在 StreamReader 外面套了一层streamReaderPacker[T]// compose/stream_reader.gotypestreamReaderPacker[T any]struct{s*schema.StreamReader[T]}它的主要作用是实现compose.StreamReader[T]接口Graph 内部用的把schema.StreamReader适配成 Graph 节点能直接处理的形式。Graph 的流式执行路径节点 A 产出*schema.StreamReader[Output]streamReaderPacker包装它传给下游节点 B 作为流式输入节点 B 通过Recv()逐 Token 消费对于需要 fan-out 的情况节点输出同时去往多个下游 CallbackGraph 内部会在节点结束后自动调用Copy(n)分发。4 条实用规则1. Recv 只能调用一次每次Recv消费一个 item不可回放。如果需要重播在第一次消费前把数据存下来。2. 总是要 Close无论正常读完还是中途 break都要调用reader.Close()。最安全的写法deferreader.Close()for{chunk,err:reader.Recv()iferr!nil{break}// io.EOF 也在这里退出// ...}3. 需要 fan-out先 Copy再 Recv// 错误先 Recv 消费了Copy 就拿不到这个 item 了chunk,_:reader.Recv()copies:reader.Copy(2)// copies 会丢失第一个 item// 正确先 Copy再各自 Recvcopies:reader.Copy(2)goconsume(copies[0])// callback handlerconsume(copies[1])// main path4. 流式 Callback 必须异步消费E45 讲过这里再强调OnEndWithStreamOutput里收到的是Copy的副本必须在 goroutine 里异步消费并调用CloseOnEndWithStreamOutput:func(ctx,info,output)context.Context{gofunc(){deferoutput.Close()// 忘记这行 goroutine 泄漏for{chunk,err:output.Recv()iferr!nil{break}}}()returnctx},小结Eino StreamReader 的设计核心是一次性读取 显式所有权机制作用Pipe[T]/stream[T]基础读写对channel-backed容量可配Copy(n)/ cpStreamElementfan-out懒读链表不复制 channelsync.Once 保证原子读MergeStreamReadersfan-in非确定性合并多 goroutine 竞速MergeNamedStreamReaders带来源标识的 fan-inSourceEOF 知道子流各自结束StreamReaderWithConvert类型转换 ErrNoValue 过滤 OnEOF 钩子SetAutomaticCloseGC finalizer 兜底不替代显式 ClosestreamReaderPackercompose 层适配器让 Graph 节点无缝衔接流式输出理解了这套机制就能理解 Eino Graph 里的流式数据为什么能在节点间、Callback 之间无缝传递而不丢数据、不发生 goroutine 泄漏。代码来源eino/schema/stream.go · eino/compose/stream_reader.go