Akka Actor通过ActorRefFactory创建子Actor时监管策略失效问题
问题分析与解决方案
你的Akka监管策略没生效的核心原因很明确:Akka的监管机制只对Actor消息处理线程中同步抛出的未捕获异常生效,而你的Mongo重复键异常是在RxJava的异步订阅线程(日志里的InnocuousThread-7)中触发的,完全脱离了Akka Actor的消息处理上下文,所以父Actor的OneForOneStrategy根本感知不到这个异常。
下面是具体的修复方案,分两种常见实现方式:
方案1:将异步操作转为Future,在Actor内处理结果
把Mongo的异步操作包装成Akka兼容的Future,这样异常会进入Actor的调度器上下文,触发监管策略。
步骤1:修改数据访问层,返回Future
调整EventInsert的方法,不再自行订阅Observable,而是返回Future给Actor处理:
import scala.compat.java8.FutureConverters._ import scala.concurrent.Future class EventInsert(mongoFactoryTrait: MongoFactoryTrait) { def eventInsert(event: Event): Future[Completed] = { val document = insertDocument(event) val collection: MongoCollection[Document] = mongoFactoryTrait.collectionMongo(Event_COLLECTION_NAME) val insertionResult: SingleObservable[Completed] = collection.insertOne(document) // 将RxJava Single转为Scala Future insertionResult.toFuture.asScala } } // 同步调整Service和Repository层的方法签名,传递Future class InsertEventService @Inject()(eventRepo: EventRepository) extends InsertEventServiceTrait { override def insertEventInMongo(event: Event): Future[Completed] = { eventRepository.insertEvent(event) } } class EventRepositoryImpl @Inject()(mongodb: MongoFactoryTrait) extends EventRepository { val eventInsert = new EventInsert(mongodb) override def insertEvent(event: Event): Future[Completed] = { eventInsert.eventInsert(event) } }
步骤2:在子Actor中处理Future并主动抛出异常
在EventWriteActor中处理Future的结果,当失败时抛出异常,让父Actor的监管策略捕获:
class EventWriteActor @Inject() (insertEventServiceTrait:InsertEventServiceTrait) extends Actor with ActorLogging { import context.dispatcher def receive: PartialFunction[Any, Unit] = { case InsertEventInMongo(event) => val originalSender = sender() val insertFuture = insertEventServiceTrait.insertEventInMongo(event) insertFuture.onComplete { case Success(_) => originalSender ! true case Failure(e: MongoException) => log.error("Mongo insertion failed", e) originalSender ! Status.Failure(e) // 主动抛出异常,触发监管策略 throw e case Failure(e) => log.error("Insertion failed with unexpected error", e) originalSender ! Status.Failure(e) throw e } } }
方案2:通过消息将异常传回Actor上下文
如果不想修改异步操作的返回类型,可以在Observer的onError方法中,向当前子Actor发送包含异常的消息,让Actor在消息处理流程中抛出异常。
步骤1:定义异常消息类型
case class InsertFailure(ex: Throwable, originalSender: ActorRef)
步骤2:修改EventInsert的订阅逻辑
在onError中向子Actor发送异常消息:
class EventInsert(mongoFactoryTrait: MongoFactoryTrait) { def eventInsert(event: Event, actorRef: ActorRef, originalSender: ActorRef):Unit = { val document = insertDocument(event) val collection: MongoCollection[Document] = mongoFactoryTrait.collectionMongo(Event_COLLECTION_NAME) val insertionResult: SingleObservable[Completed] = collection.insertOne(document) insertionResult.subscribe(new Observer[Completed] { override def onNext(result: Completed): Unit = { log.info("insertEvent: Event successfully inserted") originalSender ! true } override def onError(e: Throwable) = { log.error("insertEvent: Insertion failed", e) // 向当前子Actor发送异常消息 actorRef ! InsertFailure(e, originalSender) } override def onComplete(): Unit = {} }) } }
步骤3:在子Actor中处理异常消息并抛出
class EventWriteActor @Inject() (insertEventServiceTrait:InsertEventServiceTrait) extends Actor with ActorLogging { def receive: PartialFunction[Any, Unit] = { case InsertEventInMongo(event) => val originalSender = sender() insertEventServiceTrait.insertEventInMongo(event, self, originalSender) case InsertFailure(ex, originalSender) => originalSender ! Status.Failure(ex) // 抛出异常触发监管策略 throw ex } }
验证效果
修改完成后,当再次触发Mongo重复键异常时,你会看到ActorManager中监管策略的日志输出,并且子Actor会按照Resume策略继续运行,不会被重启或停止。
内容的提问来源于stack exchange,提问作者sarah w
相关产品推荐
相关产品推荐

