为何fs2 Stream在Deferred完成时未中断?执行顺序差异解析
FS2 Stream中断行为差异的原因分析
问题场景
场景1:Stream未触发中断
import cats.effect.* import fs2.Stream import scala.concurrent.duration.* import cats.effect.unsafe.implicits.global val test = for { cancel <- Deferred[IO, Either[Throwable, Unit]] _ <- (IO.unit.delayBy(5.seconds).map { _ => println("Completing deferred"); cancel.complete(Right(())) }).start _ <- Stream.awakeEvery[IO](1.second).map(x => println(x)).interruptWhen(cancel).compile.drain } yield () test.unsafeRunSync()
场景2:交换顺序后Stream停止
import cats.effect.* import fs2.Stream import scala.concurrent.duration.* import cats.effect.unsafe.implicits.global val test = for { cancel <- Deferred[IO, Either[Throwable, Unit]] _ <- Stream.awakeEvery[IO](1.second).map(x => println(x)).interruptWhen(cancel).compile.drain.start _ <- (IO.unit.delayBy(5.seconds).map { _ => println("Completing deferred"); cancel.complete(Right(())) }) } yield () test.unsafeRunSync()
差异原因分析
1. interruptWhen的参数类型不匹配
FS2的Stream.interruptWhen方法要求传入的Deferred类型为Deferred[F, Unit]——当这个Deferred被调用complete(())完成时,才会触发Stream中断。而你的代码中使用了Deferred[IO, Either[Throwable, Unit]],不符合方法的预期类型:
- 场景1中,即便
cancel被complete(Right(()))完成,interruptWhen的内部逻辑无法正确识别这个完成事件,因此Stream不会被中断,会持续打印时间。
2. 场景2的“中断”是程序退出的副作用
场景2中Stream停止打印,并非interruptWhen生效,而是因为:
- Stream通过
.start运行在后台Fiber中。 - 主Fiber在延迟5秒完成
cancel后,整个test流程结束,unsafeRunSync()会终止cats-effect Runtime的所有后台Fiber,Stream因此被迫停止。 - 这看起来像是
interruptWhen触发了中断,但实际是程序退出导致的结果。
修正方案
将Deferred的类型改为Deferred[IO, Unit],并调用cancel.complete(())来触发中断:
import cats.effect.* import fs2.Stream import scala.concurrent.duration.* import cats.effect.unsafe.implicits.global val test = for { cancel <- Deferred[IO, Unit] _ <- (IO.unit.delayBy(5.seconds).map { _ => println("Completing deferred"); cancel.complete(()) }).start _ <- Stream.awakeEvery[IO](1.second).map(x => println(x)).interruptWhen(cancel).compile.drain } yield () test.unsafeRunSync()
此时场景1会在5秒后正确中断Stream,程序正常退出。
内容的提问来源于stack exchange,提问作者michkhol
相关产品推荐
相关产品推荐

