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

如何为可变更行为的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 09:11:02