CC 咖啡猫的工作空间 Coding Space

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,遍历队列删除时校验)。