ScheduledThreadPoolExecutor

ScheduledThreadPoolExecutor 源码分析

定时/周期任务线程池,Timer 的替代品(Timer 单线程、任务抛异常整个定时器死掉)。
核心思路:继承 ThreadPoolExecutor,把工作队列换成按触发时间排序的二叉小顶堆 DelayedWorkQueue,任务包装成带 time/period 的 ScheduledFutureTask,周期任务跑完一轮就算好下次时间重新入队。

// 基于 JDK 21 (本地 java.base 源), java.util.concurrent.ScheduledThreadPoolExecutor
public class ScheduledThreadPoolExecutor
        extends ThreadPoolExecutor
        implements ScheduledExecutorService {

    // 构造:maximumPoolSize 写死 Integer.MAX_VALUE,但队列无界,永远不会创建超过 core 的线程
    // ——所以 max 参数在这个池里是摆设,线程数只由 corePoolSize 决定
    public ScheduledThreadPoolExecutor(int corePoolSize) {
        super(corePoolSize, Integer.MAX_VALUE,
              DEFAULT_KEEPALIVE_MILLIS, MILLISECONDS,   // keepAlive=10ms,只为 corePoolSize=0 的病态用法兜底
              new DelayedWorkQueue());
    }

    private class ScheduledFutureTask<V>
            extends FutureTask<V> implements RunnableScheduledFuture<V> {

        /** Sequence number to break ties FIFO */
        private final long sequenceNumber;     // time 相同时按提交顺序,保证堆排序的全序

        /** The nanoTime-based time when the task is enabled to execute. */
        private volatile long time;            // 绝对触发时刻(nanoTime 基准),不是剩余延迟

        // period 一个字段编码三种任务:>0 固定速率(fixedRate),<0 固定间隔(fixedDelay),=0 一次性
        private final long period;

        /** The actual task to be re-enqueued by reExecutePeriodic */
        RunnableScheduledFuture<V> outerTask = this;    // 若被 decorateTask 包装过,重新入队的是包装后的对象

        // 记住自己在堆数组里的下标 -> cancel 时 O(log n) 删除,不用 O(n) 线性找
        int heapIndex;

        // ...(构造器省略:就是给上面几个字段赋值)

        public long getDelay(TimeUnit unit) {
            return unit.convert(time - System.nanoTime(), NANOSECONDS);
        }

        public int compareTo(Delayed other) {   // 堆的排序依据:先比触发时刻,再比序号
            if (other == this) // compare zero if same object
                return 0;
            if (other instanceof ScheduledFutureTask) {
                ScheduledFutureTask<?> x = (ScheduledFutureTask<?>)other;
                long diff = time - x.time;      // 用差值比较而非 <,防 nanoTime 环绕出错
                if (diff < 0)
                    return -1;
                else if (diff > 0)
                    return 1;
                else if (sequenceNumber < x.sequenceNumber)
                    return -1;
                else
                    return 1;
            }
            long diff = getDelay(NANOSECONDS) - other.getDelay(NANOSECONDS);
            return (diff < 0) ? -1 : (diff > 0) ? 1 : 0;
        }

        public boolean isPeriodic() {
            return period != 0;
        }

        // fixedRate 和 fixedDelay 的分岔点就在这四行
        private void setNextRunTime() {
            long p = period;
            if (p > 0)
                time += p;                  // 固定速率:上次"应该"触发的时刻 + 周期,与实际执行耗时无关
            else
                time = triggerTime(-p);     // 固定间隔:now + 间隔,从"本轮跑完"起算
        }

        /**
         * Overrides FutureTask version so as to reset/requeue if periodic.
         */
        public void run() {
            if (!canRunInCurrentRunState(this))
                cancel(false);
            else if (!isPeriodic())
                super.run();                    // 一次性任务:正常 FutureTask 流程,跑完进终态
            else if (super.runAndReset()) {     // 周期任务:跑完把 state 复位回 NEW,可再跑
                setNextRunTime();               // 关键:抛异常时 runAndReset 返回 false,
                reExecutePeriodic(outerTask);   // 这两行不执行 -> 周期任务静默停止!
            }
        }
    }

    private void delayedExecute(RunnableScheduledFuture<?> task) {
        if (isShutdown())
            reject(task);
        else {
            super.getQueue().add(task);     // 先入队再保证有线程,和 TPE 的"先开线程"相反:任务还没到点,开了也白等
            if (!canRunInCurrentRunState(task) && remove(task))
                task.cancel(false);
            else
                ensurePrestart();           // 核心线程不足 core 就补一个(不带首任务启动)
        }
    }

    // 周期任务的重新入队:与 delayedExecute 同思路,但 shutdown 时静默丢弃而非 reject
    void reExecutePeriodic(RunnableScheduledFuture<?> task) {
        if (canRunInCurrentRunState(task)) {
            super.getQueue().add(task);
            if (canRunInCurrentRunState(task) || !remove(task)) {   // 入队后二次检查,防 shutdown 竞态
                ensurePrestart();
                return;
            }
        }
        task.cancel(false);
    }

    // 延迟 -> 绝对触发时刻。clamp 到 MAX_NANOS(约146年)防溢出
    long triggerTime(long delay) {
        return System.nanoTime() + Math.min(delay, MAX_NANOS);
    }

    public ScheduledFuture<?> scheduleAtFixedRate(Runnable command,
                                                  long initialDelay,
                                                  long period,
                                                  TimeUnit unit) {
        if (command == null || unit == null)
            throw new NullPointerException();
        if (period <= 0L)
            throw new IllegalArgumentException();
        ScheduledFutureTask<Void> sft =
            new ScheduledFutureTask<Void>(command,
                                          null,
                                          triggerTime(initialDelay, unit),
                                          unit.toNanos(period),     // 正数 -> fixedRate
                                          sequencer.getAndIncrement());
        RunnableScheduledFuture<Void> t = decorateTask(command, sft);
        sft.outerTask = t;
        delayedExecute(t);
        return t;
    }

    public ScheduledFuture<?> scheduleWithFixedDelay(Runnable command,
                                                     long initialDelay,
                                                     long delay,
                                                     TimeUnit unit) {
        if (command == null || unit == null)
            throw new NullPointerException();
        if (delay <= 0L)
            throw new IllegalArgumentException();
        ScheduledFutureTask<Void> sft =
            new ScheduledFutureTask<Void>(command,
                                          null,
                                          triggerTime(initialDelay, unit),
                                          -unit.toNanos(delay),     // 与 fixedRate 在源码上的唯一差别:period 取负
                                          sequencer.getAndIncrement());
        RunnableScheduledFuture<Void> t = decorateTask(command, sft);
        sft.outerTask = t;
        delayedExecute(t);
        return t;
    }
}
// 基于 JDK 21 (本地 java.base 源), ScheduledThreadPoolExecutor.DelayedWorkQueue
// 二叉小顶堆(数组存储) + ReentrantLock + Condition,堆顶永远是最先到期的任务
static class DelayedWorkQueue extends AbstractQueue<Runnable>
    implements BlockingQueue<Runnable> {

    private static final int INITIAL_CAPACITY = 16;
    private RunnableScheduledFuture<?>[] queue =
        new RunnableScheduledFuture<?>[INITIAL_CAPACITY];
    private final ReentrantLock lock = new ReentrantLock();
    private int size;

    // Leader-Follower:只有 leader 线程限时等待(awaitNanos 到堆顶到期),
    // 其余线程无限期 await。避免所有 worker 都定时醒来抢同一个堆顶(惊群+空转)
    private Thread leader;

    private final Condition available = lock.newCondition();

    private static void setIndex(RunnableScheduledFuture<?> f, int idx) {
        if (f instanceof ScheduledFutureTask)
            ((ScheduledFutureTask)f).heapIndex = idx;   // 堆元素每次挪动都同步维护下标
    }

    private void siftUp(int k, RunnableScheduledFuture<?> key) {
        while (k > 0) {
            int parent = (k - 1) >>> 1;         // 完全二叉树用数组下标定位父子,不需要指针
            RunnableScheduledFuture<?> e = queue[parent];
            if (key.compareTo(e) >= 0)
                break;
            queue[k] = e;
            setIndex(e, k);
            k = parent;
        }
        queue[k] = key;
        setIndex(key, k);
    }

    private void siftDown(int k, RunnableScheduledFuture<?> key) {
        int half = size >>> 1;
        while (k < half) {
            int child = (k << 1) + 1;
            RunnableScheduledFuture<?> c = queue[child];
            int right = child + 1;
            if (right < size && c.compareTo(queue[right]) > 0)
                c = queue[child = right];       // 取两个孩子中更早到期的
            if (key.compareTo(c) <= 0)
                break;
            queue[k] = c;
            setIndex(c, k);
            k = child;
        }
        queue[k] = key;
        setIndex(key, k);
    }

    // ...(grow 1.5 倍扩容、indexOf、remove 等省略)

    public boolean offer(Runnable x) {
        if (x == null)
            throw new NullPointerException();
        RunnableScheduledFuture<?> e = (RunnableScheduledFuture<?>)x;
        final ReentrantLock lock = this.lock;
        lock.lock();
        try {
            int i = size;
            if (i >= queue.length)
                grow();                     // 无界:只扩容,永不拒绝 -> offer 永远返回 true
            size = i + 1;
            if (i == 0) {
                queue[0] = e;
                setIndex(e, 0);
            } else {
                siftUp(i, e);
            }
            if (queue[0] == e) {            // 新任务成了堆顶(更早到期)
                leader = null;              // 旧 leader 等的期限作废,废黜它
                available.signal();         // 叫醒一个线程去竞争新 leader
            }
        } finally {
            lock.unlock();
        }
        return true;
    }

    // 出队收尾:末尾元素补到堆顶再下沉,被取走的任务 heapIndex 置 -1
    private RunnableScheduledFuture<?> finishPoll(RunnableScheduledFuture<?> f) {
        int s = --size;
        RunnableScheduledFuture<?> x = queue[s];
        queue[s] = null;
        if (s != 0)
            siftDown(0, x);
        setIndex(f, -1);
        return f;
    }

    public RunnableScheduledFuture<?> take() throws InterruptedException {
        final ReentrantLock lock = this.lock;
        lock.lockInterruptibly();
        try {
            for (;;) {
                RunnableScheduledFuture<?> first = queue[0];
                if (first == null)
                    available.await();      // 空队列:无限等,等 offer 来 signal
                else {
                    long delay = first.getDelay(NANOSECONDS);
                    if (delay <= 0L)
                        return finishPoll(first);   // 到期了,真正出队
                    first = null; // don't retain ref while waiting     // 不然任务被取消后也无法被 GC
                    if (leader != null)
                        available.await();  // 已有 leader 在掐点,自己当 follower 无限等
                    else {
                        Thread thisThread = Thread.currentThread();
                        leader = thisThread;
                        try {
                            available.awaitNanos(delay);    // leader 只睡到堆顶到期
                        } finally {
                            if (leader == thisThread)
                                leader = null;      // 醒来先卸任,回到 for 循环重新竞争出队
                        }
                    }
                }
            }
        } finally {
            if (leader == null && queue[0] != null)
                available.signal();     // 自己拿到任务走人,且没人当 leader -> 传一棒,叫醒下一个
            lock.unlock();
        }
    }
}

