本文部分图片来源于参考中的文章。

定时任务

为什么需要定时任务?定时任务有如下场景:

  • 场景1:某系统凌晨 1 点要进行数据备份
  • 场景2:每 30 s 拉取某平台的最新消息
  • 场景3:购物系统中,用户下单后 12h 未支付自动取消订单
  • 场景4:某游戏准备界面中,21s 内不点击则取消游戏
  • 场景5:某缓存系统,数千万级元素过期清理
  • ….

这些场景往往都要求我们在某个特定的时间去做某个事情,也就是定时或者延时去做某个事情。

  • 定时任务:在指定时间点执行特定的任务,例如每天早上 8 点的闹钟。定时任务可以用来做一些周期性的工作,如数据备份,日志清理,报表生成等
  • 延时任务:一定的延迟时间后执行特定的任务,例如 12h 后,21s 后等。延时任务可以用来做一些异步的工作,如订单取消,推送通知,红包撤回等

尽管二者的适用场景有所区别,但它们的核心思想都是将任务的执行时间安排在未来的某个点上,以达到预期的调度效果。

定时任务实现方案

  • Timer
  • ScheduledExecutorService
  • DelayQueue
  • Spring Task
  • 时间轮

我重点关注于重要且之前不熟悉的定时任务实现方案:时间轮,其他可看 Java 定时任务详解。

时间轮

时间轮由来

时间轮(Time Wheel)是一种高效定时任务实现方案,在论文 Hashed and Hierarchical Timing Wheels中被首次提出。主要是用于操作系统对高性能定时器机制的需求,尤其在处理大量定时器(Timer)对象时,传统方法如最小堆或链表,效率低下(复杂度 O(n) or O(logn))。为了让定时器触发更高效,时间轮作为一种近似 O(1)的延迟任务调度算法被提出。

我简单说一下传统的计时器实现方式。

传统的计时器实现方式

这里可以把计时器理解为:定时任务调度的底层支撑。我这节要说明的就是像Timer/SchduledExecutorService/DelayQueue(后续统称为传统计时器)这些定时任务实现方案下的底层支撑。

传统计时器的实现方式可以分为两类数据结构:

  • 无序列表
  • 有序列表

基于无序列表实现的计时器特点:

  • 任务添加:会将定时任务添加到列表的末尾。时间复杂度O(1)
  • 任务执行:每一次 Tick,遍历任务列表,将<=当前时间的定时任务移出列表,调度执行(放入消费线程池异步执行)。时间复杂度O(n)

该类计时器,当任务数量 n 很大时,O(n) 的复杂度会成为性能瓶颈。尤其是在 Tick 间隔较短的情况下,可能无法在一次 Tick 的间隔内完成对列表的遍历。

Tick:计算机硬件时钟定期发出的中断信号,每一次称为一个 “tick”(时钟滴答),操作系统以此为时间单位来更新系统时间并处理定时事件。我们可以简单理解为每次 tick 后总时间加 1 秒。

基于有序列表实现的计时器特点:

  • 任务添加:列表(一般是优先队列)中的任务按延迟时间升序排序,添加定时任务的时间复杂度一般是 O(log n)。
  • 任务执行:每一次 Tick,循环判断堆顶任务时间是否 <= 当前时间,以 O(1) 的时间移除堆顶任务,调度执行,直到堆顶任务时间>当前时间

可以看到,有序列表计时器,优化了无序列表计时器任务执行情况下的耗时。但是当列表维护的定数任务数量变多,每次用户添加新的定时任务时,延迟会非常高,在需要反复创建大量计时器的场合下,性能不佳。

总结:上述两类计时器都会因为定时任务的数量增加时,从任务添加/任务执行两个角度影响到整体的性能。

在涉及百万级以上的定时任务处理下,就需要更加高效定时任务调度方案,于是像 Kafka、Dubbo、ZooKeeper、Netty、Caffeine 等中间件都使用了时间轮(实现方式不同)作为定时任务处理方案。

计时器的实现与排序算法的关联

为什么我会说到排序算法呢?

这是从可优化角度入手,相比无序列表,我们可以从优化有序列表的排序算法来进一步提高性能。

有序列表(优先队列)采用的实现是推排序,复杂度为 O(nlog n)。除了堆排序,还有像快速排序、归并排序等基于比较的高效排序算法,不过这些排序算法的复杂度都是 O(nlog n),不能满足提高计时器性能的目的。

