如何处理Akka Typed Persistence Actor启动异常并接收信号
首先,我们来拆解你遇到的核心问题:你的子Actor在启动阶段就因为未配置默认持久化journal抛出了ActorInitializationException,直接终止了;而你预期的Signal(比如onPersistFailure里的逻辑)根本没机会触发——因为这些Signal是针对Actor启动后运行时的持久化失败,不是初始化阶段的致命错误。同时,父Actor因为watch了子Actor却未处理Terminated信号,触发了DeathPactException。
下面是具体的解决方案:
1. 父Actor必须处理Terminated信号,避免DeathPactException
当子Actor启动失败直接终止时,父Actor的watch机制会收到Terminated信号。如果不处理这个信号,Akka Typed就会抛出DeathPactException终止父Actor。你需要在父Actor的行为中显式捕获这个信号,并执行你的通知逻辑。
修改你的父Actor代码如下:
object MessageSupervisorSpec { def create(communicator: Option[ActorRef[RpcCmd]], logger: Option[ActorRef[LogCmd]]): Behavior[MessageCmd] = Behaviors.setup { context => context.log.info("=============> Start MessageSupervisorSpec <=============") val fault = Behaviors .supervise(Persistence.create(communicator, logger)) .onFailure[ActorInitializationException](SupervisorStrategy.stop) val store = context.spawn(fault, "StoreChilid") context.watch(store) def loop(): Behavior[MessageCmd] = Behaviors.receiveMessage[MessageCmd] { case SaveMessage(v) => println(v) store ! SaveMessage(v) Behavior.same }.receiveSignal { // 新增:处理Terminated信号 case (ctx, Terminated(ref)) if ref == store => ctx.log.error("Store子Actor启动失败或意外终止") // 这里执行你的数据库不可用通知逻辑 communicator.foreach(_ ! SendMessage("数据库不可用:存储Actor启动失败")) logger.foreach(_ ! SaveLog(Log(Error, "消息存储Actor启动失败!"))) // 可选:如果需要重试启动子Actor,可以在这里重新spawn // val newStore = ctx.spawn(fault, "StoreChilid-Retry") // ctx.watch(newStore) Behavior.same } loop() } }
2. 理解为什么onPersistFailure没触发
你在EventSourcedBehavior中配置的onPersistFailure是用来处理Actor启动成功后,在持久化事件过程中发生的失败。而你的子Actor是在初始化阶段(还没进入事件循环)就因为journal配置缺失抛出了异常,所以这个回调根本不会执行。
如果想在子Actor内部处理初始化失败,可以用Behaviors.attempt包裹初始化逻辑,把异常转化为内部处理:
object Persistence { val storeName = "connector-store" // ... 原有commandHandler和eventHandler不变 ... def create(communicator: Option[ActorRef[RpcCmd]], logger: Option[ActorRef[LogCmd]]): Behavior[MessageCmd] = Behaviors.setup { context => context.log.info("=============> Start PersistenceMessageActor <=============") Behaviors.attempt { // 尝试初始化EventSourcedBehavior val eventSourced = EventSourcedBehavior[MessageCmd, MessageEvent, MessageState]( persistenceId = PersistenceId(storeName), emptyState = MessageState(), commandHandler = commandHandler(communicator), eventHandler = eventHandler) .onPersistFailure(SupervisorStrategy.restartWithBackoff( minBackoff = 10.seconds, maxBackoff = 60.seconds, randomFactor = 0.1)) .receiveSignal { case (_, _) => println("!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!") context.log.error("PersistenceMessageActor已终止,请检查持久化配置") logger.foreach(actor => actor ! SaveLog(Log(Error, "Message store actor was stopped!"))) } Right(eventSourced) }.recover { // 捕获初始化异常 case ex: ActorInitializationException => context.log.error("持久化Actor初始化失败", ex) logger.foreach(_ ! SaveLog(Log(Error, "消息存储Actor初始化失败!"))) // 终止自身,父Actor会收到Terminated信号 Behaviors.stopped } } }
不过即便这样,子Actor最终还是会终止,父Actor依然需要通过Terminated信号来完成全局的通知逻辑。
3. 补充:解决journal配置问题(测试场景)
如果你只是想模拟数据库离线,而不是真的因为配置缺失导致失败,可以在测试配置中使用内存journal来避免初始化错误:
akka.persistence.journal.plugin = "akka.persistence.journal.inmem" akka.persistence.snapshot-store.plugin = "akka.persistence.snapshot-store.local" akka.persistence.snapshot-store.local.dir = "target/snapshots"
这样你可以通过让内存journal模拟不可用(比如自定义测试插件)来触发后续的失败逻辑,而不是卡在初始化阶段。
内容的提问来源于stack exchange,提问作者softshipper