原理串讲

scheduleAtFixedRate(task, 0, 1, SECONDS) 走一遍完整链路。
入口把 Runnable 包成 ScheduledFutureTask:time = triggerTime(initialDelay, unit) 算出绝对触发时刻,period = unit.toNanos(period) 存正数。
然后 delayedExecute 把任务直接 add 进 DelayedWorkQueue 再 ensurePrestart() 补核心线程——顺序与 ThreadPoolExecutor.execute 的”先开线程带任务跑”相反,因为任务还没到点,必须先躺进堆里等。
这也是为什么 STPE 只能用自家队列:worker 的 getTask 从队列拿任务,延迟语义完全靠队列的 take/poll 卡时间实现,换普通队列定时功能就没了。

worker 线程最终阻塞在 DelayedWorkQueue.take()。这里是 Leader-Follower 变体:堆顶未到期时,只允许一个线程(leader)awaitNanos(delay) 掐点睡,其余全部无限期 await()

代码块JAVA · 2 行收起展开
为什么这么设计?如果所有空闲 worker 都对同一个堆顶 awaitNanos,到期瞬间集体惊醒却只有一个能抢到任务,其余白白经历一次唤醒-抢锁-再睡的空转;leader 机制把定时唤醒收敛到一个线程,其他线程只在"有新堆顶""接棒"时被精确 signal。
配套的两个细节:`offer` 发现新任务成为堆顶时把 leader 置 nullsignal(旧 leader 等的期限已过时,必须有人按新期限重新等);take 成功拿到任务离开前,若无 leader 且堆非空就 signal 传棒,保证总有人在掐点。