那还有什么其他的排序算法复杂度在 O(n) or O(log n) 吗?对于 n 个元素之间的排序复杂度最低都是 O(n),不可能 O(log n)。而对于 O(n) 时间复杂度的排序算法是存在的,还不少。

  1. 计数排序,缺陷很大,对于特定算法题可能很好用
  2. 桶排序,对计数排序的优化
  3. 基数排序

上述排序算法都是不基于比较的排序算法,其实严格来说他们的平均复杂度不是 O(n)。

  1. 计数排序 O(n + k),k 是待排序元素的范围
  2. 基数排序 O(d*(n + r)),d 是数字的位数,b 是基数(按十进制排序时 b=10)
  3. 桶排序 O(nlog (n/m)​),m 是桶的个数

上述具体的算法讲解,见参考。

回到主线,我们将有序列表使用的排序算法换为桶 or 基数排序后,相当于降了一个数量级,性能提升非常明显。但这是有前提条件的,要求排序的数据分布情况均匀、数据集范围不算特别广。

举个例子,以 10w 副打乱的扑克牌为例(不算大小王),由于扑克牌只有 A-K 13 种可能,所以只需要 13 个 桶(bucket) 就可以将扑克牌全部收集,并在线性时间复杂度 O(n) 内完成排序。

可是实际情况是这样吗?

每一次 Tick 的时间间隔非常小(纳秒级别),我们使用类似基数排序的思想,使用巨大数量的 1bucket = 1ns 来存储不同过期时间的任务,这在理论上可行,但是我们空间效率却低的恐怖。

首先我们不可能每 1ns 都有定时任务执行,这将导致大量空 bucket。其次现有内存硬件也无法满足 1 纳秒对应 1 个bucket。

但如果能容忍时钟调度的时间不是那么精确,就可以极大减少所需要的 bucket 桶的数量。例如,在之前无序/有序列表中,设置的时间精度是 1s,也就是说每个定时任务我们精确到秒即可。从第 1 ~ $10^9$ns 都放在 1s 对应的 bucket 中,这样相当于空间消耗减少了一亿倍!

时间轮算法就是基于这一特点产生,即一定程度上舍弃调度时间的精确性,参考基数排序的思想,实现在常数时间内创建定时任务,并同时在常数时间内完成定时任务的调度(遍历、放入消费线程池等,不是真正意义上的执行)。

总结:知道了传统计时器实现方式下的优缺点,知道了时间轮算法参考基数排序的思想来提高效率,但是我们并不清楚时间轮具体的实现方式如何,以及还做了那些优化。

时间轮的不同实现方式

时间轮简单来说就是一个环形队列(底层一般基于数组实现,如下图),队列中的每一个元素都可以存放一个定时任务列表(一般是双端链表)。

时间轮中的每个时间格代表了时间轮的基本时间跨度或者说时间精度,假如时间一秒走一个时间格的话,那么这个时间轮的最高精度就是 1 秒(也就是说 3 s 和 3.9s 会在同一个时间格中)。

下图是一个有 12 个时间格的时间轮,转完一圈需要 12s。当我们需要新建一个 3s 后执行的定时任务,只需要将定时任务放在下标为 3 的时间格中即可。当我们需要新建一个 9s 后执行的定时任务,只需要将定时任务放在下标为 9 的时间格中即可。

  • 可以看出来,这里的每一个时间格其实就是一个 bucket
  • 任务添加:通过取模运算的方式以 O(1) 的复杂度将定时任务添加到 bucket 对应的链表上
  • 任务执行:每一次 Tick,以 O(1) 时间移出当前 bucket 中的所有定时任务(head/tail -> null),放入消费线程池中执行

可以看出来,时间轮这种模式,在处理大量定时任务时,相比于使用无序列表或有序列表,能够提供更高的效率,显著降低任务添加和执行造成的性能影响。

不过,上述这个基本的时间轮实现存在一个巨大缺陷!

如上图所示,当我们需要创建一个 13s 后执行的定时任务怎么办呢?

可以看到是 bucket 的数量限制了时间轮所能够支持的最大超时时间,如何解决呢?增加 bucket 数量吗?

答案显然不是,存储 86400秒(1天)后超时的任务就至少需要 86400 bucket,无疑这个空间消耗是巨大的,而且还会出现前文所说过的空 bucket 问题。这里有两种优化方法可以完美解决这一问题。

单层多轮次时间轮

可以让时间轮中的定时任务添加 Round(轮次/圈数) 属性。还是以 12 个 bucket 的时间轮为例,当我们创建一个 13s 后执行的定时任务时,额外计算 round:

  • round = 13 ÷ 12 = 1
  • bucket =13 % 12 = 1

所以 13s 后执行的定时任务处于下标为 1 的 bucket,其 round = 1。

