Akka中getContext().become()热切换后无法接收消息问题
在Akka中编写PingPong Actor时,使用become()切换行为后发现Actor仅能处理第一条消息,后续发送的消息均被忽略。代码及运行结果如下:
package ping_pong; import akka.actor.AbstractActor; import akka.actor.ActorRef; import akka.actor.ActorSystem; import akka.actor.Props; public class PingPongActor extends AbstractActor { public static void main(String[] args) throws Exception { ActorSystem _system = ActorSystem.create("PingPongActorApp"); ActorRef masterpp = _system.actorOf(Props.create(PingPongActor.class), "pp"); masterpp.tell(PING, masterpp); System.out.println("after first msg"); masterpp.tell(PING, masterpp); masterpp.tell(PING, masterpp); masterpp.tell(PING, masterpp); masterpp.tell(PING, masterpp); masterpp.tell(PING, masterpp); System.out.println("last msg"); } static String PING = "PING"; static String PONG = "PONG"; int count = 0; @Override public Receive createReceive() { return receiveBuilder().match(String.class, ua -> { if (ua.matches(PING)) { System.out.println("PING" + count); count += 1; Thread.sleep(100); if (count <= 10) { getSelf().tell(PONG, getSelf()); } getContext().become(receiveBuilder().match(String.class, ua1 -> { if (ua1.matches(PONG)) { System.out.println("PONG" + count); count += 1; Thread.sleep(100); getContext().unbecome(); } }).build()); if (count > 10) { System.out.println("DONE" + count); getContext().stop(getSelf()); } } }).build(); } }
运行结果:
21:36:34.098 [PingPongActorApp-akka.actor.default-dispatcher-4] INFO akka.event.slf4j.Slf4jLogger - Slf4jLogger started
after first msg
last msg
PING0
PONG1
核心原因
行为切换后的消息匹配失效
Actor处理第一条PING消息后,调用getContext().become()切换到仅处理PONG的行为。此时main方法发送的后续5条PING消息已进入Actor邮箱,但新行为只匹配PONG类型消息,这些PING无法被匹配,会触发Actor默认的unhandled()逻辑(默认记录警告日志,若日志级别未配置则无显示),最终被忽略。消息顺序与行为恢复时机不匹配
Actor处理第一条PING时发送的PONG消息会被追加到邮箱末尾,排在main发送的所有PING之后。Actor会先逐个尝试处理这些PING(均无法匹配当前行为),直到处理到PONG时才会调用unbecome()恢复初始行为,但此时邮箱已无未处理消息,后续不会再执行任何逻辑。外部消息发送逻辑违背状态流转
你的PingPong逻辑设计为Actor内部循环发送PING/PONG,但main直接批量发送PING,违背了Actor"初始状态处理PING→切换状态等待PONG"的流转预期,导致状态切换后无法处理外部涌入的不符合状态的消息。
修正方案
调整行为切换逻辑,让Actor内部控制消息流转,避免外部发送不符合状态的消息:
package ping_pong; import akka.actor.AbstractActor; import akka.actor.ActorRef; import akka.actor.ActorSystem; import akka.actor.Props; public class PingPongActor extends AbstractActor { public static void main(String[] args) throws Exception { ActorSystem _system = ActorSystem.create("PingPongActorApp"); ActorRef masterpp = _system.actorOf(Props.create(PingPongActor.class), "pp"); // 仅发送初始启动消息,后续流转由Actor内部控制 masterpp.tell(PING, masterpp); System.out.println("Initial PING sent"); } static String PING = "PING"; static String PONG = "PONG"; int count = 0; @Override public Receive createReceive() { // 初始行为:处理PING return receiveBuilder() .match(String.class, ua -> ua.equals(PING), this::handlePing) .build(); } private void handlePing(String msg) { System.out.println("PING" + count); count += 1; try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } if (count > 10) { System.out.println("DONE" + count); getContext().stop(getSelf()); return; } // 发送PONG并切换到处理PONG的行为 getSelf().tell(PONG, getSelf()); getContext().become(receiveBuilder() .match(String.class, ua -> ua.equals(PONG), this::handlePong) .build()); } private void handlePong(String msg) { System.out.println("PONG" + count); count += 1; try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } if (count > 10) { System.out.println("DONE" + count); getContext().stop(getSelf()); return; } // 发送PING并切换回处理PING的行为 getSelf().tell(PING, getSelf()); getContext().unbecome(); } }
修正说明
- 拆分
handlePing和handlePong方法,逻辑更清晰。 - 仅在
main中发送初始PING,后续消息由Actor内部循环发送,符合状态流转预期。 - 每个状态仅处理对应类型消息,避免外部消息干扰。
- 增加
Thread.sleep()的中断异常处理,避免影响Actor运行。
内容的提问来源于stack exchange,提问作者martinnovoty

