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()。
代码块收起展开
为什么这么设计?如果所有空闲 worker 都对同一个堆顶 awaitNanos,到期瞬间集体惊醒却只有一个能抢到任务,其余白白经历一次唤醒-抢锁-再睡的空转;leader 机制把定时唤醒收敛到一个线程,其他线程只在"有新堆顶"或"接棒"时被精确 signal。
配套的两个细节:`offer` 发现新任务成为堆顶时把 leader 置 null 并 signal(旧 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 规约强制手建的原因。