每一次 Tick,遍历链表,令 round 减去 1,移出 round 为 0 的定时任务,让消费线程处理。你可能会有这样的疑惑,这不又是需要遍历列表之后才执行定时任务的吗,不是和无序列表的情况又一样了吗?

其实是不对的,无序列表每次执行定时任务遍历的是所有整体列表集合,这个操作的复杂度是 O(n)。而时间轮只是遍历其中一个 bucket 中的定时任务而已,复杂度 O(k)。

虽然在极端最坏情况下(例如,所有 N 个定时任务碰巧都将在同一个 Tick 槽位到期),时间轮在那个特定的 Tick 处理时可能会遍历一个很长的列表,看起来像 O(n)。但是,从整体和平均情况来看,时间轮将 N 个任务分散到了 W 个槽位中(W 是轮子的槽位数量),每次 Tick 的处理只涉及 N/W (平均)个任务的遍历,远小于 O(n)。

Netty(4.1.119.Final) 是一个高性能的网络应用框架,其内部就使用了 HashedWheelTimer(单层多轮次实现)来高效地管理网络连接的读写超时、任务调度等。

除了增加圈数这种方法之外,论文中提到一种多层时间轮 (类似手表),Kafka(3.8.1) 采用的就是这种方案。

其实 round 属性相当于就是一个逻辑上的多层时间轮,只不过是“隐形”的时间层。

多层时间轮

针对下图的时间轮,我来举个例子便于大家理解。

上图的时间轮,总共有三层,第 1 层的时间精度为 1ms,第 2 层的时间精度为 20ms,第 3 层的时间精度为 400ms。假如我们需要添加一个 350ms 后执行的任务 A(当前时间是 0ms),这个任务会被放在第 2 层(因为第二层的时间跨度为 $20*20=400>350$)的第 $350 ÷ 20 = 17$ 个时间格。

当第一层转了 17 圈之后(第一层转一圈相当于第二层1个时间格),时间过去了 340ms ,第 2 层的指针此时来到第 17 个时间格。此时,第 2 层第 17 个格子的任务会被移动到第 1 层(350 mod 20 = 10,10 < 20)。任务 A 当前是 10ms 之后执行,因此它会被移动到第 1 层的第 10 个时间格子。

这里在层与层之间的移动也叫做时间轮的升降级。参考手表(时、分、秒)来理解就好。

高性能 Java 缓存库 Caffeine 中也是多层时间轮,它一共有五层:

这里还有一些其他问题,当多层时间轮本身的表示范围也不够时的解决方案?

  1. 在最高层(最粗粒度)轮子上增加轮次 Round,将其作为更粗粒度的计数器
  2. 使用单独的长期任务存储结构,存储延迟特别巨大(比如几个月、几年)的任务
  3. 持久化存储(极端超长的延迟任务),一个独立的调度器服务会定期从数据库中查询“即将到期”的任务

总结:时间轮比较适合任务数量比较多的定时任务场景,它的任务写入和执行的时间复杂度都是 0(1)。

上述介绍时间轮的实现方式以“文字+图片”的形式,如果有时间,可以看看下面的“已有实现+源码”的形式深入了解时间轮的工作原理。

源码分析

这个根据已有的成熟实现,分析源码,然后仿照其实现方式自己动手搓一个,这样下来对时间轮原理和实现理解的也会更加透彻些。

代码 demo

单层多轮次时间轮

这个我们通过 Netty 内部实现的时间轮进行分析。

在 pom.xml 中导入 netty-common 模块

0<dependency>
1    <groupId>io.netty</groupId>
2    <artifactId>netty-common</artifactId>
3    <version>4.1.100.Final</version>
4</dependency>

因为我想要了解的是单层次时间轮实现的核心内容,所以我会选择性的查看:

  1. 存储 TimeTask 的数据结构是什么?
  2. TimeTask 的属性参数有哪些(主要关注是否有 Round)
  3. 添加定时任务的方法
  4. 定时任务执行的代码

在依赖导入成功后,进入 io.netty.util.HashedWheelTimer,这是具体实现类。

可以看到 HashedWheelBucket[] wheel,这是时间轮的核心数据结构,使用数组模拟的环形队列(这里存放定时任务的其实还有 timeouts、cancelledTimeouts)。

然后 HashedWheelBucket 就是时间槽的数据结构,如下图:

可以看到 bucket 有个 HashedWheelTimeout 的属性参数,通过命名我们可以看出来这是用于构建 HashedWheelTimeout 链表的头尾指针。这个 HashedWheelTimeout 身份也很明显了,它是定时任务的封装类。(在 bucket 还有诸如 addTimeout/expireTimeouts 等方法)

