Akka Typed Actor如何正确向两个其他Actor请求信息?
问题分析与优化方案
原实现的问题
你的代码功能上能实现向两个Actor请求数据的需求,但存在以下明显不足:
- 错误处理缺失:未处理Actor查找失败、请求超时、Actor返回错误等异常场景,一旦出现问题会导致CompletableFuture永远无法完成
- 代码冗余:重复创建
messageAdapter和CompletableFutureholder,逻辑重复且繁琐 - 不符合Actor模型最佳实践:过度依赖CompletableFuture回调,未利用Akka Actor的状态管理能力,易引发回调地狱
优化方案
方案1:简化CompletableFuture组合逻辑
利用Akka ask API本身返回CompletableFuture的特性,结合CompletableFuture的组合式API减少冗余,同时补充完整的错误处理:
import akka.actor.typed.ActorRef; import akka.actor.typed.Behavior; import akka.actor.typed.Duration; import akka.actor.typed.javadsl.AbstractBehavior; import akka.actor.typed.javadsl.ActorContext; import akka.actor.typed.javadsl.Behaviors; import akka.actor.typed.javadsl.Receive; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class ActorA extends AbstractBehavior<ActorA.Command> { private static final Logger log = LoggerFactory.getLogger(ActorA.class); public interface Command {} public record ExampleCommand() implements Command {} private record WrappedBResponse(String response) implements Command {} private record WrappedCResponse(String response) implements Command {} private record RequestFailed(Throwable error) implements Command {} private ActorA(ActorContext<Command> context) { super(context); } public static Behavior<Command> create() { return Behaviors.setup(ActorA::new); } @Override public Receive<Command> createReceive() { return newReceiveBuilder() .onMessage(ExampleCommand.class, this::onProcessVote) .onMessage(RequestFailed.class, this::onRequestFailed) .build(); } private Behavior<Command> onProcessVote(ExampleCommand command) { ActorContext<Command> ctx = getContext(); // 一次性创建适配器,复用给所有请求 ActorRef<String> bReplyAdapter = ctx.messageAdapter(String.class, WrappedBResponse::new); ActorRef<String> cReplyAdapter = ctx.messageAdapter(String.class, WrappedCResponse::new); // 合并Actor查找与请求逻辑,链式处理 CompletionStage<String> bPropFuture = ActorUtil.findActor(ctx.getSystem(), ActorB.serviceKey) .thenCompose(actorBRef -> ctx.ask( String.class, actorBRef, Duration.ofSeconds(1), replyTo -> new ActorB.GetStringProp(bReplyAdapter), (response, error) -> { if (error != null) { ctx.self().tell(new RequestFailed(error)); return null; } return response; })); CompletionStage<String> cPropFuture = ActorUtil.findActor(ctx.getSystem(), ActorC.serviceKey) .thenCompose(actorCRef -> ctx.ask( String.class, actorCRef, Duration.ofSeconds(1), replyTo -> new ActorC.GetStringProp(cReplyAdapter), (response, error) -> { if (error != null) { ctx.self().tell(new RequestFailed(error)); return null; } return response; })); // 组合两个请求结果,统一处理 CompletableFuture.allOf(bPropFuture.toCompletableFuture(), cPropFuture.toCompletableFuture()) .thenAccept(v -> { try { String bProp = bPropFuture.toCompletableFuture().get(); String cProp = cPropFuture.toCompletableFuture().get(); log.info("Got property from ActorB: '{}' and property from ActorC: '{}'", bProp, cProp); // 执行后续业务逻辑 } catch (Exception e) { ctx.self().tell(new RequestFailed(e)); } }); return Behaviors.same(); } private Behavior<Command> onRequestFailed(RequestFailed error) { log.error("Request failed", error.error()); // 实现错误处理逻辑:重试、通知上游、记录告警等 return Behaviors.same(); } }
方案2:使用Actor状态机(更符合Actor模型设计)
当需要等待多个响应时,推荐切换到专门的等待状态,利用Actor的状态管理能力替代回调,更贴合Akka的消息驱动设计:
import akka.actor.typed.ActorRef; import akka.actor.typed.Behavior; import akka.actor.typed.Duration; import akka.actor.typed.javadsl.AbstractBehavior; import akka.actor.typed.javadsl.ActorContext; import akka.actor.typed.javadsl.Behaviors; import akka.actor.typed.javadsl.Receive; import java.util.Optional; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class ActorA extends AbstractBehavior<ActorA.Command> { private static final Logger log = LoggerFactory.getLogger(ActorA.class); public interface Command {} public record ExampleCommand() implements Command {} private record WrappedBResponse(String response) implements Command {} private record WrappedCResponse(String response) implements Command {} private record RequestTimeout() implements Command {} // 等待响应时的状态载体 private static class WaitingState { final Optional<String> bProp; final Optional<String> cProp; WaitingState(Optional<String> bProp, Optional<String> cProp) { this.bProp = bProp; this.cProp = cProp; } } private ActorA(ActorContext<Command> context) { super(context); } public static Behavior<Command> create() { return Behaviors.setup(ActorA::new); } @Override public Receive<Command> createReceive() { return newReceiveBuilder() .onMessage(ExampleCommand.class, this::onProcessVote) .build(); } private Behavior<Command> onProcessVote(ExampleCommand command) { ActorContext<Command> ctx = getContext(); ActorRef<String> bReplyAdapter = ctx.messageAdapter(String.class, WrappedBResponse::new); ActorRef<String> cReplyAdapter = ctx.messageAdapter(String.class, WrappedCResponse::new); // 发起ActorB请求,处理查找失败场景 ActorUtil.findActor(ctx.getSystem(), ActorB.serviceKey) .whenComplete((actorBRef, error) -> { if (error != null) { log.error("Failed to find ActorB", error); ctx.self().tell(new WrappedBResponse(null)); } else { actorBRef.tell(new ActorB.GetStringProp(bReplyAdapter)); } }); // 发起ActorC请求,处理查找失败场景 ActorUtil.findActor(ctx.getSystem(), ActorC.serviceKey) .whenComplete((actorCRef, error) -> { if (error != null) { log.error("Failed to find ActorC", error); ctx.self().tell(new WrappedCResponse(null)); } else { actorCRef.tell(new ActorC.GetStringProp(cReplyAdapter)); } }); // 设置全局超时 ctx.scheduleOnce(Duration.ofSeconds(1), ctx.self(), new RequestTimeout()); // 切换到等待响应的状态 return waitingForResponses(new WaitingState(Optional.empty(), Optional.empty())); } private Behavior<Command> waitingForResponses(WaitingState currentState) { return Behaviors.receive(Command.class) .onMessage(WrappedBResponse.class, (ctx, msg) -> { WaitingState newState = new WaitingState(Optional.ofNullable(msg.response()), currentState.cProp); return checkCompletion(ctx, newState); }) .onMessage(WrappedCResponse.class, (ctx, msg) -> { WaitingState newState = new WaitingState(currentState.bProp, Optional.ofNullable(msg.response())); return checkCompletion(ctx, newState); }) .onMessage(RequestTimeout.class, ctx -> { log.error("Request timed out waiting for ActorB/C responses"); // 处理超时逻辑 return create(); }) .build(); } private Behavior<Command> checkCompletion(ActorContext<Command> ctx, WaitingState state) { if (state.bProp.isPresent() && state.cProp.isPresent()) { String bProp = state.bProp.get(); String cProp = state.cProp.get(); if (bProp != null && cProp != null) { log.info("Got property from ActorB: '{}' and property from ActorC: '{}'", bProp, cProp); // 执行后续业务逻辑 } else { log.error("One or more requests returned null result"); // 处理部分请求失败的逻辑 } // 回到初始状态 return create(); } else { // 继续等待剩余响应 return waitingForResponses(state); } } }
关键知识点总结
- 优先用Actor状态管理:Actor模型的核心是状态+消息驱动,等待多响应时切换到专门状态,比回调更符合Akka设计哲学
- 错误处理必须覆盖:Actor查找、请求超时、响应异常等场景都要处理,避免出现"挂起"的请求
- 复用消息适配器:不需要为每个请求创建新的
messageAdapter,一次创建即可复用 - 合理利用ask API:ask已经封装了请求-响应逻辑,无需手动创建CompletableFuture容器
内容的提问来源于stack exchange,提问作者Steve Storck
相关产品推荐
相关产品推荐

