如何基于Cats或FS2实现可取消的超时回调?
用Cats Effect或FS2实现可取消带回调的定时器
基于Cats Effect的实现
Cats Effect 借助Fiber机制管理后台任务,我们可以通过启动后台fiber执行延迟回调,并返回该fiber的取消操作来满足需求。
先导入必要的依赖与语法:
import cats.effect.{Async, IO, IOApp, Temporal} import cats.syntax.all._ import scala.concurrent.duration._
实现核心的runTimer函数:
def runTimer[F[_]: Async: Temporal](callback: F[Unit], delay: FiniteDuration): F[F[Unit]] = // 启动后台fiber,先延迟指定时长再执行回调 Async[F].start { Temporal[F].sleep(delay) >> callback }.map(fiber => fiber.cancel) // 返回取消该fiber的操作
对应你给出的示例调用逻辑:
// 模拟用户输入询问 def askUser: IO[Boolean] = IO.println("是否取消定时器?(y/n)") >> IO.readLine.map(_.trim.toLowerCase == "y") def main: IO[Unit] = for cancel <- runTimer(IO.println("定时器触发!"), 5.seconds) shouldCancel <- askUser _ <- cancel.whenA(shouldCancel) // 仅当用户确认时执行取消 yield ()
如果用户在5秒内选择取消,后台fiber会被中断,回调不会执行;若超时未取消,回调将正常触发。
基于FS2的实现
FS2作为流式处理库,可利用其delayBy操作和fiber管理机制实现带取消能力的定时器:
导入必要依赖:
import cats.effect.{IO, IOApp, Temporal} import cats.syntax.all._ import fs2.Stream import scala.concurrent.duration._
实现runTimer函数:
def runTimer[F[_]: Temporal](callback: F[Unit], delay: FiniteDuration): F[F[Unit]] = Stream.eval(callback) .delayBy(delay) // 延迟指定时长后执行回调 .compile.drain // 将流转为副作用执行的F[Unit] .start // 启动后台fiber .map(_.cancel) // 返回取消操作
调用逻辑与Cats Effect方案一致:
def askUser: IO[Boolean] = IO.println("是否取消定时器?(y/n)") >> IO.readLine.map(_.trim.toLowerCase == "y") def main: IO[Unit] = for cancel <- runTimer(IO.println("定时器触发!"), 5.seconds) shouldCancel <- askUser _ <- cancel.whenA(shouldCancel) yield ()
该方案将回调包装为延迟执行的流,启动后台fiber后返回取消操作,逻辑与Cats Effect方案本质一致,只是用FS2的流式API封装了延迟逻辑。
内容的提问来源于stack exchange,提问作者Max Smirnov
相关产品推荐
相关产品推荐