在 HashedWheelTimeout 类中,TimeTask 是需要被执行的定时任务具体逻辑、remainingRounds 是轮次/圈数属性、next 和 prev 相互构成一个双向链表。

存储 TimeTask 的数据结构和定时任务中的属性参数,这 2 个点基本解决,我们再来看添加定时任务的方法。

 0/**  
 1 * 安排一个新的 TimerTask 在指定的延迟后执行。  
 2 * 如果底层的 AtomicLong 和集合是线程安全的,则此方法通常是线程安全的。  
 3 *  
 4 * @param task  要安排的任务。不能为 null。  
 5 * @param delay 任务执行前的延迟时间。  
 6 * @param unit  延迟的时间单位。不能为 null。  
 7 * @return 表示已安排任务的 Timeout 对象,可用于取消。  
 8 * @throws RejectedExecutionException 如果待处理任务数量超出最大允许限制。  
 9 * @throws NullPointerException 如果 task 或 unit 为 null。  
10 */  
11public Timeout newTimeout(TimerTask task, long delay, TimeUnit unit) {  
12    // 1. 检查参数是否为 null,防止后续出现 NullPointerException。  
13    ObjectUtil.checkNotNull(task, "task");  
14    ObjectUtil.checkNotNull(unit, "unit");  
15    // 2. 原子地增加待处理任务计数器,这个计数器跟踪已经安排但尚未执行或取消的任务数量。  
16    long pendingTimeoutsCount = this.pendingTimeouts.incrementAndGet();  
17    // 3. 检查是否设置了最大待处理任务限制 (maxPendingTimeouts > 0) 以及当前计数是否超出该限制。  
18    if (this.maxPendingTimeouts > 0L && pendingTimeoutsCount > this.maxPendingTimeouts) {  
19        // 4. 如果超出限制,则原子地减少计数器,因为任务被拒绝。  
20        this.pendingTimeouts.decrementAndGet();  
21        // 5. 抛出 RejectedExecutionException 异常,表示任务无法被安排。  
22        throw new RejectedExecutionException("Number of pending timeouts (" + pendingTimeoutsCount + ") is greater than or equal to maximum allowed pending timeouts (" + this.maxPendingTimeouts + ")");  
23    } else {  
24        // 6. 如果未超出限制(或未设置限制),初始化或启动处理任务调度的后台线程/资源。  
25        this.start();  
26        // 7. 计算任务的截止时间 (deadline)。  
27        long deadline = System.nanoTime() + unit.toNanos(delay) - this.startTime;  
28        // 8. 处理截止时间计算中可能出现的溢出或环绕问题。  
29        // 如果延迟是正数,但计算出的截止时间为负数(可能由于 System.nanoTime() 环绕或与非常大的 startTime 交互导致),  
30        if (delay > 0L && deadline < 0L) {  
31            deadline = Long.MAX_VALUE;  
32        }  
33        // 9. 创建 HashedWheelTimeout 对象,表示已安排的任务。  
34        HashedWheelTimeout timeout = new HashedWheelTimeout(this, task, deadline);  
35        // 10. 将新创建的 timeout 对象添加到计时器的内部集合/结构中。  
36        this.timeouts.add(timeout);  
37        return timeout;  
38    }  
39}

总的来说,这段 newTimeout() 方法用于安排一个定时任务在指定的延迟后执行。它首先检查输入参数是否有效,并限制待处理任务的总数。如果检查通过,this.start() 会确保计时器已启动,然后计算任务的相对执行截止时间(处理可能的计算问题)。最后,它创建一个表示该任务的内部对象 HashedWheelTimeout,将其添加到计时器的调度队列中,并返回该对象以便调用者进行管理(如取消/修改)。

这里有几个关键的细节:

  1. 待处理任务总数的限制
  2. this.start();
  3. 将创建的 HashedWheelTimeout 添加到 timeout 队列中而不是直接添加到 wheel

为什么我们需要有这个待处理任务总数的限制(也可以没限制)?

每个待处理的任务都会占用一定的内存资源。如果没有上限,在任务提交速度远大于处理速度时,待处理任务会无限堆积,最终 OOM。这可以作为一种简单的流量控制机制。当计时器处理不过来时,通过拒绝新的任务,可以将压力反馈给任务的提交者,让其知道系统当前繁忙,从而避免雪崩效应。

this.start() 这个方法的作用?

