Java线程池ExecutorService shutdown未按预期工作问题求助
解决ExecutorService未等待子任务完成的问题
你的核心问题在于:主线程调用shutdown()后,线程池不再接受新任务;而t0等初始任务在运行过程中提交t4、t5这类子任务时,要么被线程池拒绝(如果shutdown()已执行),要么即便提交成功,主线程的awaitTermination()仅追踪初始任务的完成状态,不会等待子任务结束。
下面是几种可行的解决办法:
方法一:用CompletableFuture关联任务依赖
将初始任务和子任务都包装为CompletableFuture,让初始任务的Future必须等待其所有子任务完成后才标记为完成,主线程只需等待所有初始任务的Future结束即可。
示例代码:
// 全局共享的ExecutorService ExecutorService executor = Executors.newFixedThreadPool(4); // 提交初始任务t0,内部包含子任务t4、t5 CompletableFuture<Void> t0Future = CompletableFuture.runAsync(() -> { System.out.println("t0开始执行"); // 提交子任务t4、t5,并等待它们完成 CompletableFuture<Void> t4Future = CompletableFuture.runAsync(() -> { try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t4执行完成"); }, executor); CompletableFuture<Void> t5Future = CompletableFuture.runAsync(() -> { try { Thread.sleep(1500); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t5执行完成"); }, executor); // 等待t4、t5完成后,t0Future才会标记为完成 CompletableFuture.allOf(t4Future, t5Future).join(); System.out.println("t0执行完成"); }, executor); // 提交其他初始任务t1、t2、t3 CompletableFuture<Void> t1Future = CompletableFuture.runAsync(() -> { try { Thread.sleep(500); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t1执行完成"); }, executor); CompletableFuture<Void> t2Future = CompletableFuture.runAsync(() -> { try { Thread.sleep(800); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t2执行完成"); }, executor); CompletableFuture<Void> t3Future = CompletableFuture.runAsync(() -> { try { Thread.sleep(1200); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t3执行完成"); }, executor); // 主线程等待所有初始任务(含其子任务)完成 CompletableFuture.allOf(t0Future, t1Future, t2Future, t3Future).join(); // 关闭线程池并清理 executor.shutdown(); if (!executor.awaitTermination(60, TimeUnit.MINUTES)) { executor.shutdownNow(); } // 写入结果到文件 writeFile(); closeFile();
方法二:用AtomicInteger跟踪活跃任务数
通过原子计数器记录所有正在运行的任务(包括初始任务和子任务),主线程循环等待计数器归0,确保所有任务执行完毕。
示例代码:
ExecutorService executor = Executors.newFixedThreadPool(4); AtomicInteger activeTasks = new AtomicInteger(0); // 封装任务提交逻辑:提交前计数器+1,任务完成后计数器-1 Runnable wrapTask(Runnable task) { return () -> { activeTasks.incrementAndGet(); try { task.run(); } finally { activeTasks.decrementAndGet(); } }; } // 提交初始任务t0,内部提交子任务时同样用wrapTask包装 executor.submit(wrapTask(() -> { System.out.println("t0开始执行"); executor.submit(wrapTask(() -> { try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t4执行完成"); })); executor.submit(wrapTask(() -> { try { Thread.sleep(1500); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t5执行完成"); })); System.out.println("t0执行完成"); })); // 提交其他初始任务t1、t2、t3 executor.submit(wrapTask(() -> { try { Thread.sleep(500); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t1执行完成"); })); executor.submit(wrapTask(() -> { try { Thread.sleep(800); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t2执行完成"); })); executor.submit(wrapTask(() -> { try { Thread.sleep(1200); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t3执行完成"); })); // 主线程等待所有任务完成 while (activeTasks.get() > 0) { try { Thread.sleep(100); // 短暂休眠避免空转 } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } // 关闭线程池 executor.shutdown(); if (!executor.awaitTermination(60, TimeUnit.MINUTES)) { executor.shutdownNow(); } writeFile(); closeFile();
方法三:改用ForkJoinPool(适合嵌套/递归任务场景)
如果你的任务是递归生成的(比如t0生成t4、t5,t4可能再生成子任务),可以使用ForkJoinPool,它天生支持跟踪嵌套任务的完成状态,调用shutdown()和awaitTermination()会等待所有子任务完成。
示例代码:
ForkJoinPool executor = new ForkJoinPool(4); // 初始任务t0,内部提交子任务时用ForkJoinTask的fork()方法 executor.submit(() -> { System.out.println("t0开始执行"); // 提交子任务t4、t5 ForkJoinTask<?> t4 = ForkJoinPool.commonPool().submit(() -> { try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t4执行完成"); }); ForkJoinTask<?> t5 = ForkJoinPool.commonPool().submit(() -> { try { Thread.sleep(1500); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t5执行完成"); }); // 等待子任务完成 t4.join(); t5.join(); System.out.println("t0执行完成"); }); // 提交其他初始任务 executor.submit(() -> { try { Thread.sleep(500); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("t1执行完成"); }); // 关闭线程池并等待所有任务完成 executor.shutdown(); if (executor.awaitTermination(60, TimeUnit.MINUTES)) { writeFile(); closeFile(); } else { executor.shutdownNow(); }
关键注意点
- 无论用哪种方法,都要确保子任务是在线程池关闭前被提交的,否则会抛出
RejectedExecutionException。 - 如果子任务依赖初始任务的执行结果,必须显式建立任务间的依赖关系(比如CompletableFuture的
join()、ForkJoinTask的join()),否则初始任务会提前结束,主线程会误判任务全部完成。
内容的提问来源于stack exchange,提问作者efoc
相关产品推荐
相关产品推荐

