【问题标题】:Stopping a thread in java CompletableFuture after timeout超时后停止java CompletableFuture中的线程
【发布时间】:2021-02-20 21:10:01
【问题描述】:

我的 java 代码中有一个异步链,我想在某个超时后停止 所以我用一些线程创建了一个线程池,并像这样调用 CompletableFuture

ExecutorService pool = Executors.newFixedThreadPool(10);

比我有一个循环方法,它从数据库加载数据并在其上执行一些任务,一旦所有 CompletableFutures 完成,它就会再次执行它

CompletableFuture<MyObject> futureTask =
                CompletableFuture.supplyAsync(() -> candidate, pool)
                .thenApply(Task1::doWork).thenApply(Task2::doWork).thenApply(Task3::doWork)
                .thenApply(Task4::doWork).thenApply(Task5::doWork).orTimeout(30,TimeUnit.SECONDS)
                .thenApply(Task6::doWork).orTimeout(30,TimeUnit.SECONDS)
                .exceptionally(ExceptionHandlerService::handle);

我的问题在于 task6,它的任务非常密集(它的网络连接任务有时会永远挂起) 我注意到我的 orTimeout 在 30 秒后被正确触发,但运行 Task6 的线程仍在运行

这样几个周期后,我所有的线程都被耗尽了,我的应用程序死了

超时后如何取消池中正在运行的线程? (不调用 pool.shutdown())

更新* 在主线程中,我做了一个简单的检查,如下所示

for (int i = TIME_OUT_SECONDS; i >= 0; i--) {
                unfinishedTasks = handleFutureTasks(unfinishedTasks, totalBatchSize);
                if(unfinishedTasks.isEmpty()) {
                    break;
                }
                if(i==0) {
                    //handle cancelation of the tasks
                    for(CompletableFuture<ComplianceCandidate> task: unfinishedTasks) {
                        **task.cancel(true);**
                        log.error("Reached timeout on task, is canceled: {}", task.isCancelled());
                    }
                    break;
                }
                try {
                    TimeUnit.SECONDS.sleep(1);
                } catch (Exception ex) {
                }
            }

我看到的是,几个周期后,所有任务都抱怨超时...... 在前 1-2 个周期中,我仍然得到预期的响应(虽然有线程来处理它)

我还是觉得线程池耗尽了

【问题讨论】:

  • 有可能,但不太可能,网络线程处于不间断睡眠状态,在这种睡眠状态下实际上不可能杀死它们。值得检查这是否没有发生。 htop 和其他类似工具可以显示线程是否处于不间​​断睡眠状态。
  • 即使评估挂在一个可中断的操作中,你也使用了错误的工具来完成这项工作。只需考虑the documentation of CompletableFuture’s cancel:“参数mayInterruptIfRunning - 此值在此实现中无效,因为中断不用于控制处理。”也就是说cancel(true)还是cancel(false)都无所谓,CompletableFuture根本不支持打断。

标签: java multithreading concurrency completable-future


【解决方案1】:

我知道你说没有打电话给pool.shutDown,但根本没有别的办法。但是,当您查看您的阶段时,它们将在“附加”它们的线程(添加那些thenApply)或您定义的池中的线程中运行。可能是一个更有意义的例子。

public class SO64743332 {

    static ExecutorService pool = Executors.newFixedThreadPool(10);

    public static void main(String[] args) {

        CompletableFuture<String> f1 = CompletableFuture.supplyAsync(() -> dbCall(), pool);

        //simulateWork(4);

        CompletableFuture<String> f2 = f1.thenApply(x -> {
            System.out.println(Thread.currentThread().getName());
            return transformationOne(x);
        });

        CompletableFuture<String> f3 = f2.thenApply(x -> {
            System.out.println(Thread.currentThread().getName());
            return transformationTwo(x);
        });

        f3.join();
    }

    private static String dbCall() {
        simulateWork(2);
        return "a";
    }

    private static String transformationOne(String input) {
        return input + "b";
    }

    private static String transformationTwo(String input) {
        return input + "b";
    }

    private static void simulateWork(int seconds) {
        try {
            Thread.sleep(TimeUnit.SECONDS.toMillis(seconds));
        } catch (InterruptedException e) {
            System.out.println("Interrupted!");
            e.printStackTrace();
        }
    }
}

上面代码的关键点是:simulateWork(4);。运行代码并将其注释掉,然后取消注释。看看哪个线程实际上将执行所有这些thenApply。它是 main 或池中的 same 线程,这意味着尽管您定义了一个池 - 它只是该池中的一个线程,它将执行所有这些阶段。

在这种情况下,您可以定义一个将运行所有这些阶段的单线程执行程序(比方说在一个方法中)。这样您就可以控制何时调用shutDownNow 并可能中断(如果您的代码响应中断)正在运行的任务。这是一个模拟的虚构示例:

public class SO64743332 {

    public static void main(String[] args) {
        execute();
    }


    public static void execute() {

        ExecutorService pool = Executors.newSingleThreadExecutor();

        CompletableFuture<String> cf1 = CompletableFuture.supplyAsync(() -> dbCall(), pool);
        CompletableFuture<String> cf2 = cf1.thenApply(x -> transformationOne(x));

        // give enough time for transformationOne to start, but not finish
        simulateWork(2);

        try {
            CompletableFuture<String> cf3 = cf2.thenApply(x -> transformationTwo(x))
                                               .orTimeout(4, TimeUnit.SECONDS);
            cf3.get(10, TimeUnit.SECONDS);
        } catch (ExecutionException | InterruptedException | TimeoutException e) {
            pool.shutdownNow();
        }

    }

    private static String dbCall() {
        System.out.println("Started DB call");
        simulateWork(1);
        System.out.println("Done with DB call");
        return "a";
    }

    private static String transformationOne(String input) {
        System.out.println("Started work");
        simulateWork(10);
        System.out.println("Done work");
        return input + "b";
    }

    private static String transformationTwo(String input) {
        System.out.println("Started transformation two");
        return input + "b";
    }

    private static void simulateWork(int seconds) {
        try {
            Thread.sleep(TimeUnit.SECONDS.toMillis(seconds));
        } catch (InterruptedException e) {
            System.out.println("Interrupted!");
            e.printStackTrace();
        }
    }
}

运行这个你应该注意到transformationOne开始了,但是因为shutDownNow而被中断了。

这样的缺点应该很明显,每次调用execute都会创建一个新的线程池……

【讨论】:

    猜你喜欢
    • 2021-03-31
    • 1970-01-01
    • 1970-01-01
    • 2020-05-29
    • 1970-01-01
    • 1970-01-01
    • 2017-05-14
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多