源码分析:

 0public void start() {
 1    // 根据当前状态执行操作 (0:初始, 1:已启动, 2:已停止)
 2    switch (WORKER_STATE_UPDATER.get(this)) {
 3        case 0: 
 4            // 原子尝试从 0 设为 1,成功则启动工作线程 (只执行一次)
 5            if (WORKER_STATE_UPDATER.compareAndSet(this, 0, 1)) {
 6                this.workerThread.start();
 7            }
 8        case 1: 
 9            break;
10        case 2:
11            throw new IllegalStateException("cannot be started once stopped");
12        default:
13            throw new Error("Invalid WorkerState");
14    }
15    // 等待工作线程初始化 startTime,确保后续 deadline 计算准确
16    while(this.startTime == 0L) {
17        try {
18            this.startTimeInitialized.await();
19        } catch (InterruptedException var2) {
20            // 忽略或处理中断
21        }
22    }
23}

这个 start() 方法的主要目的和核心操作就是用来启动 workerThread 这个专门的工作线程。this.workerThread 就是计时器的“引擎”,它在后台运行,负责推动时间轮、检查任务是否到期以及执行到期的任务。start() 方法确保了这个引擎在需要时(通常是第一个任务添加)被启动,并且只启动一次。节约了创建和维护一个线程所需的系统资源。

为什么将创建的定时任务添加到 timeout 队列中而不是直接添加到 wheel?

主要代码在:io.netty.util.HashedWheelTimer.Worker#transferTimeoutsToBuckets

this.timeouts 起到了一个任务提交缓冲区的作用。应用程序的多个线程将新任务快速、安全地提交到这个缓冲区。而计时器的后台工作线程 workerThread 则周期性地(例如,在每个时间轮刻度前进时)从这个缓冲区取出所有待处理任务,然后集中处理并将它们分配到时间轮的正确位置进行调度。这样做是为了提高并发性能、简化实现,并将任务提交与核心的调度逻辑分离开来。

添加定时任务的方法大致就是如此,还有一些没有涉及到,就如上图,workerThread 到 Wheel 的具体“处理”部分。由于篇幅有限,我们继续来看下一个部分:定时任务执行的代码。

Netty 时间轮中所有定时任务的执行逻辑(tick 推进、bucket 链表遍历、round 自减、到期任务执行等等)全部由 Worker 线程负责。主要代码在 Worker.run() 方法中。我们把 run() 方法拆成四个核心阶段来理解:

 0public void run() {  
 1	// 1、启动时间轮起始时间
 2    HashedWheelTimer.this.startTime = System.nanoTime();  
 3    if (HashedWheelTimer.this.startTime == 0L) {  
 4        HashedWheelTimer.this.startTime = 1L;  
 5    }  
 6    HashedWheelTimer.this.startTimeInitialized.countDown();  
 7    int idx;  
 8    HashedWheelBucket bucket;
 9    // 2、主循环:定时调度循环(核心部分)
10	do {
11	    long deadline = this.waitForNextTick();     // 等待下一个 tick 的时刻
12	    if (deadline > 0L) {
13	        idx = (int)(this.tick & mask);         // 当前 tick 对应的槽位
14	        this.processCancelledTasks();          // 清理取消的任务
15	        this.transferTimeoutsToBuckets();      // 将新添加的任务转移到 wheel 槽中
16	        bucket = HashedWheelTimer.this.wheel[idx];
17	        bucket.expireTimeouts(deadline);       // 遍历当前槽,执行任务
18	        ++this.tick;                           // 时间轮推进
19	    }
20	} while (WORKER_STATE_UPDATER.get(...) == 1);
21	// 3、清理 wheel[] 槽中未处理的任务(当 Timer 停止时)
22    HashedWheelBucket[] var5 = HashedWheelTimer.this.wheel;  
23    int var2 = var5.length;  
24    for(idx = 0; idx < var2; ++idx) {  
25        bucket = var5[idx];  
26        bucket.clearTimeouts(this.unprocessedTimeouts);  
27    } 
28    // 4. 清理 timeouts 队列中未被归槽的任务和 cancelledTimeouts 队列中被取消的任务 
29    while(true) {  
30        HashedWheelTimeout timeout = (HashedWheelTimeout)HashedWheelTimer.this.timeouts.poll();  
31        if (timeout == null) {  
32            this.processCancelledTasks();  
33            return;  
34        }  
35        if (!timeout.isCancelled()) {  
36            this.unprocessedTimeouts.add(timeout);  
37        }  
38    }  
39}

具体一点的代码查看具体方法即可,比如说:bucket.expireTimeouts(deadline);

其他的方法:

至此,我认为的四个核心内容都已全部查看/分析完毕。

