两个查询都返回 IO[Int],并不意味着它们已经启动,也不意味着把它们放进同一个 for 就会并行。异步程序至少有三个不同的问题:描述何时构造,任务何时开始,以及结果按什么顺序交给调用者。再加上资源限制,还要问同一时刻允许多少任务占用连接。混淆这些问题,常见后果不是类型错误,而是一个看似并行的服务实际串行,或者一个看似有界的批量操作提前启动了全部外部请求。

本文用可控信号隔离这些差异。示例版本固定为 Scala 3.3.7、Cats Effect 3.5.7、Cats 2.12.0。完整程序在 Main.scala,本次实际运行的命令、环境、源码摘要及九条断言输出在 result.json。实验不用随机等待猜测任务是否启动,也不测吞吐量。

类型中的依赖和执行中的重叠

设 loadUser: IO[User],loadOrders: User => IO[List[Order]]。第二步需要第一步产生的用户,因而组合必须等到这个值存在。loadUser.flatMap(loadOrders) 表达的是数据依赖,不是“把两个线程接起来”。第一步可以异步等待网络,也可以立即完成;这些实现差异不会消除第二步对用户值的依赖。

若输入已经包含用户编号,loadProfile(id) 和 loadPermissions(id) 都能在调用前独立构造。这时二者没有由类型表达的结果依赖,可以选择顺序运行,也可以选择并发运行。“可以并发”仍不是“应该并发”:它们可能争用同一个非线程安全句柄,可能必须遵守请求顺序,也可能受远端限流约束。数据独立只是并发的一个前提,外部协议的独立性需要另行判断。

Cats Effect 的 概念说明 区分并发结构和物理并行。两个 fiber 都处于等待状态,也是一种有意义的并发,不要求两个 CPU 核同时执行指令。本文的启动信号能证明两个计算已进入各自的逻辑区间,不能证明它们在同一纳秒占用了两个内核线程。不要把逻辑事件证据换算成硬件利用率。

IO[(A, B)] 本身也没有告诉读者里面是串行还是并行。选择发生在构造它的组合操作上。tupled 与 parTupled 可以得到同形状的结果,却有不同的运行行为。因此代码审查需要同时看类型和使用的实例、语法入口,不能只扫函数签名。

用一个未完成信号观察 flatMap

实验中的第一步先报告“已经开始”,然后停在一个尚未完成的 Deferred 上。第二步只修改一个 Ref[Boolean]。观察者等到第一步的开始信号后,读取第二步的标记。这里第一步不能自行越过门闩,所以读取标记时不存在“观察太早还是太晚”的计时猜测。

1
2
3
4
5
6
7
8
9
10
11
12
13
for
first <- Deferred[IO, Unit]
release <- Deferred[IO, Unit]
second <- Ref.of[IO, Boolean](false)
fiber <- ((first.complete(()) *> release.get)
.flatMap(_ => second.set(true))).start
_ <- first.get
before <- second.get
_ <- IO(assert(!before))
_ <- release.complete(()) *> fiber.joinWithNever
after <- second.get
_ <- IO(assert(after))
yield ()

第一条否定断言很重要:它检查第二个效果没有发生,而不是只检查最后两个步骤都发生过。只检查最终的 true,无法识别有人把第二步提前放到独立 fiber 中的回归。第二条断言则排除“第二步根本被漏掉”的错误。两个观察点合起来才描述了先后约束。

这里 start 只是让测试观察者能同时观察被测程序;被测程序内部仍然是 flatMap 顺序。若去掉外层 start,测试自身会先等待未放行的第一步,无法运行后面的放行动作。这是测试控制流的死锁,不是 flatMap 的缺陷。写并发测试时,应先明确哪个 fiber 负责被测任务,哪个负责控制信号,以及每个信号由谁完成。

Deferred 的值只完成一次,适合表达一次性阶段边界;Ref 适合保存会变化的状态。不能用反复读取普通变量替代它们,再期待测试可靠地跨线程可见。实验里的 first 表达事实已经成立,release 表达许可已经发出,它们虽然类型相同,协议角色不同。

并发启动不改变结果的位置

并行场景有两个开始信号和两个放行信号。观察者等待两者都开始,再先放行第二个任务。第二个任务记录事件 2,完成另外一个 bd 信号;观察者等到 bd 后,才断言事件列表只有 2。随后放行第一个任务,最终事件顺序为 2, 1,结果却是 (10, 20)。

核心组合如下,完整的信号声明和断言在实验源码中:

1
2
3
4
5
val both: IO[(Int, Int)] = (
a.complete(()) *> ga.get *> events.update(_ :+ 1).as(10),
b.complete(()) *> gb.get *> events.update(_ :+ 2) *>
bd.complete(()).as(20)
).parTupled

