Akka Actor处理PeerChat请求时出现死信问题求助
问题排查:PeerChat响应未被PeerManager处理
背景
业务流程如下:
- 实例A的Peer Actor生成
PeerRequestproto消息,通过HTTP发送至实例B - 实例B的
PeerServer接收并处理该请求 PeerServer通过Akka的asks(?操作符)请求PeerManagerActor处理请求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后未被处理,需排查原因。
问题原因
核心问题在于消息接收方不匹配:
- PeerManager使用
(peerEntry.get ? (self, peerChat))发起asks请求时,Akka会自动创建一个临时PromiseActorRef来接收Peer的响应,而非用PeerManager自身(self)接收。 - 但Peer Actor却把
PeerChatResponse发送给了self(PeerManager实例),而PeerManager的receive方法中没有定义处理PeerChatResponse消息的分支,导致消息未被处理,抛出未处理错误。 - 对比正常工作的
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
相关产品推荐
相关产品推荐

