逐行读取文件,看起来只需要一个惰性迭代器。可是调用方读取五行就返回时,谁关闭文件?第三行转换失败时,是否还会读第四行?消费者正在等待一个外部请求时,上游已经读了多少数据?这些不是“按需求值”四个字能够完整回答的问题。它们同时涉及需求传播、效果失败、任务取消与资源范围。

本文用 FS2 把这些问题放进一条可执行管道。依赖固定为 Scala 3.3.7、Cats Effect 3.5.7、Cats 2.12.0、FS2 3.11.0。完整 Main.scala 在临时目录创建真实文件,不依赖仓库 fixture;result.json 保存本次十四条断言、运行命令、环境与源码摘要。初次依赖获取曾遇到 DNS 失败,联网重试后才完成当前验证,没有将下载失败记作测试通过。

集合的惰性不等于效果的生命周期

一个 Iterator[Int] 可以逐个产生整数,但这个类型没有表达文件如何获取、谁拥有句柄,以及遍历中断时如何归还。把迭代器从同步资源范围里返回,甚至可能让第一次 next() 就面对已经关闭的文件。给集合增加惰性,并不会自动把资源寿命延长到最后一个消费者。

Stream[IO, Int] 同时描述输出元素与执行这些步骤需要的效果。它仍然是计算描述,不是已经在后台读取的文件。只有把流编译为 IO 并执行时,文件读取、类型转换与 finalizer 才进入运行过程。这里的 compile 不是生成机器码,而是把流的解释与收集方式转成目标效果,例如 compile.toList 或 compile.drain。

toList 收集全部输出,因而可能消耗与输出总量成比例的内存;drain 丢弃输出,但依然执行流中的效果。把 toList 改为 drain 不会自动消除上游缓冲,也不能让一次本来会产生外部写入的 evalMap 变成无副作用。输出收集策略与管道执行策略要分开理解。

本文不是在比较所有语言里名为 Stream 的类型。Java Stream、Scala 惰性集合和 FS2 的效果流有不同接口与范围协议。名字相同不是可替换的依据。判断某种抽象能否处理文件导入,应具体看它如何表示获取、读取、失败、停止和关闭,而不是看它是否有 map。

让文件句柄属于流的范围

实验将文件获取放进 Resource,再通过 Stream.resource 引入流。文件中有整数一到二十,每行一个。使用 JDK BufferedReader 是为了同时观察真实句柄的关闭状态与应用层逐行读取次数;这里运行的是实际 FS2 流,但没有调用 fs2-io 的文件 API,也不宣称验证了那个模块的内部实现。

1
2
3
4
5
6
7
8
9
10
11
def lines(path: Path, p: Probe): Stream[IO, Int] =
Stream.resource(Resource.make(
IO.blocking(Files.newBufferedReader(path, UTF_8))
.flatTap(r => p.handle.set(Some(r)))
)(r => IO.blocking(r.close()) *> p.closed.update(_ + 1)))
.flatMap { r =>
Stream.repeatEval(IO.blocking(Option(r.readLine())))
.unNoneTerminate
.evalTap(_ => p.read.update(_ + 1))
.evalMap(s => IO(s.toInt))
}

readLine() 在 EOF 返回 null,因此用 Option 把它转换为结束信号;unNoneTerminate 在没有下一行时停止。读取计数放在结束信号之后,所以它数的是实际获得的文本行,不包含 EOF 探测。解析放进 IO,使格式错误进入效果的失败通道,而不是在构造流描述时提前抛出。

Probe 由三个 Ref 组成:已读行数、成功关闭次数、实际句柄。它不是业务对象,是测试观察面。关闭动作先执行 r.close(),成功后才增加次数,因此次数一表示这一段关闭动作确实返回,而不是仅进入了 finalizer。之后再对保存的实际句柄调用 readLine,要求得到 IOException,避免只用一个人为布尔值宣称资源已关闭。

保存句柄仅用于关闭后的否定检查,生产代码不应把它暴露给范围外调用方。这个测试也不需要在关闭后再次写入文件或依赖平台特定的删除行为。文件删除在外层临时文件资源的 finalizer 里执行;内层每个场景各自打开并关闭 reader,便于逐场景计量,不共享已消费的读取位置。

分批处理先定义批次的含义

正常路径将行流按四个元素组成一块,每块求和,最终输出 10、26、42、58、74。这覆盖真实文件打开、逐行读取、解析、分批计算、输出收集和关闭。求和是一个可复核的导入处理替身,不是已经完成数据库批量写入的宣称;实验没有连接数据库,也没有测试事务提交。

