DelayQueue 是 Java 并发包中基于优先级队列实现的无界阻塞延迟队列,仅允许延迟时间到期的元素被取出,核心应用于定时任务调度(如订单超时取消、缓存自动清理)。其底层通过 "优先级队列+锁+条件变量" 实现线程安全的延迟排序与阻塞唤醒机制。
一、核心源码解析
1. 类结构与核心属性
DelayQueue 实现 BlockingQueue 接口,内部组合以下关键组件:
public class DelayQueue<E extends Delayed> extends AbstractQueue<E> implements BlockingQueue<E> {
private final transient ReentrantLock lock = new ReentrantLock(); // 控制并发的独占锁
private final PriorityQueue<E> q = new PriorityQueue<>(); // 按延迟时间排序的优先级队列(小顶堆)
private Thread leader = null; // 等待队首元素的线程(Leader-Follower模式)
private final Condition available = lock.newCondition(); // 条件变量:通知元素可用
}
- 优先级队列(
PriorityQueue):按元素延迟时间升序排序,确保堆顶为最早到期元素。 - Leader-Follower模式:通过
leader标记当前等待队首元素的线程,避免多线程无效等待,提升效率。
2. 入队操作(offer 方法)
添加元素时,通过锁保证线程安全,若新元素成为堆顶(最早到期),则唤醒等待线程:
public boolean offer(E e) {
final ReentrantLock lock = this.lock;
lock.lock(); // 获取锁
try {
q.offer(e); // 加入优先级队列(按延迟时间排序)
if (q.peek() == e) { // 若新元素成为堆顶
leader = null; // 重置leader线程
available.signal(); // 唤醒等待的消费线程
}
return true; // 无界队列,总是成功
} finally {
lock.unlock();
}
}
- 写时复制与排序:元素入队时触发优先级队列的堆排序,确保延迟最小的元素始终在队首。
- 无界特性:队列容量无上限(受内存限制),
add/put方法均直接调用offer,不会阻塞。
3. 出队操作(take 方法)
阻塞获取到期元素,未到期则等待至超时或被唤醒:
public E take() throws InterruptedException {
final ReentrantLock lock = this.lock;
lock.lockInterruptibly(); // 可中断锁
try {
for (;;) {
E first = q.peek(); // 获取堆顶元素
if (first == null) {
available.await(); // 队列为空,阻塞等待
} else {
long delay = first.getDelay(TimeUnit.NANOSECONDS); // 剩余延迟时间
if (delay <= 0) { // 已到期,直接出队
return q.poll();
}
first = null; // 释放引用,帮助GC
if (leader != null) {
available.await(); // 已有leader线程,当前线程进入等待
} else {
Thread thisThread = Thread.currentThread();
leader = thisThread; // 标记当前线程为leader
try {
available.awaitNanos(delay); // 等待剩余延迟时间
} finally {
if (leader == thisThread) {
leader = null; // 唤醒后重置leader
}
}
}
}
}
} finally {
if (leader == null && q.peek() != null) {
available.signal(); // 唤醒其他等待线程
}
lock.unlock();
}
}
- 阻塞逻辑:若堆顶元素未到期,
leader线程等待精确的延迟时间,其他线程则无条件等待,避免 CPU 空转。 - 条件变量唤醒:元素入队或到期时,通过
available.signal()唤醒等待线程,确保及时处理新到期任务。
二、元素要求:Delayed 接口
DelayQueue 元素必须实现 Delayed 接口,定义延迟时间与排序规则:
public interface Delayed extends Comparable<Delayed> {
long getDelay(TimeUnit unit); // 返回剩余延迟时间(<=0 表示到期)
int compareTo(Delayed o); // 比较延迟时间,用于优先级排序
}
示例实现(订单超时任务):
public class OrderTask implements Delayed {
private final String orderId;
private final long expireTime; // 到期时间戳(毫秒)
public OrderTask(String orderId, long delayMs) {
this.orderId = orderId;
this.expireTime = System.currentTimeMillis() + delayMs;
}
@Override
public long getDelay(TimeUnit unit) {
return unit.convert(expireTime - System.currentTimeMillis(), TimeUnit.MILLISECONDS);
}
@Override
public int compareTo(Delayed o) {
return Long.compare(this.expireTime, ((OrderTask) o).expireTime); // 按到期时间升序排序
}
}
getDelay:计算剩余延迟时间,需随时间递减至 <=0。compareTo:确保优先级队列按延迟时间排序,避免因排序错误导致任务无法到期。
三、面试考察点
1. 底层实现原理
- 核心组件:优先级队列(排序)+ 锁(线程安全)+ 条件变量(阻塞唤醒)。
- Leader-Follower模式:通过
leader线程减少无效等待,对比普通wait/notify机制的优势。
2. 与 PriorityQueue 的区别
- 线程安全:
DelayQueue是阻塞队列,通过锁实现线程安全;PriorityQueue非线程安全。 - 延迟特性:
DelayQueue仅允许到期元素出队,PriorityQueue无时间限制。
3. 应用场景与局限性
- 典型场景:订单超时取消、缓存自动清理、定时重试机制。
- 局限性:
- 无界队列:大量任务堆积可能导致 OOM(需结合容量控制)。
- 精度问题:依赖系统时间,若时间回拨可能导致任务提前执行;不支持微秒级精度。
- 不支持删除:无法高效移除未到期元素(需遍历队列,时间复杂度 O(n))。
4. 手写延迟任务示例
需完整实现 Delayed 接口,并结合 DelayQueue 模拟任务调度:
public class DelayQueueDemo {
public static void main(String[] args) throws InterruptedException {
DelayQueue<OrderTask> queue = new DelayQueue<>();
// 添加任务:延迟 3s、1s、2s
queue.put(new OrderTask("ORDER_001", 3000));
queue.put(new OrderTask("ORDER_002", 1000));
queue.put(new OrderTask("ORDER_003", 2000));
// 消费任务(单独线程)
new Thread(() -> {
while (true) {
try {
OrderTask task = queue.take(); // 阻塞至任务到期
System.out.println("取消订单: " + task.orderId);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}).start();
}
}
// 输出顺序:ORDER_002(1s)→ ORDER_003(2s)→ ORDER_001(3s)
四、总结
DelayQueue 是延迟任务调度的轻量级实现,通过优先级队列实现按时间排序,借助锁与条件变量保证线程安全与高效唤醒。其核心优势在于无锁读操作与低CPU消耗,但需注意无界队列的内存风险与时间精度限制。面试中需重点掌握其实现原理、Delayed 接口设计及与其他并发队列的差异。
思考:若需实现支持任务取消的延迟队列,如何基于 DelayQueue 扩展?(提示:通过 ConcurrentHashMap 记录任务ID,遍历队列删除时校验)。