任务到期被 finishPoll 取出后,worker 调 ScheduledFutureTask.run()
周期任务走 runAndReset():执行用户代码但不设置结果、不进终态,成功则把 FutureTask 状态复位回 NEW。
返回 true 才执行 setNextRunTime() + reExecutePeriodic(outerTask)——同一个任务对象改个 time 字段重新入队,全程零新对象分配。
setNextRunTime 就是两种周期语义的全部实现:fixedRate 用 time += p(基于上次应触发时刻,执行耗时不影响节奏,但耗时超过周期时下一轮立即触发,只会顺延不会并发);fixedDelay 用 time = triggerTime(-p)(基于当前时刻,即”跑完歇 delay 再跑”)。
而回到 scheduleAtFixedRate/scheduleWithFixedDelay 两个方法本身,逐行对比会发现唯一差别就是 unit.toNanos(period)-unit.toNanos(delay) 的一个负号——用 period 的符号位当类型标记,省掉一个 boolean 字段。

坑在 runAndReset 的返回值:用户代码抛出未捕获异常时,FutureTask 把异常记进 outcome 并进入 EXCEPTIONAL 终态,runAndReset 返回 false,于是 setNextRunTime 和 reExecutePeriodic 都不会执行——周期任务从此静默停止,不打日志、不报警,异常只存在返回的 Future 里,而周期任务几乎没人去 get()(一 get 还会永久阻塞到异常发生为止)。
所以工程上周期任务体必须自己 try-catch 全包,或者包一层带日志的 wrapper。

