Scala E02:FS2 的需求、背压与资源作用域
从一百万条订单里取前五条,不应先把一百万条读入集合。可是把集合改成 Iterator,仍然没有说明谁负责关闭文件,也没有说明把处理交给异步任务后,上游可以领先多少。流式处理至少要回答三个问题:何时请求下一份数据,允许保留多少未消费数据,以及消费提前结束时谁关闭资源。
FS2 将这些问题放进带效果的流描述。一个 Stream[IO, A] 可以分批产生 A,在生产过程中执行 IO,并把资源作用域与消费过程连接起来。它不是预先装满的容器;compile 把流整理为目标效果,随后效果运行时才进行本次消费。
本篇使用 Scala 3.3.7、FS2 3.11.0、Cats Effect 3.5.7 与 JDK 21。实验位于 examples/scala-lab/electives/E02,只处理本机临时文件,没有远端服务、消息队列或吞吐量基准。最终通过记录是 evidence/20261002-d/file-backpressure.log。
从需求端解释读取次数
输入文件包含 1 到 100,每行一个整数。生产端每次执行一次 readLine,把成功读到的行计入 Ref。消费端先等待 5 毫秒,再把该行转换为整数;整个流只要前五个元素。
1 | |
IO(s.toInt) 将解析放进效果执行阶段,*> 按顺序运行睡眠与解析。若写成 IO.sleep(5.millis).as(s.toInt),as 的严格参数会在构造该 IO 时先解析字符串,随后才睡眠;合法数字的最终值相同,却改变了解析与异常发生的时序。Cats Effect 3.5.7 IO 源码
这里没有 parEvalMap、后台生产 fiber 或显式预取。evalMap 必须先完成当前元素的效果,才能把该元素交给下游并继续串行消费。下游 take(5) 达到上限后终止,不再请求第六行。实际输出的列表为 1、2、3、4、5,成功读取计数也为 5。
五次读取与五次消费相等,是这段具体管道的观测结果。它不是“任何 FS2 管道永远只多读零个元素”的保证。若上游一次读取一个大块,或中间加入预取缓冲,下游只保留五个元素,也可能已经请求更多数据。必须把需求单位说明白:元素、块和字节的边界不一样。
实验使用 Java BufferedReader,因此 readLine 计数并不等于系统调用次数,也不等于从磁盘读入的字节数。reader 自身会缓冲字符,操作系统还可能预读磁盘页。本文的 buffer=0 表示没有增加 FS2 异步队列,不表示机器没有任何缓存。
这个区分决定了验证方式。要测应用级积压,就在生产与消费的语义边界计数;要测堆内存,需要记录对象和块大小;要测磁盘读取,则要借助系统层观察。不能只看到前五条业务结果,就断言物理层只读取五行对应的字节。
慢消费者怎样限制上游
背压意味着下游处理能力能限制上游继续推进。串行拉取管道通过“当前效果完成后才继续”建立这个约束;显式有界队列则通过“队列满时生产者等待”建立约束。两者都可能限制积压,但暂停位置、并发度与等待对象不同。
本实验把 IO.sleep 放在 evalMap 内部。睡眠属于本次消费的效果,流需要等待它完成。如果改成普通函数里启动一个后台任务,再立即返回 Unit,流只能观察到“任务已提交”,不会自动等待实际处理完成。此时输入读取得很快、任务排队很长,并不违反流描述自身的顺序。
因此检查背压时,要找到真正完成业务工作的那个效果。若数据库客户端返回 Future,桥接必须等待 Future 的结果;若客户端只返回“写入本地发送队列成功”,这个完成信号就不能证明服务端完成写入。流组合器能传播自己观察到的完成,无法推断库内部隐藏队列的状态。
有界也需要具体单位。限制为十个元素时,如果元素是整个文件的 Array[Byte],内存仍可能很大;限制为十个固定大小块,才能根据块大小估算队列占用。一个元素含有可继续增长的对象图时,数量上界还不能直接换算成字节上界。
[PATTERN] 判断背压要同时标出完成信号、等待位置和容量单位。只有“用了流库”这条事实,无法给出任何积压上界。
资源范围覆盖消费过程
读取器由 Resource 获得,关闭动作也绑定到同一个资源。Stream.resource 把这个资源范围提升到流的消费范围,因此提前 take、正常读到 EOF 或消费抛错时,都有明确的释放路径。
1 | |
closed 标记放在 close 之后。如果关闭抛异常,标记不会被写成 true,实验也不会把“尝试关闭”算作关闭成功。外层临时文件资源则负责删除文件;层次关系是先结束 reader 的流作用域,再退出临时文件的使用区。
获取资源不能写成先在外部执行 Files.newBufferedReader,再把已经打开的句柄塞入纯 Stream。那样文件在流真正运行前就打开了;如果描述构造成功但从未消费,关闭动作可能没有对应的生命周期。资源描述应当保存获得动作,让获得与消费发生在同一套效果协议内。
另一种错误是在资源使用区返回惰性迭代器。Resource.use(r => IO.pure(r.lines())) 的结果虽然包含后续读取方法,use 已经结束,reader 随即关闭。引用能逃逸并不意味着文件仍有效。正确结构是让完整的迭代消费留在 use 内,或者让返回对象本身携带受管理的资源协议。
本实验检查的是 take(5) 的提前结束路径。对于真实文件解析器,还应注入格式错误、磁盘读错和取消,观察资源释放与错误传播。一次提前结束成功不能代替所有故障分支;同样,没有抛异常也不是已经释放句柄的证据。
Iterator 与流的职责差异
Iterator 的 next() 同样能够按需计算下一个值。若生产和消费都在一个同步调用栈中,它已经足以避免提前计算全部元素。为一个十行文件引入效果运行时,不会自动提高正确性,代码结构可能反而更复杂。
差异出现在效果与生命周期需要组合时。标准 Iterator 没有一个统一的“取消当前消费并关闭全部内层资源”协议,也不携带异步等待完成的类型。它可以被放入 try/finally 中安全消费,但跨函数、跨异步任务传递之后,必须额外建立关闭责任。
FS2 把效果类型、资源作用域和流终止放入组合结构,从而能够表达读取、解析、异步写出及取消。但它仍然依赖每个边界正确封装。如果某个函数内部启动了不受管理的任务,或把完整输入先转成 List,流的外观不能修复已经丢失的约束。
尤其要注意 compile.toList。它会收集实际消费的所有元素,本实验因为前面有 take(5),列表只有五项;删除 take 后列表就随文件内容增长。流式读取并不保证最终结果集合有界,消费端选用的终结操作也参与空间复杂度。
订单聚合通常可以用 fold 保存固定大小的累计值,而不保存每一条订单。若需求是按用户分组,状态大小又取决于用户数量。是否能做到常量空间来自算法和业务基数,不能由返回类型里出现 Stream 推导。
缓冲与并行需要新的验收
给当前管道添加缓冲,可以让读取与慢处理重叠,但会改变提前终止时已经读取的数量。缓冲大小控制某一层队列容量,生产端持有的当前块、下游正在处理的元素和 reader 内部缓冲仍要另算。对整个系统建立内存预算时,应把这些部分相加。
并行处理则同时引入顺序问题。要求结果保持输入顺序时,一个很慢的前序任务可能使后面已完成的结果等待;允许乱序输出时,需要业务能够根据订单标识重新关联结果。提高并发度会影响数据库连接占用,不能只观察吞吐量就决定参数。
如果使用无界队列把上游与下游彻底分离,上游可能读完整个文件,而慢消费者只完成少量订单。程序仍使用 FS2 类型,但有界积压这一性质已经被配置破坏。真实验收需要计数最高在途量,而不是只比较最终列表是否相等。
本篇没有加入这些并发操作,因此不宣称已经验证并行吞吐、队列满时的取消行为或多消费者公平性。它提供的是一个可以从输出手算需求次数的串行基线,后续每加一层并发,就能对照基线解释新增的在途工作。
有界性需要写出计量单位
本例文件只有一百行,每行也是短整数,因此“只读五行”同时方便人工核验。若真实输入允许一行包含数百兆字节,按行读取即使没有队列,也不能保证单条元素很小。内存上界至少需要同时考虑队列容量、单元素大小、解码器状态与下游保留的结果。只写缓冲容量为十,却不限制元素尺寸,仍然缺少完整的空间约束。
此外,compile.toList 会保留所有输出。在本例中它位于 take(5) 之后,所以最多收集五个整数;如果移除 take,输出集合就随文件行数增长。流式读取只能限制中间过程,不会阻止终端收集器保存全部结果。需要处理大文件时,可以选择逐条写入、折叠固定大小状态,或把结果存入具有明确容量策略的外部系统。
下游失败同样应进入测试计划。例如第三行解析失败时,期望前两行已被处理,后续行不再读取,读句柄仍关闭。这个场景能够同时检查停止传播与资源范围。本文没有把正常截断实验扩大成所有错误路径的证据;新增失败场景时,应保留解析错误、读取计数和关闭标记三个独立观察值。
结果与可执行练习
运行退出码为 0,断言同时检查结果、读取次数和关闭标记。
| 观测项 | 实测值 | 边界 |
|---|---|---|
| 输入逻辑行数 | 100 | 临时文件中的整数行 |
| 下游结果 | 1、2、3、4、5 | 保留原顺序 |
| 成功 readLine 次数 | 5 | 不等于磁盘系统调用数 |
| reader 关闭标记 | true | 在 close 成功后设置 |
| FS2 异步缓冲 | 未增加 | reader 自身仍有缓冲 |
手算题:把 take(5) 移到 evalMap 前面,串行代码的结果与读取计数是否变化?在这段具体程序里仍为五项、五次读取。若把 take 完全删除,结果变为一百项,compile.toList 必须保存这一百项。
修改练习:在读取第五行后的解析阶段抛出异常,外层通过 attempt 观察失败,再断言关闭标记为 true。不要在关闭器中吞掉异常来“保证测试通过”;应当保留能区分读取失败与关闭失败的输出。如果加入缓冲,重新记录最高读取领先量,不能沿用本篇的五次结论。
| 可迁移规则 | 设计时的问题 |
|---|---|
| 需求以具体单位传播 | 一次请求产生一行还是一个块 |
| 异步完成需要进入效果 | 返回代表已提交还是已完成 |
| 资源范围覆盖全部消费 | 提前结束由谁关闭句柄 |
| 终结操作影响空间复杂度 | 是累计一个值还是收集全部值 |
完整运行命令见实验说明。
参考资料
- FS2 guide:流、效果与资源。
- FS2 3.11.0 官方指南源码:冻结发行的组合器、分块与资源说明。
- Cats Effect Resource:获得与释放范围。
- Cats Effect tutorial:阻塞文件操作的效果边界。
