流式传输引擎: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/8/16 8:09:00 阅读更多 →
AI开发C语言应用按步走,表达式计算器calc的第十一步,Token 联合体、哈希表满报错、批处理变量赋值

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

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

2026/8/8 1:46:01 阅读更多 →
基于YOLOv11的焊接缺陷检测系统设计与优化

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

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

2026/7/30 7:50:07 阅读更多 →

最新新闻

跳出单一改写误区:一文读懂双重优化,AIGCBiye 如何同步实现论文降重与降 AIGC

跳出单一改写误区:一文读懂双重优化,AIGCBiye 如何同步实现论文降重与降 AIGC

随着各大高校陆续将 AIGC 内容检测纳入论文审核标准,越来越多毕业生陷入双重难题:修改后的文稿重复率达标,却在 AI 检测中标记高风险;单纯消除 AI 行文特征,又造成文本相似度反弹。很多创作者至今没有分清,…

2026/8/16 8:08:38 阅读更多 →
【AI开源】ponytail 中文版:让 AI 代理少写无效代码

【AI开源】ponytail 中文版:让 AI 代理少写无效代码

Ponytail 中文版发布:一个让 AI 代理像最懒资深开发那样思考的开源项目。核心不是写更少代码,而是用原生 HTML 元素替代过度封装——真实基准测试中,日期选择器从 404 行降至 23 行,整体代码量平均减少 54%,成本降 20%…

2026/8/16 8:08:38 阅读更多 →
亲测突破夸克网盘下载限制,提高百倍下载速度的方法

亲测突破夸克网盘下载限制,提高百倍下载速度的方法

当我们想要下载夸克网盘里面的一些文件时,当你不想开会员又想下载速度很快的时候,应该怎么办呢?不妨来看看我这方法---》: 下载速度看你的网速和宽带跑个10几M/秒不是问题,亲测有效,接下来就是教程部分 打开…

2026/8/16 8:08:38 阅读更多 →
从零构建高质量文本转语音系统:原理、选型与实战优化指南

从零构建高质量文本转语音系统:原理、选型与实战优化指南

1. 从“看”到“听”:为什么我们需要让文字“发声”?我们生活在一个信息爆炸的时代,每天被海量的文字信息包围——新闻、报告、邮件、电子书、社交媒体动态……眼睛的负担越来越重。你有没有过这样的体验:通勤路上想“读”点东西&…

2026/8/16 8:08:38 阅读更多 →
课程思政元素收集遴选系统-ssm

课程思政元素收集遴选系统-ssm

本项目为前几天收费帮学妹做的一个项目,在工作环境中基本使用不到,但是很多学校把这个当作编程入门的项目来做,故分享出本项目供初学者参考。 一、项目描述 基于ssm课程思政元素收集遴选系统通过Mysql数据库连接数据库 http://localhost:808…

2026/8/16 8:08:38 阅读更多 →
运维常见面试题_05_NFS 与共享存储

运维常见面试题_05_NFS 与共享存储

Q102. NFS v3 和 NFS v4 的主要区别是什么? 1. 是什么:NFS(Network File System)是网络文件共享协议。v3(RFC 1813)是经典的 RPC 版 NFS;v4(RFC 7530,含 v4.1/4.2&#x…

2026/8/16 8:07:38 阅读更多 →

日新闻

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

如果你是一名开发者,最近可能已经感受到了AI大模型正在从“玩具”变成“生产力工具”的强烈信号。从代码补全到智能Agent,从本地部署到云端API,我们正处在一个技术栈快速重构的节点。然而,面对层出不穷的模型、框架和工具&#xf…

2026/8/16 0:00:54 阅读更多 →
工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

第四篇:反射——高频能量撞墙之后会发生什么? —— 你以为信号已经过去了,其实它正在回来打你 老Q的现场笔记 第五季,我们正式进入工业神经系统层。这里不再是单个设备的战斗,而是整个工厂“经脉”层面的秩序之战。从这一篇开始,你将第一次看清:看似简单的信号传播,背…

2026/8/16 0:00:55 阅读更多 →
【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

✅作者简介:热爱科研的Matlab仿真开发者,擅长毕业设计辅导、数学建模、数据处理、建模仿真、程序设计、完整代码获取、论文复现及科研仿真。🍎 往期回顾关注个人主页:Matlab科研工作室👇 关注我领取海量matlab电子书和…

2026/8/16 0:03:55 阅读更多 →

周新闻

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

如果你是一名开发者,最近可能已经感受到了AI大模型正在从“玩具”变成“生产力工具”的强烈信号。从代码补全到智能Agent,从本地部署到云端API,我们正处在一个技术栈快速重构的节点。然而,面对层出不穷的模型、框架和工具&#xf…

2026/8/16 0:00:54 阅读更多 →
工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

第四篇:反射——高频能量撞墙之后会发生什么? —— 你以为信号已经过去了,其实它正在回来打你 老Q的现场笔记 第五季,我们正式进入工业神经系统层。这里不再是单个设备的战斗,而是整个工厂“经脉”层面的秩序之战。从这一篇开始,你将第一次看清:看似简单的信号传播,背…

2026/8/16 0:00:55 阅读更多 →
【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

✅作者简介:热爱科研的Matlab仿真开发者,擅长毕业设计辅导、数学建模、数据处理、建模仿真、程序设计、完整代码获取、论文复现及科研仿真。🍎 往期回顾关注个人主页:Matlab科研工作室👇 关注我领取海量matlab电子书和…

2026/8/16 0:03:55 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/16 6:00:23 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/16 6:00:24 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/16 6:00:27 阅读更多 →