CompletableFuture调用阻塞问题排查与正确实现咨询
问题描述
调用test(1)时程序完全阻塞,main方法无法完成执行;但传入2时运行正常。需要明确问题的根本原因,以及正确的编码实现方式。同时,采用Holger提供的解决方案后运行恢复正常,疑惑为何该方案能捕获RuntimeException,而原方案无法做到?
原问题代码
public class Test { private static final ScheduledExecutorService EXECUTOR = Executors.newScheduledThreadPool(1, r -> { Thread t = defaultThreadFactory().newThread(r); t.setDaemon(true); return t; }); public static void main(String[] args) { Test t = new Test(); t.test(1).toCompletableFuture().join(); System.out.println("DONE"); } public CompletionStage<Void> run(int i) { if (i == 1) throw new RuntimeException(); CompletableFuture<Void> future = new CompletableFuture<>(); future.completeExceptionally(new RuntimeException()); return future; } public CompletionStage<Void> test(int i) { CompletableFuture<Void> future = new CompletableFuture<>(); EXECUTOR.schedule(() -> run(i).handle((output, error) -> { if (error instanceof CompletionException) { error = error.getCause(); } if (error != null) { CompletableFuture<Void> failedFuture = new CompletableFuture<>(); failedFuture.completeExceptionally(error); return failedFuture; } return completedFuture(output); }).thenCompose(u -> u).thenApply(future::complete).exceptionally(future::completeExceptionally), 0, TimeUnit.SECONDS); return future; } }
Holger提供的解决方案代码
public class Test { private static final ScheduledExecutorService EXECUTOR = Executors.newScheduledThreadPool(1, r -> { Thread t = defaultThreadFactory().newThread(r); t.setDaemon(true); return t; }); public static void main(String[] args) { Test t = new Test(); System.out.println("STARTED"); t.test(1).toCompletableFuture().join(); System.out.println("DONE"); } public CompletionStage<Void> run(int i) { if (i == 1) throw new RuntimeException(); CompletableFuture<Void> future = new CompletableFuture<>(); future.completeExceptionally(new RuntimeException()); return future; } public CompletionStage<Void> test(int i) { return completedFuture(i).thenComposeAsync(this::run, r -> EXECUTOR.schedule(r, 0, TimeUnit.SECONDS)) .handle((res, ex) -> { if (ex == null) return completedFuture(res); if (ex instanceof CompletionException) { ex = ex.getCause(); } CompletableFuture<Void> failedFuture = new CompletableFuture<>(); failedFuture.completeExceptionally(ex); return failedFuture; }).thenCompose(u -> u); } }
问题根本原因
原代码阻塞的核心问题
当调用test(1)时,run(1)直接抛出RuntimeException,且这个异常发生在**schedule任务的同步执行阶段**——还未进入handle方法就已经抛出。
原代码中,EXECUTOR.schedule的任务逻辑是先执行run(i),再链式调用handle等方法。但run(1)抛异常会导致整个schedule任务抛出未捕获异常,而你使用的是守护线程池:守护线程抛出未捕获异常时,不会触发默认异常处理逻辑(不打印堆栈、不终止程序),只会直接死亡。
这就导致test方法中手动创建的CompletableFuture<Void> future永远不会被complete或completeExceptionally,调用join()的main线程会一直阻塞等待这个Future完成;同时程序中没有其他非守护线程,JVM无法自动退出,最终造成永久阻塞。
而调用test(2)时,run(2)不会直接抛异常,而是返回一个已完成异常的CompletableFuture,后续handle等逻辑能正常执行,最终会调用future.completeExceptionally,join()捕获到异常后程序继续执行,打印"DONE"。
为何Holger的方案能捕获异常?
Holger的方案通过completedFuture(i).thenComposeAsync(this::run, ...)执行run方法,核心差异在于:
thenComposeAsync会将run方法的执行完全包装到异步任务逻辑中,无论run是正常返回CompletionStage还是直接抛异常,thenComposeAsync都会捕获异常并包装成CompletionException,传递给后续的handle方法处理。- 也就是说,
run(1)抛出的异常不会导致线程直接死亡,而是被转化为CompletionStage的异常状态,后续handle方法能正常处理该异常,最终返回一个完成异常的CompletionStage,join()能正确捕获异常,程序不会阻塞。
此外,方案中没有手动创建未完成的CompletableFuture,而是通过CompletionStage的链式调用自然生成最终结果,避免了手动管理Future时可能出现的"遗漏完成"问题。
内容的提问来源于stack exchange,提问作者Adhiraj