1
2
3
4
lines(path, p)
.chunkN(4)
.evalMap(chunk => IO.pure(chunk.toList.sum))
.compile.toList

四行一块意味着处理函数一次得到一组元素,不意味着四行共享一个数据库事务。事务范围需要由实际数据库接口表达。类似地,失败后从哪一行恢复、是否允许重复批次、批次号如何持久化,都不由 chunkN 自动决定。流解决组合和范围问题,业务提交协议仍然必须设计。

当前输入二十行刚好被四整除,正常路径没有验证尾部不足一块时的业务处理要求。若导入规定尾批也应提交,应增加二十一行的案例并验证最后一个和为二十一;若规定尾批必须拒绝,则应在业务校验处表达。本文不从整除输入的一次通过推断所有批次边界都已覆盖。

分批还影响失败粒度。若把解析放在分块前,解析到坏行时可能已经得到本批前面的合法行,但尚未执行整批处理。若先读文本块再在块内解析,则错误收集和内存保留方式会变化。选择哪一种需要看是否允许局部导入、是否要报告所有坏行,以及失败时已经发生哪些外部动作。

提前取五行,必须检查没有第六次逻辑读取

早停路径执行 lines(path, p).take(5).compile.toList。输出必须是 1 到 5,读取计数必须为五,关闭次数必须为一,实际句柄必须拒绝后续读取。只看五个输出不足以证明需求没有提前扩大;上游可能已经把二十行全读进内存,再返回前五个值。

这一场景没有显式预取或并行操作,源通过一次次 repeatEval 提供行。实验的否定断言针对这个具体管道:没有第六个非 EOF 行被交付到计数点。它不是 FS2 中任意 take(5) 管道的普遍读取上界,后面的分块反例会改变这个数字。

也不能把逻辑行数五写成“操作系统只读了五行对应的字节”。BufferedReader 可能预先填充字符缓冲,底层文件系统也有页缓存。探针位于应用层 readLine 之后,只能观察返回的逻辑行。讨论背压时必须说清计量单位:元素、字符、字节、块、请求数以及占用连接数,不是同一个量。

旧文 FS2 的需求与资源作用域 已经用真实 reader 说明过提前取五行与关闭。本文保留同样的最小观察面,但将定时等待替换为信号,并增加下游失败、取消、分块需求扩大与并发上限。这些新增路径不能由旧文的早停例子直接推出。

下游失败与跳过错误行的差别

失败场景让下游在收到第三个整数时返回 IO.raiseError。随后用 attempt 收集整个编译效果的错误,检查错误消息为 downstream,并断言读取计数为三。第四行没有经过读取计数点,reader 仍只关闭一次,实际关闭检查仍然失败于 IOException。

1
2
3
4
lines(path, p).evalMap { n =>
if n == 3 then IO.raiseError[Int](new Exception("downstream"))
else IO.pure(n)
}.compile.toList.attempt

这里错误产生在下游处理函数,不是文件读取本身。这样能检验消费者失败是否导致上游资源范围结束,避免只验证源自身能关文件。attempt 放在整个编译结果外部,所以一旦流失败,结果是失败值,而不是一个装着前两行的成功列表。

若改成在每行内部执行 job(n).attempt,错误就可能成为流中的普通元素,后面还会继续读取。这可以用于收集逐行错误,但会改变导入语义。此时“失败后没有第四行”的断言理应不再成立。错误处理器放置的位置不仅改变返回类型,也改变需求是否继续向上游传播。

此外,前两行已经执行过的外部效果不会因第三行失败自动回滚。当前处理只是返回整数,因而没有持久化副作用需要补偿。如果改为发送消息或写数据库,应明确每行、每批还是整个文件的事务范围。资源最终被关闭,和此前业务写入被撤销,没有逻辑等价关系。

消费者停住时,用信号测需求而不是猜时长

取消场景让消费者处理第一行时完成 entered,随后停在 IO.never。观察者等待 entered 后,读取计数必须为一。这表示消费者当前不再请求后继元素;然后取消承载 compile.drain 的 fiber,等待取消返回,检查 outcome 和实际句柄关闭。

1
2
3
4
5
6
7
8
9
fiber <- lines(path, p)
.evalMap(_ => entered.complete(()) *> IO.never[Unit])
.compile.drain.start
_ <- entered.get
n <- p.read.get
_ <- IO(assert(n == 1))
_ <- fiber.cancel
outcome <- fiber.join
_ <- IO(assert(outcome.isCanceled))

