CompletableFuture超时或异常后程序仍执行问题排查
问题分析与解决
我编写了如下Java代码,通过CompletableFuture遍历两个对象列表调用Test类的mymethod()方法。Test实例构造参数为-1时抛出异常,参数大于0时会先等待10秒再执行。期望程序在出现异常或超时时立即终止,但目前超时后程序仍持续执行,求分析原因并给出解决思路。
示例代码
import java.util.ArrayList; import java.util.List; import java.util.concurrent.*; class Test { int waitSec=0; String className; public Test(int waitSec, String className) { this.waitSec = waitSec; this.className = className; } public void mymethod() throws Exception { System.out.println("Executing mymethod in "+className); System.out.println("before "+className); if(waitSec != 0) { new CompletableFuture<>().completeAsync(() -> null, CompletableFuture.delayedExecutor(10, TimeUnit.SECONDS)).join(); } if(waitSec == -1) { throw new Exception("aaa"+className); } System.out.println("after"+className); } } public class Main2 { public static void main(String[] args) throws ExecutionException, InterruptedException { ArrayList<List<Test>> list= new ArrayList<List<Test>>(); List<Test> myTest1 = new ArrayList<>(); myTest1.add(new Test(10,"Class1")); myTest1.add(new Test(0, "Class2")); List<Test> myTest2 = new ArrayList<>(); myTest2.add(new Test(0, "Class3")); myTest2.add(new Test(0, "Class4")); list.add(myTest1); list.add(myTest2); ThreadPoolExecutor executor = (ThreadPoolExecutor) Executors.newFixedThreadPool(2); List<CompletableFuture<Void>> cp = new ArrayList<>(); for (List<Test> myTest : list) { for (Test test : myTest) { CompletableFuture<Void> future = CompletableFuture.runAsync(() -> { try { test.mymethod(); } catch (Exception e) { System.out.print(">>>>>>>>>>>>>>>>exception<<<<<<<<<<<<<<<<"); handleFailure(e,"exception"); } }, executor); future.exceptionally(t -> {handleFailure(t,"timeout"); System.out.print(">>>>>>>>>>>>>>>>timeout<<<<<<<<<<<<<<<<"); return null; }); //future.orTimeout(5,TimeUnit.SECONDS); cp.add(future); } CompletableFuture<Void> allOf = CompletableFuture.allOf(cp.toArray(new CompletableFuture[0])); allOf.orTimeout(5,TimeUnit.SECONDS); allOf.get(); } executor.shutdown(); } private static void handleFailure(Throwable t, String myString){ System.out.println(">>>>>>>>>>FAILED<<<<<<<<<<"+ myString); } }
原因分析
orTimeout逻辑未生效:CompletableFuture是不可变对象,allOf.orTimeout(5, TimeUnit.SECONDS)会返回一个新的带超时逻辑的Future实例,但你没有接收这个新实例,而是继续调用原allOf的get(),导致超时判断完全没起作用。- 异常未传递到Future:在
runAsync的任务内部,你用try-catch捕获了mymethod()抛出的异常,直接处理但没有将异常传递给CompletableFuture,导致future.exceptionally永远不会触发,异常无法终止流程。 - 线程池未强制终止:即使超时逻辑生效,默认的
executor.shutdown()会等待所有已提交任务执行完毕,不会主动中断正在运行的长耗时任务,所以等待10秒的任务会继续执行。
解决思路
1. 正确使用orTimeout
必须接收orTimeout返回的新Future实例,并在该实例上调用get(),确保超时逻辑生效:
CompletableFuture<Void> allOfWithTimeout = CompletableFuture.allOf(cp.toArray(new CompletableFuture[0])) .orTimeout(5, TimeUnit.SECONDS); allOfWithTimeout.get();
2. 让异常传递到CompletableFuture
移除任务内部的try-catch,或者捕获异常后主动抛出CompletionException(CompletableFuture能识别的异常类型),让exceptionally能捕获到异常:
CompletableFuture<Void> future = CompletableFuture.runAsync(() -> { try { test.mymethod(); } catch (Exception e) { // 将异常包装为CompletionException传递给Future throw new CompletionException(e); } }, executor);
3. 超时/异常时强制终止线程池
当触发超时或异常时,调用executor.shutdownNow()中断所有正在执行的任务,而不是等待任务完成:
try { allOfWithTimeout.get(); } catch (TimeoutException e) { System.out.println("超时触发,强制终止所有任务"); executor.shutdownNow(); executor.awaitTermination(2, TimeUnit.SECONDS); return; } catch (ExecutionException | InterruptedException e) { System.out.println("任务异常,强制终止所有任务"); executor.shutdownNow(); executor.awaitTermination(2, TimeUnit.SECONDS); return; }
4. 单个任务设置超时(可选)
如果需要单个任务也有超时限制,给每个future单独添加orTimeout:
CompletableFuture<Void> future = CompletableFuture.runAsync(...) .orTimeout(5, TimeUnit.SECONDS);
修改后的完整代码
import java.util.ArrayList; import java.util.List; import java.util.concurrent.*; class Test { int waitSec = 0; String className; public Test(int waitSec, String className) { this.waitSec = waitSec; this.className = className; } public void mymethod() throws Exception { System.out.println("Executing mymethod in " + className); System.out.println("before " + className); if (waitSec != 0) { new CompletableFuture<>().completeAsync(() -> null, CompletableFuture.delayedExecutor(10, TimeUnit.SECONDS)).join(); } if (waitSec == -1) { throw new Exception("aaa" + className); } System.out.println("after" + className); } } public class Main2 { public static void main(String[] args) { ArrayList<List<Test>> list = new ArrayList<List<Test>>(); List<Test> myTest1 = new ArrayList<>(); myTest1.add(new Test(10, "Class1")); myTest1.add(new Test(0, "Class2")); List<Test> myTest2 = new ArrayList<>(); myTest2.add(new Test(0, "Class3")); myTest2.add(new Test(0, "Class4")); list.add(myTest1); list.add(myTest2); ThreadPoolExecutor executor = (ThreadPoolExecutor) Executors.newFixedThreadPool(2); List<CompletableFuture<Void>> cp = new ArrayList<>(); for (List<Test> myTest : list) { for (Test test : myTest) { CompletableFuture<Void> future = CompletableFuture.runAsync(() -> { try { test.mymethod(); } catch (Exception e) { throw new CompletionException(e); } }, executor) .orTimeout(5, TimeUnit.SECONDS); future.exceptionally(t -> { if (t instanceof CompletionException) { handleFailure(t.getCause(), "exception"); System.out.println(">>>>>>>>>>>>>>>>exception<<<<<<<<<<<<<<<<"); } else if (t instanceof TimeoutException) { handleFailure(t, "timeout"); System.out.println(">>>>>>>>>>>>>>>>timeout<<<<<<<<<<<<<<<<"); } return null; }); cp.add(future); } try { CompletableFuture<Void> allOfWithTimeout = CompletableFuture.allOf(cp.toArray(new CompletableFuture[0])) .orTimeout(5, TimeUnit.SECONDS); allOfWithTimeout.get(); } catch (TimeoutException e) { System.out.println("整体超时,终止所有任务"); executor.shutdownNow(); executor.awaitTermination(2, TimeUnit.SECONDS); return; } catch (ExecutionException | InterruptedException e) { System.out.println("任务异常,终止所有任务"); executor.shutdownNow(); executor.awaitTermination(2, TimeUnit.SECONDS); return; } } executor.shutdown(); } private static void handleFailure(Throwable t, String myString) { System.out.println(">>>>>>>>>>FAILED<<<<<<<<<<" + myString + ": " + t.getMessage()); } }
内容的提问来源于stack exchange,提问作者user3756792
相关产品推荐
相关产品推荐

