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

Akka Typed Actor如何正确向两个其他Actor请求信息?

问题分析与优化方案

原实现的问题

你的代码功能上能实现向两个Actor请求数据的需求,但存在以下明显不足:

  • 错误处理缺失:未处理Actor查找失败、请求超时、Actor返回错误等异常场景,一旦出现问题会导致CompletableFuture永远无法完成
  • 代码冗余:重复创建messageAdapter和CompletableFuture holder,逻辑重复且繁琐
  • 不符合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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 01:32:02