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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 14:22:15