总的来说,Netty 的单层次哈希时间轮利用了时间轮的结构和独立的后台线程 Worker,辅以中间队列处理并发提交,实现了一个高性能、可扩展的定时任务调度器。

仿照的简单实现。

多层时间轮

这个参考 Kafka 的实现来进行分析。

在 pom.xml 中导入 kafka-server 模块

0<dependency>  
1    <groupId>org.apache.kafka</groupId>  
2    <artifactId>kafka-server</artifactId>  
3</dependency>

对于多层时间轮,关注:

  1. 存储 TimeTask 的数据结构是?
  2. 如何实现多层结构?
  3. 添加定时任务的方法
  4. 定时任务执行的代码(时间轮的升/降级)

依赖导入成功后,进入 org.apache.kafka.server.util.timer.SystemTimer 这是主要的实现类。

可以看到在 SystemTimer 中,DelayQueue 和 TimingWheel 和定时任务存储相关。很明显 TimingWheel 是时间轮的核心数据结构(名字….),我们进入 TimingWheel 类中查看。

在 TimingWheel 中同样实现了一个延迟队列 queue(具体作用看下文),然后就是时间轮的存储槽 buckets。可以看到具体的 bucket 里是定时任务列表 TimerTaskList。

TimerTaskList 封装了具体定时任务的 root 节点,实现了 Delayed 接口,可以放入一个 DelayQueue(参考前文的 queue 的定义)。继续看定时任务实体 TimerTaskEntry。

oh my gah,这又是一层封装类,有 next 和 prev 指针用于构建双向链表,具体定时任务实现逻辑在 TimerTask 中。

可以看到在 TimerTask 中只有 delayMs 延迟时间的定义,没有 Round 的概念。

现在我们对 kafka 的定时任务存储结构基本上了解,和 netty 一样,kafka 的时间轮存储也是数组 + 双向链表的组合使用。

第一个需求解决,继续看下一个需求:“时间轮多层次结构如何实现的”。

回到 TimingWheel 这张图:

这里有几个属性非常重要:

  • tickMs:时间轮的时间刻度
  • interval:代表不同层的时间刻度 $interval = tickMs * wheelSize$
  • overflowWheel:上层时间轮的引用

总的来说,Kafka 的时间轮之所以是“多层次”的,是通过构建由多个 TimingWheel 实例组成的链条来实现的,每个实例代表一个不同时间刻度的时间轮。

具体的代码:

 0/**
 1 * 创建并添加上层时间轮,用于处理延迟时间大于当前时间轮覆盖范围的任务。
 2 * 该方法是线程安全的,通过 synchronized 保证并发环境下不会重复创建 overflowWheel。
 3 */
 4private synchronized void addOverflowWheel() {
 5    if (this.overflowWheel == null) {
 6        this.overflowWheel = new TimingWheel(
 7            this.interval,         // 层次结构的关键:上一层的刻度大小等于下一层的总跨度,确保了时间粒度的逐层变粗
 8            this.wheelSize,        // 槽数量相同
 9            this.currentTimeMs,    // 当前时间,用于校准新时间轮的初始时间
10            this.taskCounter,      // 全局任务数计数器
11            this.queue             // 用于统一任务调度的延迟队列
12        );
13    }
14}

第三个需求:“添加定时任务的方法”。

0// SystemTimer 类
1public void add(TimerTask timerTask) {
2    this.readLock.lock();  
3    try {  
4	    // 主要行为
5        this.addTimerTaskEntry(new TimerTaskEntry(timerTask, timerTask.delayMs + Time.SYSTEM.hiResClockMs()));  
6    } finally {  
7        this.readLock.unlock();  
8    }  
9}
 0// TimingWheel 类
 1public boolean add(TimerTaskEntry timerTaskEntry) {
 2    // 主要功能块 1: 检查任务状态和任务时间是否符合范围
 3    long expiration = timerTaskEntry.expirationMs; // 获取任务的到期时间戳
 4    // 如果任务已取消,或到期时间早于当前轮的下一个刻度,则不添加到此轮。
 5    if (timerTaskEntry.cancelled() || expiration < this.currentTimeMs + this.tickMs) {
 6        return false;
 7    }
 8    // 主要功能块 2: 调度到当前时间轮
 9    // expiration属于[currentTimeMs + tickMs, currentTimeMs + tickMs * wheelSize)
10    else if (expiration < this.currentTimeMs + this.interval) {
11        // 计算桶ID,将任务添加到对应 bucket 中
12
13        long virtualId = expiration / this.tickMs; // 任务对应的虚拟槽编号
14        int bucketId = (int)(virtualId % (long)this.wheelSize); // 将虚拟槽编号映射到实际的槽数组下标
15        TimerTaskList bucket = this.buckets[bucketId];
16        bucket.add(timerTaskEntry);
17		
18		// 如果桶的过期时间被更新(通常是首次设置),则将该桶添加到延迟队列中,以便 Timer 主循环检测并处理到期的桶
19        if (bucket.setExpiration(virtualId * this.tickMs)) {
20            this.queue.offer(bucket);
21        }
22        return true;
23    }
24    // 主要功能块 3: 委托给上一层时间轮 (溢出处理)
25    else {
26        // 确保上一层轮存在,然后将任务添加到上一层轮
27        if (this.overflowWheel == null) {
28            this.addOverflowWheel();
29        }
30        return this.overflowWheel.add(timerTaskEntry); // 递归调用,向上层传递(时间轮升级)
31    }
32}

