流式传输引擎: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

相关新闻

零基础玩转bWAPP靶场(十六):SQL 注入(POST/选择型)

零基础玩转bWAPP靶场(十六):SQL 注入(POST/选择型)

摘要:本文是 bWAPP 靶场系列的第十六篇,聚焦于 SQL Injection (POST/Select)(POST 型下拉选择框 SQL 注入)。页面看起来和第十四篇(GET/Select)一模一样,区别只是请求方式从 GET 变成了 POST。我…

2026/7/24 14:01:51 阅读更多 →
AI开发C语言应用按步走,表达式计算器calc的第十一步,Token 联合体、哈希表满报错、批处理变量赋值

AI开发C语言应用按步走,表达式计算器calc的第十一步,Token 联合体、哈希表满报错、批处理变量赋值

calc11 — Token 联合体、哈希表满报错、批处理变量赋值 1. 概述 本次迭代基于 suggestions.md 中的建议,完成了三项代码改进: Token 联合体 — 将 value 和 name 字段合并为联合体,减少内存占用哈希表满报错 — 变量符号表满时输出错误提示&…

2026/7/24 14:01:51 阅读更多 →
基于YOLOv11的焊接缺陷检测系统设计与优化

基于YOLOv11的焊接缺陷检测系统设计与优化

1. 项目背景与核心价值 焊接缺陷检测一直是工业质检领域的痛点问题。传统人工目检效率低下且容易漏检,而基于机器视觉的自动化方案往往面临复杂场景适应性差的问题。这个毕业设计项目选择基于YOLOv11构建焊接缺陷检测系统,正是瞄准了这一行业刚需。 我去…

2026/7/24 14:01:51 阅读更多 →

最新新闻

AM574x异构SoC硬件调试:JTAG与TPIU时序配置与工程实践

AM574x异构SoC硬件调试:JTAG与TPIU时序配置与工程实践

1. 项目概述与核心价值在嵌入式系统开发,尤其是涉及像TI AM574x这类集成了双核Cortex-A15、双C66x DSP以及多个协处理器的复杂异构SoC时,硬件级的调试与跟踪能力不再是“锦上添花”,而是“雪中送炭”的必需品。想象一下,当你的系统…

2026/7/24 14:10:55 阅读更多 →
Java 企业级 SaaS 架构 B2C 微信小程序电商项目实战训练营07-全栈测试部署与运维实战

Java 企业级 SaaS 架构 B2C 微信小程序电商项目实战训练营07-全栈测试部署与运维实战

文章目录 一、概述 二、部署架构总览 2.1 生产环境拓扑 2.2 服务清单 三、后端构建与打包 3.1 Maven Profile 多环境管理 3.2 构建命令 3.3 构建产物分析 3.4 生产环境配置 3.5 启动命令 四、前端构建与部署 4.1 管理端构建 4.2 店铺端构建 4.3 移动端 H5 构建 4.4 微信小程序发…

2026/7/24 14:10:55 阅读更多 →
【具身智能】Claude Sonnet 5写IMU代码竟有“人格分裂”?用《旋生万物》公理 I2=−N 根治Agent几何幻觉(附螺旋积分器)

【具身智能】Claude Sonnet 5写IMU代码竟有“人格分裂”?用《旋生万物》公理 I2=−N 根治Agent几何幻觉(附螺旋积分器)

摘要:2026年7月,CSDN首页热议“TVA-具身智能如何跨越电子与原子鸿沟”。然而实测发现,被称为“最懂代码”的Claude Sonnet 5在编写IMU(惯性测量单元)姿态解算代码时,频繁出现2π相位跳变和螺距漂移——仿佛…

2026/7/24 14:10:55 阅读更多 →
CC2652RB射频与模拟外设深度解析:从数据手册到设计实战

CC2652RB射频与模拟外设深度解析:从数据手册到设计实战

1. 从数据手册到设计实战:CC2652RB射频与模拟外设深度解析在物联网设备开发中,选型一颗无线MCU时,我们最关心的往往是那几个硬核指标:通信距离有多远?功耗能有多低?传感器数据采得准不准?这些问…

2026/7/24 14:10:55 阅读更多 →
Claude AI编程助手:从自然语言到代码生成的开发效率提升指南

Claude AI编程助手:从自然语言到代码生成的开发效率提升指南

在日常开发中,很多开发者初次接触 Claude 时会产生误解,认为它是一款新型编译器或编程工具。实际上,Claude 是一个基于人工智能的对话助手,它能理解自然语言并协助完成编程任务,这与传统编译器有着本质区别。本文将详细…

2026/7/24 14:10:54 阅读更多 →
大模型接入的认证与计费:多模型供应商的统一网关设计

大模型接入的认证与计费:多模型供应商的统一网关设计

大模型接入的认证与计费:多模型供应商的统一网关设计 一、三个团队各用各的 API Key,月底账单混乱,财务说"分不清谁花了多少钱" 多模型供应商共存的企业场景下(OpenAI Claude 自建 vLLM),接入层…

2026/7/24 14:09:54 阅读更多 →

日新闻

用Highcharts 创建可拖拽三维散点立方体3D图表

用Highcharts 创建可拖拽三维散点立方体3D图表

该案例基于Highcharts scatter3d 三维散点图实现空间立方体散点可视化,核心特色:三维 X/Y/Z 三轴空间,所有散点分布在 0~10 立方体空间内;散点使用径向渐变实现立体 3D 圆球质感;支持鼠标 / 触屏拖拽画布,…

2026/7/24 0:00:29 阅读更多 →
AppCertDlls:进程创建路径上的 DLL 入口

AppCertDlls:进程创建路径上的 DLL 入口

AppCertDlls:进程创建路径上的 DLL 入口 AppCertDlls 位于 HKLM\System\CurrentControlSet\Control\Session Manager\AppCertDlls。本文的程序功能是只读列出这个键在 64 位和 32 位注册表视图中的全部值,并显示每条值的来源、名称、类型和可安全显示的数…

2026/7/24 0:00:29 阅读更多 →
我的编程之路:第一篇博客

我的编程之路:第一篇博客

大家好,我是一名编程初学者,同时这也是我编程学习之路上的第一篇博客。在这里,我想要向大家介绍我的一些想法和规划。a.自我介绍我是一个刚刚接触编程的新手,目前在学习c语言,我对编程世界充满了强烈的好奇。当然&…

2026/7/24 0:00:29 阅读更多 →

周新闻

Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/24 3:59:20 阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/24 1:23:39 阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/23 17:49:47 阅读更多 →

月新闻