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

Akka Actor处理PeerChat请求时出现死信问题求助

问题排查:PeerChat响应未被PeerManager处理

背景

业务流程如下:

  • 实例A的Peer Actor生成PeerRequest proto消息,通过HTTP发送至实例B
  • 实例B的PeerServer接收并处理该请求
  • PeerServer通过Akka的asks(?操作符)请求PeerManager Actor处理请求
  • PeerManager通过asks(?)向指定ID的Peer Actor(ID来自PeerRequest字段)发送PeerRequest.PeerChat消息
  • Peer Actor完成解密、聊天内容入库、加密响应操作后,回复PeerChatResponse

对应代码实现:
PeerServer处理逻辑

else if (peerRequest.request.isPeerChat) {
  // forward the PeerChat to PeerManager
  val future: Future[PeerResponse] = (peerManager ? peerRequest.getPeerChat).mapTo[PeerResponse]
  
  try {
    val result = Await.result(future, timeout.duration)
    resultingRoute = complete(result)
  } catch {
    case e: TimeoutException =>
      resultingRoute = complete(PeerResponse().withStatusResponse(HydraStatusCodes.PEER_MANAGER_TIMEOUT.getStatusResponse))
  }
}

PeerManager处理PeerChat的逻辑

case peerChat@PeerChat(_, serverID, _, _) =>
  val peerEntry = peerMap.get(serverID)

  if (peerEntry.isDefined) {
    try {
      val state = Await.result((peerEntry.get ? QueryState()).mapTo[PeerState], timeout.duration)

      if (state != null && state.name == "WaitMessages") {
        val peerResponse = Await.result((peerEntry.get ? (self, peerChat)).mapTo[PeerChatResponse], timeout.duration)
        sender() ! PeerResponse().withPeerChatResponse(peerResponse)
      }
    } catch {
      case e: TimeoutException =>
        sender() ! PeerResponse().withStatusResponse(HydraStatusCodes.PEER_MANAGER_TIMEOUT.getStatusResponse)
    }
  } else {
    sender() ! PeerResponse().withStatusResponse(HydraStatusCodes.PEER_NOT_READY.getStatusResponse)
  }

Peer处理(ActorRef, PeerChat)的逻辑

case (sender: ActorRef, peerChat@PeerChat(chatID, _, encryptedByteString, _)) =>
  system.log.info("processing PeerChat")
  try {
    // decrypt and store the chat in the database
    val decodedData = Base64.getDecoder.decode(encryptedByteString)
    val unencryptedBytes = peer.getEncryptionManager.decrypt(decodedData)
    // re-form the bytes back into a ChatRequest
    val chatRequest: ChatRequest = ChatRequest.parseFrom(unencryptedBytes)
    // verify we don't already have this chat in the database

    val database = DatabaseUtil.getInstance
    val chats = database.queryByFields(classOf[Chat], Map(
      "serverID" -> Seq(chatRequest.serverID),
      "timestamp" -> Seq(chatRequest.timestamp),
      "text" -> Seq(chatRequest.chatText)
    ))

    if (chats.isEmpty) {
      // store the chat
      val chat = new Chat(chatRequest)
      chat.sent = true
      chat.confirmed = true
      database.addWithGeneratedID(chat)
    }
    // tell the sender it was successful
    val encryptedResponse = peer.getEncryptionManager.encrypt(ByteBuffer.allocate(4).putInt(chatID).array())
    val encodedData = Base64.getEncoder.encodeToString(encryptedResponse)
    val peerChatResponse = PeerChatResponse(encodedData)

    system.log.info(s"Replying with PeerChatResponse:${peerChatResponse.toProtoString}")

    sender ! peerChatResponse
  } catch {
    case e: Exception =>
      system.log.error(e, "Failed processing PeerChat")
  }

注:传递PeerManager的self是因为Peer的receive行为定义在State类中,无法直接访问默认sender()。

问题描述

当Peer Actor收到来自PeerManager的(ActorRef, PeerChat)消息并回复PeerChatResponse时,出现以下错误:

Message [peer.chat.PeerChatResponse.PeerChatResponse] to Actor[akka://system-actor/user/peer_manager#1099549428] was unhandled.

其他同类流程(如传递SessionRequest消息)可正常执行,Await.result能正确处理Peer的响应,但PeerChat流程中,Peer回复的消息直接发送给PeerManager后未被处理,需排查原因。

问题原因

核心问题在于消息接收方不匹配:

  1. PeerManager使用(peerEntry.get ? (self, peerChat))发起asks请求时,Akka会自动创建一个临时PromiseActorRef来接收Peer的响应,而非用PeerManager自身(self)接收。
  2. 但Peer Actor却把PeerChatResponse发送给了self(PeerManager实例),而PeerManager的receive方法中没有定义处理PeerChatResponse消息的分支,导致消息未被处理,抛出未处理错误。
  3. 对比正常工作的SessionRequest流程,Peer回复的消息是发送给asks生成的临时ActorRef,因此能被Await.result正确捕获。

解决方法

方案一:修正Peer的回复目标(推荐)

不需要传递PeerManager的self,直接让Peer回复给asks请求的默认sender(临时PromiseActorRef)。如果Peer的State类无法访问默认sender(),调整消息结构即可:

  • 修改PeerManager发送的消息,去掉self参数:
// PeerManager中的代码修改
val peerResponse = Await.result((peerEntry.get ? peerChat).mapTo[PeerChatResponse], timeout.duration)
  • 修改Peer的receive逻辑,直接处理PeerChat消息并回复默认sender:
// Peer中的代码修改
case peerChat@PeerChat(chatID, _, encryptedByteString, _) =>
  // ... 原有处理逻辑不变
  sender() ! peerChatResponse

方案二:在PeerManager中添加消息处理分支(不推荐)

如果必须传递self,则需要在PeerManager的receive方法中添加PeerChatResponse的处理逻辑,同时需要手动跟踪请求上下文(比如关联请求ID),但这种方式会增加代码复杂度,违背Akka异步设计原则。

额外优化建议

避免在Actor的receive方法中使用Await.result,这会阻塞Actor线程,推荐改用链式Future+pipeTo实现异步非阻塞处理:

// PeerManager中的优化示例
case peerChat@PeerChat(_, serverID, _, _) =>
  val peerEntry = peerMap.get(serverID)
  val originalSender = sender()

  if (peerEntry.isDefined) {
    (peerEntry.get ? QueryState()).mapTo[PeerState].flatMap { state =>
      if (state != null && state.name == "WaitMessages") {
        (peerEntry.get ? peerChat).mapTo[PeerChatResponse].map { resp =>
          PeerResponse().withPeerChatResponse(resp)
        }
      } else {
        Future.successful(PeerResponse().withStatusResponse(/* 对应状态码 */))
      }
    }.recover {
      case _: TimeoutException =>
        PeerResponse().withStatusResponse(HydraStatusCodes.PEER_MANAGER_TIMEOUT.getStatusResponse)
    }.pipeTo(originalSender)
  } else {
    originalSender ! PeerResponse().withStatusResponse(HydraStatusCodes.PEER_NOT_READY.getStatusResponse)
  }

内容的提问来源于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 10:25:21