【问题标题】:Running a CompletableFuture from an Actor从 Actor 运行 CompletableFuture
【发布时间】:2017-05-05 11:40:50
【问题描述】:

我正在使用 Akka (2.5.1) 的 Java (8) Spring Boot (1.5.2.RELEASE) 应用程序中使用响应式模式。它进展顺利,但现在我被困在试图从演员那里运行 CompletableFuture 。为了模拟这一点,我创建了一个返回 CompletableFuture 的非常简单的服务。但是,当我尝试将结果返回给调用控制器时,我收到有关死信的错误,并且没有返回任何响应。

我得到的错误是:

[INFO] [05/05/2017 13:12:25.650] [akka-spring-demo-akka.actor.default-dispatcher-5] [akka://akka-spring-demo/deadLetters] 消息 [从 Actor[akka://akka-spring-demo/user/$a#-1561144664] 到 Actor[akka://akka-spring-demo/deadLetters] 的 java.lang.String] 未交付。 [1] 遇到死信。可以使用配置设置“akka.log-dead-letters”和“akka.log-dead-letters-during-shutdown”关闭或调整此日志记录。

这是我的代码。这是调用actor的控制器:

@Component
@Produces(MediaType.TEXT_PLAIN)
@Path("/")
public class AsyncController {
    @Autowired
    private ActorSystem system;

    private ActorRef getGreetingActorRef() {
        ActorRef greeter = system.actorOf(SPRING_EXTENSION_PROVIDER.get(system)
                  .props("greetingActor"));

        return greeter;
    }

    @GET
    @Path("/foo")
    public void test(@Suspended AsyncResponse asyncResponse, @QueryParam("echo") String echo) {
        ask(getGreetingActorRef(), new Greet(echo), 1000)
            .thenApply((greet) -> asyncResponse.resume(Response.ok(greet).build()));
    }
}

这里是服务:

@Component
public class GreetingService {
    public CompletableFuture<String> greetAsync(String name) {
        return CompletableFuture.supplyAsync(() -> "Hello, " + name);
    }
}

那么这里是接收呼叫的演员。起初我有这个:

@Component
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public class GreetingActor extends AbstractActor {
    @Autowired
    private GreetingService greetingService;

    @Autowired
    private ActorSystem system;

    @Override
    public Receive createReceive() {
        return receiveBuilder()
                .match(Greet.class, this::onGreet)
                .build();
    }

    private void onGreet(Greet greet) {
        greetingService.greetAsync(greet.getMessage())
            .thenAccept((greetingResponse) -> getSender().tell(greetingResponse, getSelf()));
    }

}

这导致 2 个调用被正确处理,但之后我会收到死信错误。然后我在这里阅读了可能导致我的问题的原因: http://doc.akka.io/docs/akka/2.5.1/java/actors.html

警告 在使用未来回调时,在 Actor 内部,您需要小心避免关闭包含 Actor 的引用,即不要从回调内部调用方法或访问封闭 Actor 上的可变状态。这会破坏actor封装,并可能引入同步错误和竞争条件,因为回调将同时调度到封闭的actor。不幸的是,目前还没有一种方法可以在编译时检测这些非法​​访问。另请参阅:Actor 和共享可变状态

所以我认为这个想法是将结果通过管道传递给 self(),然后你可以执行 getSender().tell(response, getSelf())。

所以我把我的代码改成了这样:

@Component
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public class GreetingActor extends AbstractActor {
    @Autowired
    private GreetingService greetingService;

    @Autowired
    private ActorSystem system;

    @Override
    public Receive createReceive() {
        return receiveBuilder()
                .match(Greet.class, this::onGreet)
                .match(String.class, this::onGreetingCompleted)
                .build();
    }

    private void onGreet(Greet greet) {
        pipe(greetingService.greetAsync(greet.getMessage()), system.dispatcher()).to(getSelf());
    }

    private void onGreetingCompleted(String greetingResponse) {
        getSender().tell(greetingResponse, getSelf());
    }

}

正在使用来自 GreetingService 的响应调用 onGreetingCompleted 方法,但当时我再次收到死信错误,因此由于某种原因它无法将响应发送回调用控制器。

请注意,如果我将服务更改为:

@Component
public class GreetingService {
    public String greet(String name) {
        return "Hello, " + name;
    }
}

和演员中的onGreet:

private void onGreet(Greet greet) {
    getSender().tell(greetingService.greet(greet.getMessage()), getSelf());
}

然后一切正常。所以看起来我的基本 Java/Spring/Akka 设置正确,只是当我试图从我的演员那里调用 CompletableFuture 时,问题才开始。

任何帮助将不胜感激,谢谢!

【问题讨论】:

  • deadLetters 表示您正在向无效(过时)引用发送消息。使用 ask apttern,如果请求在响应之前超时(通过实际创建临时 actor 来请求工作),就会发生这种情况。也许发布错误会有所帮助。虽然我必须承认我没有很多想法。
  • 我已经编辑了我的帖子以显示我遇到的错误。我知道死信是什么意思以及你如何得到它们。我只是不知道为什么我会在这种特殊情况下得到它们。
  • 谢谢,我需要查看发送给 dl 的邮件类型以及实际预期的收件人是谁。希望答案能解释你的疑惑。

标签: java spring akka


【解决方案1】:

getSender 方法仅在在消息的同步处理期间可靠地返回发送者的 ref。

在第一种情况下,您有:

 greetingService.greetAsync(greet.getMessage())
        .thenAccept((greetingResponse) -> getSender().tell(greetingResponse, getSelf()));

这意味着一旦未来完成,getSender() 就会被异步调用。已经不靠谱了。您可以将其更改为:

 ActorRef sender = getSender();
 greetingService.greetAsync(greet.getMessage())
        .thenAccept((greetingResponse) -> sender.tell(greetingResponse, getSelf()));

在你的第二个例子中,你有

pipe(greetingService.greetAsync(greet.getMessage()), system.dispatcher()).to(getSelf());

您正在将响应传递给“getSelf()”,即您的工作角色。原始发件人永远不会得到任何东西(因此请求过期)。您可以将其修复为:

pipe(greetingService.greetAsync(greet.getMessage()), system.dispatcher()).to(getSender());

在第三种情况下,您在消息处理过程中同步执行了 getSender(),因此它可以工作。

【讨论】:

  • 谢谢迭戈!!这确实有效。周末愉快!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2015-09-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-01-05
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多