entered 放在消费者已经拿到元素之后,因而读取一次是已发生事实。消费者后面是不会自行完成的等待,所以不存在“观察稍晚就多读几行”的计时窗口。若用一个很短的 sleep 模拟慢消费者,机器负载变化可能让它在断言前结束,测试就把调度偶然性当成了需求协议。

取消路径还检查实际关闭,与早停和失败各自独立打开文件。不要只在正常路径检查 close,然后把同一个关闭次数当成所有终止路径的证明。不同路径的 bug 可能藏在编译效果的取消传播、异常转换或范围嵌套中。每次重建独立探针,也避免上个场景遗留的计数掩盖漏关闭。

这个取消实验没有阻塞在操作系统文件读取上,而是阻塞在下游的可取消等待上。因此它证明的是消费者取消能结束流资源范围,不证明正在进行的任意磁盘读取可以被立即打断。前一章对阻塞与协作取消的限制在这里仍然成立,换成流并不会删除底层 API 的边界。

分块后取五个元素,上游可以已经读了八行

背压不是“消费者需要一个值,上游恰好只计算一个值”的统一承诺。考虑在 take(5) 前增加 chunkN(4),再把块展开为元素:

1
2
3
4
5
lines(path, p)
.chunkN(4)
.flatMap(Stream.chunk)
.take(5)
.compile.toList

为了提供第五个元素,需要第二块;形成第二块时,上游已经读到第八行。因此实验同时断言输出仍是前五行,逻辑读取计数却是八。这个反例不表示需求控制失效,而是需求单位在中间变成了块。关闭检查继续通过,说明“多读了一部分已获准的块”和“早停后泄漏资源”是两类不同问题。

这个差异在真实导入中会影响成本。若读取之后立即附带某种外部效果,块内虽未最终输出的元素也可能已触发该效果。应尽量让容易提前发生的上游步骤只做读取或纯转换,把必须严格控制次数的业务动作放在符合需求的位置。即便如此,也应通过实际管道测试来判断,不依赖直觉推导任意组合的预取范围。

块大小不是越小越安全或越大越高效的单向参数。小块减少一次需求带来的元素数量,但可能增加调度、调用和批处理次数;大块可以摊薄部分开销,却扩大暂存与失败恢复单位。本文没有做性能基准,因此只报告四行块造成八次逻辑读取,不提供吞吐提升比例或所谓最佳块大小。

有界并发限制任务,不直接给出总内存上界

最后一个场景在实际文件流上使用 parEvalMap(2)。处理函数的资源范围增加当前活跃数并更新最大值,退出时减少;前两个任务停在门闩上。观察者等到两个任务进入后,检查当前与最大值都是二,然后放行全部工作。

1
2
val processed: Stream[IO, Int] = lines(path, p).parEvalMap(2)(job)
val result: IO[List[Int]] = processed.compile.toList

结束后,当前活跃数必须归零,历史最大值仍为二,输出为一到二十的原输入顺序,文件已关闭。这里的保序是输出顺序,不是工作完成顺序。实验没有人为制造每个任务的反向完成过程;第28章用双门闩检验的是业务事件逆序写入时元组结果仍保持位置,没有证明两个子 fiber 的终态逆序。那个实验的业务事件序列也不能冒充本章文件任务的实际序列。

本场景没有断言并行暂停点的总读取行数等于二。并发处理、内部交接、结果保序和上游分块都可能影响已读取但尚未交付的数据量;当前探针只精确限定 job 的活跃区间为二。若要建立严格的总缓冲上界,应在选定版本与完整算子链上增加对应观察点,并把源缓冲、处理中、完成但待排序的结果以及最终收集器分别计量。

保序还可能产生队首阻塞:较早的任务慢时,后面的已完成结果不能越过它直接输出。是否采用不保序处理,取决于业务是否能用编号恢复关联,不能仅凭“更快”选择。即使换成不保序操作,资源归属、错误传播和取消行为也要重新核对,不应只更新最终列表的断言。

本文最后收集二十个整数,规模明确且很小。生产导入若有数百万行,照抄 compile.toList 会保留所有输出;限制两个处理任务不解决这个问题。若只需要处理成功数量,可在流中做汇总并返回一个计数,或者在所有必要效果后 drain。需要保存所有失败记录时,也应设计其存储边界,而不是让一个不断增长的内存列表承担持久化职责。

把流类型当作接口契约的一部分