主要有三部分功能:

  1. 检查任务状态和任务时间是否符合范围
  2. 调度任务到当前时间轮
  3. 委托给上一层时间轮 (溢出处理)

下面有几个问题需要理解。

1、不太理解的是为什么需要设置 bucket 的过期时间?将 bucket 加入延迟队列的目的?(该问题和后续定时任务处理有关)

  • 设置 Bucket 过期时间,标记该时间槽(bucket)在何时到期就绪,这是 DelayQueue 判断何时取出它的依据
  • 将 Bucket 加入延迟队列,是为了 Timer 可以批量处理到期桶里的任务

2、为什么将 bucket 加入 DelayQueue,而不是 TimerTaskEntry ?

  • DelayQueue 的底层是基于优先队列
  • 如果每一个 TimerTaskEntry 都被放入 DelayQueue,DelayQueue 会变得非常庞大。这会导致插入和删除效率降低、内存开销增加(前文对传动计时器-有序队列的效率分析中有过讲解)

3、就功能块 3 而言:“如果任务的到期时间 $>$ 最高层时间轮范围,就会创建更高层次的时间轮”。这样不会导致时间轮数量过多吗?

  • 就功能块 3 而言,理论上不断有到期时间越来越大的任务被添加进来,确实会导致 overflowWheel 层层向上创建,形成无限多的层级
  • 不过 kafka 肯定考虑到这个问题,一个完整的基于分层时间轮的定时器通常会在 SystemTimer 初始化时就确定好:底层时间轮的 tickMs 、 wheelSize 以及时间轮的层数(3、5)。最高层时间轮确定后,其 overflowWheel 引用就会保持为 null,或者指向一个特殊的“无效”时间轮

4、如果定时任务 1 的到期时间是 5s(假设当前时间 0s),第一个时间轮:tickMs = 1s; wheelSize = 10。这时候创建定时任务 2(特殊情况),到期时间是 $10^6$s。根据add()来看,$tickMs = 10^6 ÷ 10 = 10^5$,这样的话就会创建 6 层时间轮。如果后续定时任务的到期时间都集中于 1s ~ 100s 以内,其它高层次时间轮不就浪费了吗?

  • 这是惰性创建(Lazy Creation)高层时间轮机制的一个问题:潜在的浪费。不过惰性创建只在需要时才创建,总比一开始就创建所有层级要好
  • 在问题 3 中指明了 Kafka 会对分层时间轮做层数限制,这保证了一定数量的时间轮是 kafka 可以接受的范围。其次浪费的是轮和桶的结构(相对少量),而不是存储任务本身占用的空间(链表,这个很大)
  • 创建少量额外的层级来支持较长延迟,通常比使用一个拥有海量桶的单层时间轮要高效得多

最后一个需求了:“定时任务执行的代码(时间轮的升/降级)”

任务最终的执行是通过 SystemTimer 中的 ExecutorService taskExecutor 来完成的,ExecutorService 是线程池用于对定时任务异步消费。

有两部分,分别是 SystemTimer 类的 addTimerTaskEntry() 、advanceClock() 方法

 0/**
 1 * 将 TimerTaskEntry 加入时间轮,如果不能加入(表示已过期),并且未被取消,则立即通过线程池 taskExecutor 执行该任务。
 2 * @param timerTaskEntry 定时任务包装对象,包含任务本体与过期时间等信息
 3 */
 4private void addTimerTaskEntry(TimerTaskEntry timerTaskEntry) {
 5    // 尝试将任务加入时间轮
 6    if (!this.timingWheel.add(timerTaskEntry) && !timerTaskEntry.cancelled()) { 
 7        // 任务已经过期且任务未被取消,异步执行
 8        this.taskExecutor.submit(timerTaskEntry.timerTask);
 9    }
10}

