CompletableFuture 是异步的。但它是非阻塞的吗?
关于 CompletableFuture 的一个正确之处在于它是真正的异步的,它允许您从调用者线程异步运行您的任务,并且 API (例如 thenXXX)允许您在结果可用时对其进行处理。另一方面,CompletableFuture 并不总是非阻塞的。例如,当你运行以下代码时,会在默认的ForkJoinPool上异步执行:
CompletableFuture.supplyAsync(() -> {
try {
Thread.sleep(1000);
}
catch (InterruptedException e) {
}
return 1;
});
很明显ForkJoinPool中执行任务的Thread最终会被阻塞,也就是说我们不能保证调用是非阻塞的。
另一方面,CompletableFuture 公开了 API,允许您使其真正实现非阻塞。
例如,您始终可以执行以下操作:
public CompletableFuture myNonBlockingHttpCall(Object someData) {
var uncompletedFuture = new CompletableFuture(); // creates uncompleted future
myAsyncHttpClient.execute(someData, (result, exception -> {
if(exception != null) {
uncompletedFuture.completeExceptionally(exception);
return;
}
uncompletedFuture.complete(result);
})
return uncompletedFuture;
}
如您所见,CompletableFuturefuture 的 API 为您提供了 complete 和 completeExceptionally 方法,可以在需要时完成您的执行,而不会阻塞任何线程。
单声道与 CompletableFuture
在上一节中,我们概述了 CF 行为,但 CompletableFuture 和 Mono 之间的核心区别是什么?
值得一提的是,我们也可以阻塞 Mono。没有人阻止我们编写以下内容:
Mono.fromCallable(() -> {
try {
Thread.sleep(1000);
}
catch (InterruptedException e) {
}
return 1;
})
当然,一旦我们订阅了future,调用者线程就会被阻塞。但是我们总是可以通过提供一个额外的subscribeOn 运算符来解决这个问题。尽管如此,Mono 更广泛的 API 并不是关键特性。
为了了解CompletableFuture 和Mono 之间的主要区别,让我们回到前面提到的myNonBlockingHttpCall 方法实现。
public CompletableFuture myUpperLevelBusinessLogic() {
var future = myNonBlockingHttpCall();
// ... some code
if (something) {
// oh we don't really need anything, let's just throw an exception
var errorFuture = new CompletableFuture();
errorFuture.completeExceptionally(new RuntimeException());
return errorFuture;
}
return future;
}
在CompletableFuture 的情况下,一旦调用该方法,它就会急切地执行对另一个服务/资源的HTTP 调用。即使我们在验证了一些前置/后置条件后并不真正需要执行结果,但它会开始执行,并且将为这项工作分配额外的 CPU/DB-Connections/What-Ever-Machine-Resources。
相比之下,Mono 类型根据定义是惰性的:
public Mono myNonBlockingHttpCallWithMono(Object someData) {
return Mono.create(sink -> {
myAsyncHttpClient.execute(someData, (result, exception -> {
if(exception != null) {
sink.error(exception);
return;
}
sink.success(result);
})
});
}
public Mono myUpperLevelBusinessLogic() {
var mono = myNonBlockingHttpCallWithMono();
// ... some code
if (something) {
// oh we don't really need anything, let's just throw an exception
return Mono.error(new RuntimeException());
}
return mono;
}
在这种情况下,在订阅最终的mono 之前什么都不会发生。因此,只有当myNonBlockingHttpCallWithMono方法返回的Mono被订阅时,提供给Mono.create(Consumer)的逻辑才会被执行。
我们可以走得更远。我们可以让我们的执行更加懒惰。您可能知道,Mono 扩展了 Reactive Streams 规范中的 Publisher。 Reactive Streams 的尖叫特性是背压支持。因此,使用Mono API,我们可以仅在真正需要数据并且我们的订阅者准备好使用它们时执行:
Mono.create(sink -> {
AtomicBoolean once = new AtomicBoolean();
sink.onRequest(__ -> {
if(!once.get() && once.compareAndSet(false, true) {
myAsyncHttpClient.execute(someData, (result, exception -> {
if(exception != null) {
sink.error(exception);
return;
}
sink.success(result);
});
}
});
});
在此示例中,我们仅在订阅者调用 Subscription#request 时执行数据,因此它声明已准备好接收数据。
总结
-
CompletableFuture 是异步的,可以是非阻塞的
-
CompletableFuture 很渴望。你不能推迟执行。但是你可以取消它们(这总比没有好)
-
Mono 是异步/非阻塞的,可以通过使用不同的运算符组合主 Mono,轻松地对不同的 Thread 执行任何调用。
-
Mono 真的很懒惰,它允许根据订阅者的存在及其消费数据的准备情况推迟执行启动。