You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 09:02:13