You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 18:30:49