元组左位置属于第一个计算,右位置属于第二个计算;它不是按到达时间填充的两个槽。业务事件的写入顺序由门闩控制,结果位置由组合结构决定。将结果位置保持稳定,可以让调用方继续按字段语义解释返回值,而不必猜哪个服务先响应。

这里还有一个容易遗漏的先后关系:bd.complete 放在写入事件之后。若把它放到前面,观察者读事件时,第二个任务可能尚未完成写入,测试会凭空引入竞争。信号的名字即使叫 done,也不能让它自动代表后面尚未执行的动作。测试里的信号必须紧跟所要证明的事实,事件记录与信号之间的程序顺序也是证据的一部分。

bd 完成时,第二个 fiber 尚不一定进入终态;观察者得到信号后就可以放行第一个任务。因此 parallel-both-start-reverse-business-events 检查的是业务事件先 2 后 1,不证明两个 fiber 的 terminal Outcome 也按这个次序出现。最后等待的是整个 parTupled 的结果,没有分别观测两个子 fiber 的终态顺序。

结果保序不保证所有副作用保序。日志、服务调用开始、远端提交以及完成回调都可能交错。若业务要求按编号依次写入,不能因为 parTraverse 最后返回一个有序列表,就声称写入也是按序发生的。可将独立计算并发执行,再将有顺序要求的提交放在明确的顺序阶段,但这个拆分是否安全,还取决于提交前结果是否允许过期。

在途上限要包住真正占用资源的区间

六个任务配合上限二,比只跑两个任务更能检查限制。实验为每个任务记录三个数:当前活跃数、历史最大活跃数、总启动数。进入任务时原子地增加计数,退出时减少当前活跃数;门闩关闭期间,前两个任务不能结束,因此其余四个不应进入这个区间。

1
2
3
4
5
6
7
8
9
10
11
12
import cats.effect.syntax.all.*

val task = (i: Int) => Resource.make(
counts.modify { case (active, maximum, starts) =>
val next = (active + 1, maximum.max(active + 1), starts + 1)
(next, next._3)
}.flatTap(n => if n == 2 then two.complete(()).void else IO.unit)
)(_ => counts.update { case (active, maximum, starts) =>
(active - 1, maximum, starts)
}).use(_ => gate.get.as(i * 10))

val program = (1 to 6).toList.parTraverseN(2)(task)

上面的 Resource 是测试探针的作用域,不是实际数据库连接。它让当前计数的减少跟随任务退出,即使以后给任务加入失败或取消,也不必在多个分支重复减数。需要注意 cats.effect.syntax.all.*:parTraverseN 的效果语法不能仅靠 Cats 通用语法导入得到。本章首次编译确实漏了这一导入,修正后才获得当前通过的证据;依赖更新提示不是该编译错误的原因。

门闩未放行时快照必须是 (2, 2, 2)。这不只是检查最大值二,还检查没有第三个任务短暂进入又退出。全部完成后必须是 (0, 2, 6),结果为 10 到 60 的输入顺序。当前数归零说明所计量的区间全部退出,总启动数六排除了漏执行或重复执行,最大数二才说明峰值符合上限。只保留其中一个数字,会丢失另外两类错误。

计数探针必须放在被限制的实际操作内部。假设先构造六个已启动的 Future,再对六个“等待 Future 的 IO”执行 parTraverseN(2),受限的只是等待者,外部六个请求可能早已发出。反过来,若在许可区间里才创建请求,并让区间覆盖响应和必要清理,限制才与业务所说的在途请求接近。检查限流实现时,应先找真正的提交点,而不是先找名字里有 N 的方法。

有界并发也不等于限速。最多两个同时执行的任务,如果每个很快完成,一秒内仍能产生大量调用。并发上限约束占用数量,速率限制约束一段时间内的启动数量,队列容量约束等待数量。它们可以共同出现,但不能相互代替。本实验只验证一个进程内六个 IO 的活跃区间,不证明集群级配额或第三方接口的每秒限额。

Future 的提交与 IO 的描述不是同一个阶段

旧文 Future 与 ExecutionContext 已经解释了 Future 放在 for 内外时的提交差异。本章把它和 CompletableFuture、冷的 IO 放到同一计数器旁,关注适配边界,不重复线程池配置教程。

实验使用直接执行器,让提交的 Runnable 就在调用线程运行。这样做不是推荐生产中用直接执行器处理阻塞任务,而是消除“已提交但暂未调度”这一干扰因素。构造两个 future 后计数立即为二,构造 IO 后仍为二。等待同一个 Scala Future 两次不会再次增加计数,执行同一 IO 描述两次才从二变成四。

