Akka Actor通过ActorRefFactory创建子Actor时监督策略失效问题
你的监督策略没生效的核心原因是:MongoDB异步操作的异常发生在独立的线程池(日志里的InnocuousThread),并没有在EventWriteActor的消息处理线程中抛出。Akka的SupervisorStrategy只能捕获Actor自身消息处理流程中(也就是receive方法执行期间)的未捕获异常,异步回调里的异常脱离了Actor的执行上下文,父ActorActorManager的监督机制根本感知不到。
看你的代码,EventInsert里的insertionResult.subscribe是异步执行的,当Mongo抛出重复键异常时,异常被回调的onError方法捕获了,但这个回调是在Mongo自己的线程里执行的,和EventWriteActor的消息处理线程完全无关,所以父Actor的监督策略根本不会触发。
要让监督策略生效,你需要把异常重新引入到EventWriteActor的消息处理上下文里,这里有两种最直接的实现方式:
方式一:利用Akka的Future包装异步操作
Akka对异步Future有很好的支持,你可以把Mongo的异步操作转换成Scala Future,让异常回到Actor的执行上下文:
修改EventInsert的代码
import scala.compat.java8.FutureConverters._ import scala.concurrent.Future def eventInsert(event: Event): Future[Completed] = { val document = insertDocument(event) val collection: MongoCollection[Document] = mongoFactoryTrait.collectionMongo(Event_COLLECTION_NAME) collection.insertOne(document).toScala // 转换为Scala Future }
修改EventWriteActor的处理逻辑
import akka.pattern.pipe import scala.concurrent.ExecutionContext.Implicits.global class EventWriteActor @Inject() (insertEventServiceTrait:InsertEventServiceTrait) extends Actor { def receive: PartialFunction[Any, Unit] = { case InsertEventInMongo(event) => val originalSender = sender() // 将Future结果回传给发送方,异常会被包装为Status.Failure insertEventServiceTrait.insertEventInMongo(event) .map(_ => true) .recover { case mongoEx: MongoException => // 抛出异常,触发父Actor的监督策略 throw mongoEx } .pipeTo(originalSender)(self) } }
同时调整服务层方法签名,返回Future而非Unit,确保异常能传递到Actor上下文。
方式二:在异步回调中向当前Actor发送异常消息
如果不想改动异步操作的结构,可以在onError回调里把异常发送给EventWriteActor自身,让Actor的receive方法处理异常,从而触发监督:
修改EventInsert的代码(需传入当前Actor引用)
def eventInsert(event: Event, currentActor: ActorRef, senderRef: Option[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") senderRef.foreach(_ ! true) } override def onError(e: Throwable) = { log.error("insertEvent: Mongoexception in inserting Event", e) // 将异常发送给当前Actor,让消息处理逻辑触发监督 currentActor ! akka.actor.Status.Failure(e) senderRef.foreach(_ ! akka.actor.Status.Failure(e)) } override def onComplete(): Unit = { senderRef.foreach(_ ! true) } }) }
修改EventWriteActor的receive方法
class EventWriteActor @Inject() (insertEventServiceTrait:InsertEventServiceTrait) extends Actor { def receive: PartialFunction[Any, Unit] = { case InsertEventInMongo(event) => val originalSender = sender() insertEventServiceTrait.insertEventInMongo(event, self, Some(originalSender)) case akka.actor.Status.Failure(e) => // 抛出异常,触发父Actor的监督策略 throw e } }
这样当异步回调发生异常时,会给EventWriteActor发送Status.Failure消息,Actor处理该消息时抛出异常,父Actor的SupervisorStrategy就能捕获并执行你配置的Resume策略。
- 你配置的
Resume策略意味着Actor会继续处理后续消息,不会重启;如果需要重启Actor,可以改为Restart。 - 从代码看,子Actor确实是由
ActorManager的context创建的,父子关系正确,这部分无需调整。
内容的提问来源于stack exchange,提问作者sarah w

