【发布时间】: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 的邮件类型以及实际预期的收件人是谁。希望答案能解释你的疑惑。