ExecutorService未等待所有线程执行完成就执行后续逻辑求助
问题分析与修复方案
你的代码出现提前执行后续逻辑的情况,大概率是这几个原因:
- 任务执行过程中抛出异常,导致
countDownLatch.countDown()没被调用,Latch无法归0,10秒超时后直接进入后续逻辑 - 任务处理耗时超过10秒,
await超时后不再等待 executorService.shutdown()只是停止接受新任务,不会等待已提交的任务完成
下面是几种靠谱的修复方式:
方式一:修复CountDownLatch的异常处理
确保不管任务成功还是失败,countDown()都会执行,同时根据实际情况调整超时时间(如果任务确实可能超过10秒,就调大或者去掉超时):
public void doParallelStuff() throws InterruptedException { final int THREAD_COUNT = 5; ExecutorService executorService = Executors.newFixedThreadPool(THREAD_COUNT); Set<String> myIdSet = Set.of("123", "234", "345", "897", "893"); CountDownLatch countDownLatch = new CountDownLatch(myIdSet.size()); myIdSet.forEach(id -> executorService.execute(() -> { try { // retrieve & do processing the entity related to id } finally { // 不管是否异常,都要countDown countDownLatch.countDown(); } })); // 根据实际任务耗时调整超时时间,或者用无参的await()无限等待 if (!countDownLatch.await(30L, TimeUnit.SECONDS)) { // 可以在这里处理超时情况,比如打印日志、中断任务等 System.out.println("部分任务超时未完成"); executorService.shutdownNow(); } else { executorService.shutdown(); } // 现在这里会在所有任务完成(或超时处理后)执行 }
方式二:用ExecutorService的invokeAll方法(更简洁)
invokeAll会自动等待所有提交的Callable任务完成,不需要手动管理Latch:
import java.util.concurrent.Callable; import java.util.concurrent.ExecutionException; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; public void doParallelStuff() throws InterruptedException, ExecutionException { final int THREAD_COUNT = 5; ExecutorService executorService = Executors.newFixedThreadPool(THREAD_COUNT); Set<String> myIdSet = Set.of("123", "234", "345", "897", "893"); // 将每个id的处理逻辑包装成Callable List<Callable<Void>> tasks = myIdSet.stream() .map(id -> (Callable<Void>) () -> { // retrieve & do processing the entity related to id return null; }) .collect(Collectors.toList()); // invokeAll会等待所有任务完成(或超时) executorService.invokeAll(tasks, 30L, TimeUnit.SECONDS); executorService.shutdown(); // 后续逻辑在这里执行 }
方式三:shutdown + awaitTermination组合
先调用shutdown()停止接受新任务,然后用awaitTermination等待所有已提交任务完成:
public void doParallelStuff() throws InterruptedException { final int THREAD_COUNT = 5; ExecutorService executorService = Executors.newFixedThreadPool(THREAD_COUNT); Set<String> myIdSet = Set.of("123", "234", "345", "897", "893"); myIdSet.forEach(id -> executorService.execute(() -> { // retrieve & do processing the entity related to id })); executorService.shutdown(); // 停止接受新任务 // 等待所有任务完成,设置足够长的超时时间 if (!executorService.awaitTermination(30L, TimeUnit.SECONDS)) { // 超时后强制中断未完成的任务 executorService.shutdownNow(); } // 后续逻辑在这里执行 }
额外提醒:如果你的任务是IO密集型,20k+任务用固定5线程的池可能效率很低,可以考虑调整线程池大小(比如用newCachedThreadPool或者根据CPU核心数设置合理的线程数)。
内容的提问来源于stack exchange,提问作者CodeIntro
相关产品推荐
相关产品推荐

