在Akka Actor遇错误时取消fs2 Streams的实现方法
如何在Akka Actor中取消fs2定时Stream?
你当前的代码存在两个核心问题:一是fs2 Stream的实现逻辑不符合定时发送自消息的需求,二是没有保留Stream的运行句柄,导致无法主动取消。下面是具体的修复和取消方案:
1. 先修正fs2 Stream的实现逻辑
你当前的getATaskAsStream写法有误,source.concurrently(sink)会导致callback仅执行一次就终止Stream,无法实现定时触发的效果。正确的写法应该让每个定时信号都触发一次callback:
def getATaskAsStream(pingInterval: FiniteDuration)(callback: Int => Unit): Stream[IO, Unit] = { // 每隔指定间隔生成信号,每个信号触发一次callback Stream.awakeEvery[IO](pingInterval) .map(_ => 0) .evalMap(x => IO(callback(x))) .drain }
2. 在Akka Actor中保存Stream的取消句柄
要取消fs2 Stream,必须持有它的运行Fiber(即Stream的执行句柄)。修改Actor代码,添加成员变量保存句柄,并在Actor启动时运行Stream:
import cats.effect.{IO, Fiber} import scala.concurrent.duration.FiniteDuration import akka.actor.{Actor, ActorLogging, ActorRef, PoisonPill} import play.api.libs.json.{JsError, JsValue, Json} class YourActor(bindings: Bindings, sink: ActorRef[JsValue]) extends Actor with ActorLogging { // 保存Stream的执行句柄,用于后续取消 private[this] var streamFiber: Option[Fiber[IO, Throwable, Unit]] = None override def preStart(): Unit = { super.preStart() if (bindings.appConfig.isPersistentWSConn) { val stream = StreamUtils.getATaskAsStream(bindings.appConfig.pingInterval)(x => self ! x) // 非阻塞启动Stream并获取执行句柄 val fiber = stream.compile.drain.start.unsafeRunSync() streamFiber = Some(fiber) log.info("WebSocket ping stream started") } else { log.info("Not creating a Persistent WebSocket connection") } } override def receive: Receive = { case jsValue: JsValue => jsValue.validate[ValidateSomething].asEither match { case Right(ocppCall) => // 原业务逻辑... case Failure(fail) => sink ! JsError(s"${fail.getMessage}") case Success(succ) => sink ! Json.toJson(succ) case Left(errors) => sink ! Json.toJson(s"error -> ${errors.head._2}") } case x: Int => log.info(s"Elem: $x") Future.successful(heartbeatResponse(2,"HeartbeatRequest")).onComplete { case Failure(fail) => sink ! JsError(s"${fail.getMessage}") case Success(succ) => sink ! Json.toJson(succ) } case msg: Any => log.warn(s"Received unknown message ${msg.getClass.getTypeName} that cannot be handled, eagerly closing websocket connection") // 主动取消Stream streamFiber.foreach { fiber => fiber.cancel.unsafeRunAsync { case Left(err) => log.error("Failed to cancel ping stream", err) case Right(_) => log.info("Ping stream cancelled successfully") } } // 关闭Actor self ! PoisonPill } // 可选:Actor停止时自动取消Stream,避免资源泄漏 override def postStop(): Unit = { streamFiber.foreach(_.cancel.unsafeRunSync()) super.postStop() } }
3. 关键细节说明
- Fiber的作用:fs2 Stream启动后返回的
Fiber是其执行的控制句柄,调用fiber.cancel即可主动终止Stream的所有后续执行。 - 线程安全:Akka Actor是单线程模型,因此在
preStart、receive、postStop中操作streamFiber无需额外同步,天然线程安全。 - unsafeRun系列方法:在Actor内部使用
unsafeRunSync或unsafeRunAsync是安全的,这些操作都在Actor的专属线程中执行,不会阻塞其他Actor或系统线程。 - postStop钩子:添加该钩子可确保Actor因任何原因停止时,都能自动取消Stream,避免资源泄漏。
内容的提问来源于stack exchange,提问作者joesan
相关产品推荐
相关产品推荐

