如何为可变更行为的AKKA Actor添加默认Receive行为?
解决方案
针对你遇到的问题,这里提供几个简洁的实现思路,避免在每个PeerState子类中重复编写状态查询逻辑:
方案一:在PeerState父类中添加默认状态查询行为
利用Akka Receive是PartialFunction的特性,在父类中定义通用的状态查询处理逻辑,让子类的行为自动叠加这个逻辑,无需每个子类重复实现。
步骤1:定义状态查询相关消息
// 用于查询Peer当前状态的消息 case object QueryState // 携带当前状态的响应消息 case class CurrentState(state: PeerState)
步骤2:修改PeerState父类
abstract case class PeerState(var order: Int, var complete: Boolean, var waitForPing: Boolean) { // 基础行为:统一处理状态查询请求 private val baseBehavior: Receive = { case QueryState => sender() ! CurrentState(this) } // 子类只需实现自身业务行为,自动叠加基础行为 def behavior: Receive = PartialFunction.empty orElse baseBehavior }
步骤3:子类实现自身行为时叠加基础逻辑
以WaitAccept为例,子类的行为只需关注自身业务,通过super.behavior orElse叠加父类的基础处理:
class WaitAccept(order: Int, peer: Peer)(implicit system: ActorSystem) extends PeerState(order, false, true) { override def behavior: Receive = super.behavior orElse { case Tick => system.log.info("Waiting for accept") if (complete) peer.getContext.unbecome() case rA@RegistrationAccept(serverID, serverIP, _) => system.log.info(s"Processing RegistrationAccept: ${rA.toProtoString}") complete = true } }
步骤4:修改PeerManager的状态查询逻辑
用ask模式向Peer请求状态,替代原来的queryPeerState:
case peerRequest@PeerRequest(serverID, SessionRequest(sessionRequest), _) => peerMap.get(serverID) match { case Some(peer) => try { // 先查询Peer当前状态 val currentState = Await.result((peer ? QueryState).mapTo[CurrentState], timeout.duration) if (currentState.state.isInstanceOf[SessionHandshake]) { // 状态符合,转发请求并等待响应 val response = Await.result((peer ? (self, sessionRequest)).mapTo[PeerResponse.Response], timeout.duration) sender() ! response } else { sender() ! PeerResponse.Response.StatusResponse(HydraStatusCodes.PEER_NOT_READY.getStatusResponse) } } catch { case _: TimeoutException => sender() ! PeerResponse.Response.StatusResponse(HydraStatusCodes.PEER_MANAGER_TIMEOUT.getStatusResponse) } case None => sender() ! PeerResponse.Response.StatusResponse(HydraStatusCodes.PEER_NOT_READY.getStatusResponse) }
方案二:让Peer Actor统一维护当前状态并处理查询
在Peer内部维护当前活跃状态的引用,同时在所有行为中保留状态查询的处理逻辑,无需修改PeerState子类。
步骤1:在Peer中添加状态跟踪变量
class Peer(...) extends Actor { private var currentActiveState: Option[PeerState] = None // 通用状态查询处理 private val stateQueryBehavior: Receive = { case QueryState => sender() ! CurrentState(currentActiveState.orNull) } override def receive: Receive = stateQueryBehavior orElse { case FirstTick => system.log.info("receive FirstTick") timers.startTimerWithFixedDelay(TickKey, Tick, FiniteDuration(2, TimeUnit.SECONDS)) case Tick => system.log.info("receive Tick") var statesLeft = false states = states.sortWith{ (s1, s2) => s1.order > s2.order } for (state <- states) { if (!state.complete) { statesLeft = true currentActiveState = Some(state) // 更新当前活跃状态 val targetBehavior = if (state.waitForPing) pingService(state) else state.behavior // 切换行为时叠加状态查询逻辑 context.become(targetBehavior orElse stateQueryBehavior) } } if (!statesLeft) { currentActiveState = None context.become(pingService(null) orElse stateQueryBehavior) } } }
步骤2:PeerManager的状态查询逻辑同方案一
通过ask获取CurrentState后判断状态即可。
方案三:取消提前状态检查,让Peer自行处理不匹配请求
简化PeerManager逻辑,直接转发请求,由Peer在状态不匹配时返回错误,无需单独查询状态。
步骤1:修改PeerManager逻辑
case peerRequest@PeerRequest(serverID, SessionRequest(sessionRequest), _) => peerMap.get(serverID) match { case Some(peer) => val responseFuture = (peer ? (self, sessionRequest)).mapTo[PeerResponse.Response] responseFuture.onComplete { case Success(resp) => sender() ! resp case Failure(_: TimeoutException) => sender() ! PeerResponse.Response.StatusResponse(HydraStatusCodes.PEER_MANAGER_TIMEOUT.getStatusResponse) case Failure(_) => sender() ! PeerResponse.Response.StatusResponse(HydraStatusCodes.PEER_NOT_READY.getStatusResponse) }(context.dispatcher) case None => sender() ! PeerResponse.Response.StatusResponse(HydraStatusCodes.PEER_NOT_READY.getStatusResponse) }
步骤2:在PeerState父类中添加通用不匹配处理
abstract case class PeerState(var order: Int, var complete: Boolean, var waitForPing: Boolean) { private val baseBehavior: Receive = { // 非SessionHandshake状态收到SessionRequest时返回错误 case (manager: ActorRef, _: SessionRequest) if !this.isInstanceOf[SessionHandshake] => manager ! PeerResponse.Response.StatusResponse(HydraStatusCodes.PEER_NOT_READY.getStatusResponse) } def behavior: Receive = PartialFunction.empty orElse baseBehavior }
这样所有非SessionHandshake的状态收到SessionRequest时,都会自动返回错误,无需子类重复编写逻辑。
内容的提问来源于stack exchange,提问作者Kris Rice
相关产品推荐
相关产品推荐

