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

Scala Cats-Effect 2中IO可靠取消方法及CE3等效实现问询

Cats-Effect 2 可靠触发IO取消及Cats-Effect 3等效实现

一、Cats-Effect 2 中可靠触发IO取消的正确方式

Cats-Effect的取消是协作式的——只有当IO执行可取消的步骤(如IO.sleep、IO.cancelable或基于资源的操作)时,取消信号才能被捕获。你的示例代码无输出,核心问题在于forever未内置取消响应逻辑,且runCancelable的回调仅触发于IO正常完成或错误终止,不会响应主动取消。

1. 让目标IO支持取消

先构建能响应取消的forever逻辑,比如用IO.foreverM包裹带可取消步骤的操作:

import scala.concurrent.duration._
import cats.effect.IO

val forever = IO.foreverM {
  IO.println("Running...") >> IO.sleep(1.second)
}

2. 注册取消回调并执行取消

通过IO.onCancel注册取消时的触发逻辑,再调用runCancelable获取取消句柄:

val cancelableForever = forever.onCancel(IO.println("Cancelled!"))

val start = cancelableForever.runCancelable {
  case Left(ee) => IO.println(s"FAILURE: ${ee.toString}")
  case Right(_) => IO.println("SUCCESS!") // 此分支永远不会触发,因forever永不结束
}

val cancel = start.unsafeRunSync()

Thread.sleep(3000)
cancel.unsafeRunSync()

println("finished")

3. 强制取消(进程终止)

若IO逻辑完全不支持协作式取消,只能通过终止进程强制停止:

  • 调用Runtime.getRuntime().halt(0)直接终止JVM
  • 将IO运行在独立线程,调用线程的interrupt()(依赖线程内部响应中断)

二、Cats-Effect 3.5.x 中的等效实现

Cats-Effect 3移除了runCancelable,改用Fiber、race和Resource等API处理取消逻辑。

1. 用Fiber手动延迟取消

通过start方法启动IO并获取Fiber实例,后续调用其cancel方法触发取消:

import cats.effect.{IO, IOApp}
import scala.concurrent.duration._

object CancelExample extends IOApp {
  val forever = IO.foreverM {
    IO.println("Running...") >> IO.sleep(1.second)
  }.onCancel(IO.println("Cancelled!"))

  override def run(args: List[String]): IO[ExitCode] = {
    for {
      fiber <- forever.start
      _ <- IO.sleep(3.seconds)
      _ <- fiber.cancel
      _ <- IO.println("finished")
    } yield ExitCode.Success
  }
}

2. 用race实现超时自动取消

通过IO.race让超时IO与目标IO竞速,超时后自动取消目标IO:

val timedCancel = IO.race(IO.sleep(3.seconds), forever).flatMap {
  case Left(_) => IO.println("Timeout, cancelled!")
  case Right(_) => IO.println("Success (unreachable)")
}

timedCancel.as(ExitCode.Success)

3. 用Resource处理取消时的资源释放

利用Resource的自动释放机制,取消时会自动执行资源的释放逻辑:

val resourceForever = Resource.make(IO.println("Starting"))(_ => IO.println("Cancelled/Released")) >>= { _ =>
  IO.foreverM(IO.println("Running...") >> IO.sleep(1.second))
}

resourceForever.use(_ => IO.sleep(3.seconds)).as(ExitCode.Success)

内容的提问来源于stack exchange,提问作者tribbloid

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 05:12:36