这是在添加定时任务时可能会触发的任务执行(只有定时任务在最底层时间轮时才会被执行)。

下面时时间轮推进当前时间、触发到期任务的核心逻辑之一。

 0/**
 1 * 推进时间轮的时钟,并处理到期的 bucket(时间槽)中的任务
 2 *
 3 * @param timeoutMs 最多等待 DelayQueue 中任务的时间(阻塞时间)
 4 * @return 如果成功推进了时间轮,返回 true;否则(队列为空超时)返回 false
 5 */
 6public boolean advanceClock(long timeoutMs) throws InterruptedException {
 7    
 8    // 从 DelayQueue 中拉取一个即将过期的时间槽(TimerTaskList),等待时间为 timeoutMs
 9    TimerTaskList bucket = (TimerTaskList) this.delayQueue.poll(timeoutMs, TimeUnit.MILLISECONDS);
10    // 如果没有拉取到(可能是 timeoutMs 时间内队列没有数据)
11    if (bucket == null) {
12        return false; // 没有任务到期,无法推进时钟
13    } else {
14        // 加写锁,避免多线程同时修改时间轮结构(例如添加/推进时钟)
15        this.writeLock.lock();
16        try {
17            // 处理当前 bucket 及后续可能立即可用的 bucket
18            while (bucket != null) {
19                // 更新 TimingWheel 的 currentTime,确保内部逻辑基于最新时钟
20                this.timingWheel.advanceClock(bucket.getExpiration());
21                // 遍历 bucket 中的任务,将它们重新添加到 TimingWheel,addTimerTaskEntry:如果任务仍未过期,继续加入时间轮;否则立即执行
22                bucket.flush(this::addTimerTaskEntry);
23                // 继续从 DelayQueue 中尝试拉取下一个到期的 bucket(非阻塞)
24                bucket = (TimerTaskList) this.delayQueue.poll();
25            }
26        } finally {
27            // 释放写锁
28            this.writeLock.unlock();
29        }
30        // 至少处理了一个 bucket,说明时钟已推进
31        return true;
32    }
33}

流程:

  • Timer 主循环调用 advanceClock,从 delayQueue 中取出一个到期的 bucket
  • 如果取出的 bucket 不在最底层时间轮,那么这个 bucket 中的任务就需要被“降级”:bucket.flush()
  • 降级逻辑:遍历桶中的所有任务 –> 将每个任务从这个链表中移除 –> 对移除的每个任务调用回调函数(降级) –> 最后重置 bucket 的过期时间
  • 定时任务何时执行:回调函数是 addTimerTaskEntry,因为 TimingWheel 的 currentTime 已经前进了,所以定时任务在 add 的过程因为过期而被消费

总的来说,Kafka 的多层时间轮是一个高性能的定时器实现,旨在高效处理大量具有广泛延迟范围的定时任务;它构建了一个由不同时间粒度层级组成的时间轮结构,将任务组织到各自时间槽的桶中;通过一个延迟队列管理到期的桶,驱动任务随着时间临近从粗粒度层级逐步“降级”到细粒度层级,最终由独立的线程池进行异步执行。

解决了空转问题的多层次时间轮

再次回顾,什么是空转问题?

时间轮当前 tick 推进时,即使当前层的 bucket 为空或无任务到期,也仍然会推进并消耗一次轮询(tick),但没有实际执行任何任务。

还是以 kafka 实现的多层时间轮为例(Netty 的这个版本也解决了空转问题,在waitForNextTick())。

解决了定时器“空转”(Busy-waiting,即线程在没有任务到期时不断地循环检查时间,浪费 CPU)问题的关键在于使用了延迟队列 DelayQueue。

关键代码:

0TimerTaskList bucket = (TimerTaskList)this.delayQueue.poll(timeoutMs, TimeUnit.MILLISECONDS);
  • delayQueue.poll(timeoutMs, TimeUnit.MILLISECONDS) 方法会使当前线程(通常是 Timer 的主循环线程)阻塞住。
  • 它会一直等待,直到:
    • DelayQueue 中有 bucket 的延迟到期
    • 或者指定的 timeoutMs 时间到了

通过这种阻塞式的等待,Timer 的主线程在没有任务到期时会进入休眠状态,而不会浪费 CPU 资源去反复检查当前时间。只有当 DelayQueue 通知有桶到期了,线程才会被唤醒进行处理。这就优雅地解决了定时器常见的空转问题。

代码实现demo

参考