1
2
3
4
5
val direct: Executor = (r: Runnable) => r.run()
given ExecutionContext = ExecutionContext.fromExecutor(direct)
val future = Future(count.incrementAndGet())
val cf = CompletableFuture.supplyAsync(() => count.incrementAndGet(), direct)
val cold = IO(count.incrementAndGet())

supplyAsync 的 “async” 不能代替对执行器的理解。这个调用把任务交给指定执行器,执行器允许在当前线程直接执行它。由此得到的结论是构造会提交,而不是“每个 CompletableFuture 都新建一个线程”。同样,Future 的 eager 启动语义也不意味着构造时一定已拿到结果,普通线程池可能稍后才执行;本文特意选择直接执行器,才有确定的构造后数值断言。

IO.fromFuture(IO.pure(future)) 将等待结果接入 IO,并没有回到过去撤销已经发生的提交。若希望每次运行才创建新的 Future,需要把创建表达式本身放进延迟的外层,例如 IO.fromFuture(IO(Future(...)))。这又改变了共享策略:重复运行可能提交多个请求。适配代码同时决定启动时机和重复运行语义,应在接口契约中写清楚,而不是只说“返回类型已经改成 IO”。

取消也不能从返回类型自动推导。Scala Future 没有因此获得取消正在执行代码的通用协议。将已有 future 包进 IO,只能说明某种等待被组合起来;外部操作是否支持停止、如何登记回调、取消后是否还会写入,仍是适配器要解决的边界。下一章用原生 fiber 检验生命周期,不能把那些保证未经验证地移植到任意异步 SDK。

启动、结果和错误需要分别定义

批量查询常把“保序”“全部完成”“尽快失败”写在同一个需求里。保序只规定成功结果的排列;尽快失败会影响其他正在执行的任务;全部完成则可能要求收集每个任务的失败值。若把每个 IO[A] 改为 IO[Either[Throwable, A]],失败被转为普通结果,批量组合的中断行为就可能变化。这是业务语义的改变,不只是为了少写一个异常处理器。

本文没有把失败路径混入上限实验,因为那会同时引入“谁触发取消”和“谁等待清理”两个问题。当前通过的九条断言证明了顺序、启动重叠、位置保持、活跃上限和 eager 对照;竞争失败的清理由下一章独立构造信号检查。明确各实验的证明范围,比把多个现象塞进一次成功执行更容易定位回归。

如果任务数量很大,先构造完整输入列表仍然占用内存。parTraverseN 的上限不会自动让已有列表变成外部流,也不会替调用方限制输入读取速度。需要边读边处理时,应将数据获取也放入效果流,并在流上选择有界并发。这是第30章要补的边界:并发任务数与上游已读取、已缓冲的数据量不是同一个量。

还有一种隐蔽的容量问题是返回值积累。即使最多两个任务同时运行,IO[List[LargeResult]] 仍会保存六个结果直到整个列表交付。限制计算中的在途数量,不意味着限制最终输出的总大小。若调用方只需要总和或成功计数,可以在合适的流式结构中逐步汇总;若确实需要全部结果,就应把其大小纳入容量预算。

许可范围也可能制造依赖死锁

上限越小不一定越安全。假设两个父任务各占一个许可,每个父任务都启动一个子任务,并等待子任务完成;子任务又必须从同一个上限二的许可池中申请许可。此时两个许可都由父任务持有,父任务只有等子任务完成才释放,而子任务没有许可就无法开始。这是等待关系形成闭环,增加超时时间不能解除它。

本章没有嵌套申请,所以不会触发这个例子;这是一项设计反例,而非已运行的死锁测试。它说明并发参数不能脱离作用域审查。若父任务只是组织工作,不直接占用受限资源,可以让实际操作才获取许可;若父子分别占用不同资源,则应分开限制,并规定获取顺序。简单地把上限从二调到三,也可能只把问题推迟到更多嵌套时发生。

另一种类似问题是任务之间存在隐式信号依赖。对输入列表使用有界并发时,前两个任务若等待第五个任务发送信号,但前两个不结束第五个就无法启动,也会永久等待。数据类型可能仍是互相独立的 IO[Int],隐藏在闭包中的 Deferred 却已经建立运行依赖。决定是否使用 parTraverseN 前,应检查函数是否真的能在任意允许的子集里推进。

实验的 gate 由外部观察者完成,观察者不属于六个受限任务,因此没有这个闭环。若为了“把测试写得更紧凑”而将放行动作改成第六个任务的一部分,测试会挂住。信号由谁负责、是否能在当前上限下运行,是并发测试协议的必要组成部分,不是代码排版上的细节。

