【发布时间】: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’scancel:“参数:mayInterruptIfRunning- 此值在此实现中无效,因为中断不用于控制处理。”也就是说cancel(true)还是cancel(false)都无所谓,CompletableFuture根本不支持打断。
标签: java multithreading concurrency completable-future