返回 Stream[IO, Row] 的函数应说明每次编译是否重新打开数据源,数据是否可以重复读取,以及消费者提前停止会触发什么清理。本文每次调用 lines 都在运行范围内打开新的 reader;重复编译同一描述并不意味着复用已关闭的句柄。若应用将资源提前分配再返回流,寿命规则会改变,类型外还需要相应契约。

消费端也应说明错误策略。整流失败、逐行累积错误和跳过错误行是三种不同操作。跳过策略若没有错误计数或记录,可能让导入返回成功但静默遗漏数据;累积策略若无限保留异常,又可能使一个可流式处理的文件变成内存问题。流组合让这些选择显式,但不会替业务确定哪一种正确。

效果流并不要求所有逻辑都藏在 evalMap 里。可将字符串解析、字段校验与记录转换保留为纯函数,只有读取、持久化和必要观察放在效果边界。这样既能独立测试转换,也能在流层集中检验失败如何终止需求、取消如何清理资源。纯核心和效果范围各自有可观察契约,通常比一个包揽全部动作的处理函数更易验证。

EOF、空输入与失败输入不能共用一个成功判据

当前文件含二十行合法整数,正常路径会读到 EOF;提前停止则不会为了寻找 EOF 而继续读完整个文件。两种路径都应关闭资源,但结束原因不同。若把关闭判断写成“只有读到 null 才 close”,正常文件可能通过,早停与下游失败却会泄漏。将 close 绑定到范围退出,而不是绑定到某个数据值,才能覆盖这些出口。

空文件还需要独立思考:它没有任何元素交给消费者,但只要 reader 已获得,仍有释放责任。如果取消测试将 ready 放在第一行的消费者中,却改用空文件,观察者就会永远等不到 ready。这不是空流不能关闭,而是测试准备条件不可能成立。换输入时必须同时检查控制协议的前提,不能只复用相同门闩。

格式错误文件与下游错误也不同。当前失败注入发生在整数解析之后,验证的是消费者拒绝第三个值。如果把第三行改成非数字,错误会发生在 lines 的解析阶段;此时计数器已增加,因为该文本行已成功读取。要验证这一路径,应断言解析错误、读到第三行且关闭,而不是继续要求错误消息为 downstream。两种故障虽然都终止流,负责产生错误的阶段不同。

EOF、空输入、解析失败和消费者失败分别对应没有下一项、从未有项、项无法构造以及项无法处理。将这些状态分清,可以避免把所有异常都映射成“文件为空”,也让故障测试有明确靶点。本文已运行消费者失败,空文件和坏文本属于扩展验收项,不列入当前十四条通过结论。

复跑与练习

1
node examples/functional-programming/run.mjs 30

十四条断言覆盖正常分批与真实关闭两条、早停两条、下游失败两条、取消需求与结果及关闭三条、分块早停两条、并发范围与结果及关闭三条。所有文件由运行时创建并清理,仓库没有额外输入夹具、缓存或临时数据。当前运行仅说明这些有限场景成立,不包括磁盘损坏、读取权限失败、reader 关闭失败或真实数据库事务。

手算题:从 readLine(): String 到 IO[Option[String]],再到 Stream[IO, String]、Stream[IO, Int]、IO[List[Int]],逐步写出每次转换消除了什么结构、保留了什么效果。为什么把 compile.toList 放在资源范围外不能修复一个已经关闭的共享 reader?再比较原始流和四行分块流,解释为什么同样输出五个元素会对应五与八两种读取次数。

修改题:把文件变为二十一行,并将批大小改成六。先手算正常批次和,再调整断言;随后在原始逐行流的第七个元素制造下游失败,要求计数七、没有第八次逻辑读取、真实句柄已关闭。最后在分块展开后做同样失败注入,重新观察读取范围,不应机械保留计数七的预期。

另一个修改题是给取消场景的关闭动作增加可控门闩,检查取消调用在关闭完成前不返回。这会把第29章的生命周期中间态与本章真实 reader 结合起来。新增实验需要单独复跑;当前十四条记录没有包含人为阻塞 reader finalizer 的场景。

参考资料

本文实际阅读 FS2 3.11.0 Guide 的流构造、编译、错误、资源和分块说明,以及 Cats Effect Resource 与 Sync。文档用于核对抽象与 API 边界,具体五行、八行、活跃上限二和关闭次数一来自本文所链接的实验,不能从这些数值推导任意管道的内存或性能结论。

系列导读 · 上一篇:29