有界组合限制的是单次组合范围。如果服务同时处理十个请求,每个请求都允许两个任务,应用层可能存在二十个受限区间。是否还需要共享连接池或应用级许可,取决于资源在哪里共享。把一个局部的二误读成整个进程的二,会让单请求测试全通过而总负载仍超过预期。当前实验只有一次组合,没有证明跨请求限制。

用三个计数器定位不同回归

设放行前得到的快照不是 (2,2,2),应先根据三个数定位问题。若为 (2,3,3),说明当前只剩两个,但曾有第三个进入并退出,不能用当前值二掩盖峰值违规。若为 (1,1,1),可能第二个任务尚未进入,也可能组合被错误改成顺序执行;在本实验里观察者等待第二次开始信号,因此后一种改动会卡在信号处而不是产生这个快照。

若完成后当前计数不是零,可能退出逻辑没有运行,或者计数增加了两次而只减少一次。若总启动数不是六,可能输入被过滤、重复遍历,或者探针没有包住真正的入口。三个字段不能互相替代,它们分别对应同时占用、历史峰值与累计执行。把一次综合断言拆成这三项原因说明,会比只打印“并发异常”更利于维护。

还要避免计数器自身造成观察漏洞。Ref.modify 在一个原子更新中同时计算新活跃值、最大值和启动值,因此不会先增加 active,稍后才更新 maximum,留下两个字段暂时矛盾的窗口。若改成三次独立更新,读者可能在两次更新之间取到不一致快照,即使任务限制本身没有出错也会产生假失败。

退出时只减少 active,不减少总启动数或最大值,这是三个指标的定义决定的。将所有字段在任务结束时归零,会使最后无法检查历史峰值;仅用一个当前计数器,也无法发现中途短暂超限。这个模式适合边界测试,但不能直接宣称等同于生产监控:生产指标的采样、导出失败和多实例汇总还需要独立设计。

本章选择六个立即可计算的整数,以确保测试关注调度协议,而不是计算成本。如果把每个任务改成真实远程调用,测试还要控制响应、错误和清理,不能继续依赖远端恰好按某种顺序返回。适配器的单元测试可以使用可控端口,真实集成测试再验证协议匹配;两者证据应分别标注,避免把可控信号当成远程系统已经实现相同语义。

并发测试也需要处理“不再推进”的失败形态。如果把本章 parTupled 错改成顺序组合,观察者会等不到第二个启动信号;这时不会自动出现一条失败断言,而是任务无法继续。测试执行环境应有独立的运行截止机制,将挂起报告为失败,不能将无输出解释为通过。这个外部截止仅负责发现测试未完成,不用于证明任务之间的正常顺序。

如何复跑与修改实验

在仓库根目录执行:

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

runner 将代码复制到临时目录编译运行,证据记录实际命令与源码摘要。当前九条 PASS 分属四组:顺序组两条,并发业务事件次序与结果位置两条,有界组两条,Future/CompletableFuture 与 IO 的启动共享组三条。退出码为零且这些断言存在,才对应本文所说的通过;单独看到编译器的版本提示不构成失败,也不构成运行证明。

手算题:first: IO[Int] 与 next: Int => IO[String] 应组合成什么类型?若改成 other: IO[String],为什么可以选择 parTupled,却不能只凭类型断言业务允许并发?再为实验画出 b开始 → b放行 → 写入2 → bd完成 → 观察Vector(2) 的顺序,指出删掉 bd 等待会使哪条断言变得不可靠。

修改题:把六个任务改成七个,上限改成三,并把到达信号的阈值也改成三。放行前应断言 (3,3,3),结束后应断言 (0,3,7);仅修改 parTraverseN 参数而不修改控制协议,会让测试含义不一致。随后把返回值改成输入编号的平方,保持输入顺序断言,避免把“值刚好递增”误当成保序的定义。

进一步的反例题是把 Future 的构造移到有界组合外面,并给提交点单独计数。不要预设它会遵守上限二;先描述哪些动作已经发生,再决定 parTraverseN 实际限制了什么。这个练习的答案应包含提交计数和等待计数两个观察面,而不是只看最终列表是否正确。

参考资料

本文实际阅读了 Cats Effect Concepts、Spawn 和 IO API。前两者用于核对 fiber 与并发组合的概念,API 用于核对 IO 组合边界。Future 对照参阅 Scala Futures and Promises,Java 执行器及取消边界参阅 JDK 21 CompletableFuture。网页可能随文档版本更新,本文的具体执行结论以冻结依赖的实验和对应证据为限,不把文档中的一般描述当作本次多核性能测试。

系列导读 · 上一篇:27 · 下一篇:29