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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 19:50:26