为什么取消要靠 heapIndex?堆结构只保证父子有序,remove(task) 若靠 equals 线性扫描是 O(n);每个 ScheduledFutureTask 随 siftUp/siftDown 实时记录自己的数组下标,取消时直接定位、末位补洞再单次 sift,降到 O(log n)。
定时任务场景里大量任务是”设置了但会被取消”(如超时保护),这个优化直接决定取消风暴下的吞吐。

设计取舍

  • period 的符号编码任务类型(>0 rate / <0 delay / =0 一次性),一个 long 干了枚举+字段两件事,代价是可读性差。
  • 队列无界 + maximumPoolSize 无意义:STPE 永远只有 core 个线程;任务堆积不会触发拒绝策略,只会内存涨。
  • fixedRate 不会并发执行同一任务:单个任务耗时超过 period 只会推迟下一轮(run 完才 reExecutePeriodic 入队),不存在两轮重叠。
  • time 用 System.nanoTime 相对时钟,不受系统时间回拨影响;代价是无法表达”每天凌晨 3 点”这种日历语义,那是 Quartz/cron 的活。
  • 周期任务异常即静默死亡:这是 runAndReset 返回值语义的副作用,不是 bug,但必须在业务代码里自行兜底。

Executors 工厂方法:为什么建议手建 ThreadPoolExecutor

// 基于 JDK 21 (本地 java.base 源), java.util.concurrent.Executors
public class Executors {

    public static ExecutorService newFixedThreadPool(int nThreads) {
        return new ThreadPoolExecutor(nThreads, nThreads,
                                      0L, TimeUnit.MILLISECONDS,
                                      new LinkedBlockingQueue<Runnable>());  // 无参构造 = 容量 Integer.MAX_VALUE
    }                                                                        // 任务堆积无上限 -> OOM 风险

    public static ExecutorService newSingleThreadExecutor(ThreadFactory threadFactory) {
        return new AutoShutdownDelegatedExecutorService                     // 包装一层,屏蔽 setCorePoolSize 等调参入口
            (new ThreadPoolExecutor(1, 1,
                                    0L, TimeUnit.MILLISECONDS,
                                    new LinkedBlockingQueue<Runnable>(),    // 同样是无界队列
                                    threadFactory));
    }

    public static ExecutorService newCachedThreadPool() {
        return new ThreadPoolExecutor(0, Integer.MAX_VALUE,                 // 线程数无上限
                                      60L, TimeUnit.SECONDS,
                                      new SynchronousQueue<Runnable>());    // 零容量队列:没有空闲线程就立刻开新线程
    }                                                                       // 突发流量 -> 线程数爆炸

    public static ScheduledExecutorService newScheduledThreadPool(int corePoolSize) {
        return new ScheduledThreadPoolExecutor(corePoolSize);   // 即上文:DelayedWorkQueue 无界 + max 形同虚设
    }

    public static ScheduledExecutorService newSingleThreadScheduledExecutor() {
        return new DelegatedScheduledExecutorService
            (new ScheduledThreadPoolExecutor(1));
    }
}
工厂方法全是 ThreadPoolExecutor(或其子类)的参数预设,问题就出在预设值把风险藏起来了:newFixedThreadPool / newSingleThreadExecutor 用无界 LinkedBlockingQueue,消费速度跟不上时任务无限堆积直到 OOM,且队列永远塞不满导致 maximumPoolSize 和拒绝策略形同虚设;
newCachedThreadPool 反过来,SynchronousQueue 不存任务,max 是 Integer.MAX_VALUE,高峰期直接创建海量线程把内存和调度打穿;

newScheduledThreadPool 同样是无界堆。
共同点:要么队列无界、要么线程数无界,两个”无界”总占其一,而拒绝策略这道最后防线永远触发不了。
手建 ThreadPoolExecutor 就是强迫你把 core/max/队列容量/拒绝策略/ThreadFactory(线程命名,排查问题必备)五件事显式写出来,让资源上限和过载行为在代码评审时可见——这也是阿里 Java 规约强制手建的原因。

延伸阅读