如何等待ScheduledExecutorService所有任务(含后续新增任务)执行完成
实现方案
核心逻辑是跟踪全链路任务的执行状态,全程不主动关闭线程池,既允许任务执行过程中提交新的后续任务,又能在所有任务(含递归提交的关联任务)全部执行完成时返回,同时支持超时阈值避免永久阻塞。
推荐方案:基于Phaser实现动态任务计数
Phaser原生支持参与者动态注册、注销,完美适配任务执行中动态新增子任务的场景,不需要手动处理复杂的状态同步、虚假唤醒问题。
实现步骤
- 初始化Phaser时先注册1个主线程参与者,避免未提交任务时Phaser提前进入终止状态
- 统一封装线程池的任务提交入口:所有任务(包括初始任务、任务执行中提交的后续任务)提交前,先向Phaser注册1个参与者
- 对提交的原始任务做包装:任务执行逻辑放在try块中,无论任务正常执行完成、抛出异常、被中断,都在finally块中调用方法注销当前任务对应的参与者
- 主线程等待时,先注销自身的初始参与者,调用Phaser的带超时等待方法,等所有任务参与者都注销(即所有任务执行完成)后返回
参考代码
import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicBoolean; public class RecursiveTaskScheduler { private final ScheduledExecutorService scheduler; private final Phaser taskTracker = new Phaser(1); private final AtomicBoolean shutdownFlag = new AtomicBoolean(false); public RecursiveTaskScheduler(int corePoolSize) { this.scheduler = Executors.newScheduledThreadPool(corePoolSize); } // 统一调度入口,所有任务必须通过该方法提交 public ScheduledFuture<?> schedule(Runnable task, long delay, TimeUnit timeUnit) { if (shutdownFlag.get()) { throw new RejectedExecutionException("调度器已关闭"); } // 新任务提交前注册计数 taskTracker.register(); // 包装原始任务,执行完成后自动注销计数 Runnable wrappedTask = () -> { try { task.run(); } finally { taskTracker.arriveAndDeregister(); } }; return scheduler.schedule(wrappedTask, delay, timeUnit); } /** * 等待所有任务执行完成 * @param timeout 最大等待超时时间 * @param timeUnit 时间单位 * @return true为所有任务正常执行完成,false为等待超时 */ public boolean waitForAllTasks(long timeout, TimeUnit timeUnit) throws InterruptedException { // 主线程退出初始参与者计数,等待剩余任务全部完成 taskTracker.arriveAndDeregister(); try { return taskTracker.awaitAdvanceInterruptibly(taskTracker.getPhase(), timeout, timeUnit); } finally { // 等待结束后重新注册主线程,支持后续再次调用waitForAllTasks taskTracker.register(); } } // 所有业务逻辑执行完成后,手动调用关闭线程池 public void shutdown() { if (shutdownFlag.compareAndSet(false, true)) { scheduler.shutdown(); } } }
使用注意事项
- 所有任务(包括任务逻辑内部提交的后续任务)必须走封装后的
schedule方法提交,不能直接调用原生ScheduledExecutorService的提交接口,否则会出现计数遗漏,导致等待逻辑提前返回 - 任务包装逻辑中必须用finally块做计数注销,避免任务抛异常、被取消时出现计数泄漏,导致主线程永久阻塞
- 超时参数根据业务容忍度设置即可,出现任务死循环、无限递归提交新任务的异常场景时,等待逻辑会在超时后直接返回,不会永久卡死
不推荐的实现方式
以下方案存在边界问题,不要使用:
- 基于
ThreadPoolExecutor.getActiveCount()+工作队列大小判断任务是否完成:这两个API返回的是近似值,且任务执行中提交新任务时存在临界区判断错误,比如活跃线程数为0、队列为空的瞬间,任务刚好提交新任务,会导致等待逻辑误判所有任务已完成 - 提前调用
shutdown()+awaitTermination():调用shutdown后线程池会拒绝新提交的任务,不符合任务执行中提交后续任务的需求 - 手动维护普通原子计数器+wait/notify:需要自行处理虚假唤醒、超时、多线程状态同步问题,代码复杂度高,容易出并发bug
内容的提问来源于stack exchange,提问作者Hk0
相关产品推荐
相关产品推荐

