本文目录导读:

我来为你提供几个Java实现延迟队列的完整案例,从简单到复杂。
基于 JDK DelayQueue 的实现
延迟消息实体类
import java.util.concurrent.Delayed;
import java.util.concurrent.TimeUnit;
/**
* 延迟消息实体
*/
public class DelayMessage implements Delayed {
private String messageId;
private String content;
private long delayTime; // 延迟时间(毫秒)
private long expireTime; // 过期时间(毫秒)
public DelayMessage(String messageId, String content, long delayTime) {
this.messageId = messageId;
this.content = content;
this.delayTime = delayTime;
this.expireTime = System.currentTimeMillis() + delayTime;
}
@Override
public long getDelay(TimeUnit unit) {
long diff = expireTime - System.currentTimeMillis();
return unit.convert(diff, TimeUnit.MILLISECONDS);
}
@Override
public int compareTo(Delayed o) {
if (this.expireTime < ((DelayMessage) o).expireTime) {
return -1;
} else if (this.expireTime > ((DelayMessage) o).expireTime) {
return 1;
}
return 0;
}
// getters and setters
public String getMessageId() {
return messageId;
}
public void setMessageId(String messageId) {
this.messageId = messageId;
}
public String getContent() {
return content;
}
public void setContent(String content) {
this.content = content;
}
public long getDelayTime() {
return delayTime;
}
public void setDelayTime(long delayTime) {
this.delayTime = delayTime;
}
public long getExpireTime() {
return expireTime;
}
public void setExpireTime(long expireTime) {
this.expireTime = expireTime;
}
@Override
public String toString() {
return "DelayMessage{" +
"messageId='" + messageId + '\'' +
", content='" + content + '\'' +
", delayTime=" + delayTime +
", expireTime=" + expireTime +
'}';
}
}
延迟队列服务类
import java.util.concurrent.DelayQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
/**
* 延迟队列管理器
*/
public class DelayQueueManager {
private static final DelayQueue<DelayMessage> delayQueue = new DelayQueue<>();
private static final ExecutorService executorService = Executors.newFixedThreadPool(5);
// 启动消费者
public void start() {
System.out.println("延迟队列服务启动...");
for (int i = 0; i < 5; i++) {
executorService.execute(() -> {
while (true) {
try {
// 获取并移除延迟队列的头元素,如果没有到期则返回 null
DelayMessage message = delayQueue.poll();
if (message != null) {
System.out.println("消费者收到消息: " + message);
// 处理消息
processMessage(message);
} else {
// 没有消息,等待一段时间再检查
Thread.sleep(1000);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
});
}
}
// 处理消息
private void processMessage(DelayMessage message) {
// 这里可以放具体的业务处理逻辑
System.out.println("处理延迟消息: " + message.getContent() +
", 当前时间: " + System.currentTimeMillis());
}
// 添加延迟消息
public void addMessage(DelayMessage message) {
System.out.println("添加延迟消息: " + message +
", 当前时间: " + System.currentTimeMillis());
delayQueue.put(message);
}
// 停止服务
public void stop() {
executorService.shutdown();
System.out.println("延迟队列服务已停止");
}
public static void main(String[] args) throws InterruptedException {
DelayQueueManager manager = new DelayQueueManager();
// 启动服务
manager.start();
// 添加测试消息,延迟时间分别为 3秒、5秒、10秒
DelayMessage msg1 = new DelayMessage("001", "订单1超时关闭", 3000);
DelayMessage msg2 = new DelayMessage("002", "订单2超时关闭", 5000);
DelayMessage msg3 = new DelayMessage("003", "订单3超时关闭", 10000);
manager.addMessage(msg1);
manager.addMessage(msg2);
manager.addMessage(msg3);
// 运行一段时间后停止
Thread.sleep(15000);
manager.stop();
}
}
基于 Redis 的延迟队列实现
依赖准备
<dependency>
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
<version>4.4.0</version>
</dependency>
Redis延迟队列实现
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPool;
import redis.clients.jedis.Tuple;
import com.alibaba.fastjson.JSON;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
/**
* Redis延迟队列实现(使用 ZSet)
*/
public class RedisDelayQueue {
private static final String DELAY_QUEUE_KEY = "delay:queue";
private static final String READY_QUEUE_KEY = "ready:queue";
private JedisPool jedisPool;
private ScheduledExecutorService scheduledExecutor;
public RedisDelayQueue() {
// 初始化Redis连接池
jedisPool = new JedisPool("localhost", 6379);
scheduledExecutor = Executors.newScheduledThreadPool(4);
}
/**
* 添加延迟任务
* @param task 任务内容
* @param delay 延迟时间
* @param unit 时间单位
*/
public String addTask(String task, long delay, TimeUnit unit) {
String taskId = UUID.randomUUID().toString();
long expireTime = System.currentTimeMillis() + unit.toMillis(delay);
try (Jedis jedis = jedisPool.getResource()) {
// 将任务添加到延迟队列,score为过期时间
jedis.zadd(DELAY_QUEUE_KEY, expireTime, taskId);
// 存储任务内容
TaskData taskData = new TaskData(taskId, task, expireTime);
jedis.hset("task:data", taskId, JSON.toJSONString(taskData));
System.out.println("添加延迟任务: " + taskId + ", 过期时间: " + expireTime);
return taskId;
}
}
/**
* 检查并转移到期任务
*/
private void transferExpiredTasks() {
try (Jedis jedis = jedisPool.getResource()) {
// 获取到期的任务
Set<String> expiredTasks = jedis.zrangeByScore(DELAY_QUEUE_KEY,
0, System.currentTimeMillis());
if (!expiredTasks.isEmpty()) {
for (String taskId : expiredTasks) {
// 将任务从延迟队列移动到就绪队列
String taskData = jedis.hget("task:data", taskId);
if (taskData != null) {
// 添加到就绪队列
jedis.lpush(READY_QUEUE_KEY, taskData);
// 移除延迟队列中的任务
jedis.zrem(DELAY_QUEUE_KEY, taskId);
// 清理任务数据
jedis.hdel("task:data", taskId);
System.out.println("任务已到期,转移到就绪队列: " + taskId);
}
}
}
}
}
/**
* 启动延迟队列服务
*/
public void start() {
// 定时检查到期的任务(每秒执行一次)
scheduledExecutor.scheduleAtFixedRate(this::transferExpiredTasks,
0, 1, TimeUnit.SECONDS);
// 启动消费者线程处理就绪队列
for (int i = 0; i < 3; i++) {
scheduledExecutor.execute(() -> {
while (!Thread.currentThread().isInterrupted()) {
try (Jedis jedis = jedisPool.getResource()) {
// 从就绪队列获取任务(阻塞式)
String taskData = jedis.brpop(0, READY_QUEUE_KEY);
if (taskData != null) {
// 处理任务
TaskData data = JSON.parseObject(taskData, TaskData.class);
processTask(data);
}
} catch (Exception e) {
e.printStackTrace();
}
}
});
}
System.out.println("Redis延迟队列服务已启动");
}
/**
* 处理任务
*/
private void processTask(TaskData taskData) {
System.out.println("处理任务: " + taskData.getTaskId() +
", 内容: " + taskData.getTask() +
", 当前时间: " + System.currentTimeMillis());
// 在这里执行具体的业务逻辑
}
/**
* 停止服务
*/
public void stop() {
scheduledExecutor.shutdown();
jedisPool.close();
System.out.println("Redis延迟队列服务已停止");
}
// 任务数据类
static class TaskData {
private String taskId;
private String task;
private long expireTime;
public TaskData() {}
public TaskData(String taskId, String task, long expireTime) {
this.taskId = taskId;
this.task = task;
this.expireTime = expireTime;
}
public String getTaskId() { return taskId; }
public void setTaskId(String taskId) { this.taskId = taskId; }
public String getTask() { return task; }
public void setTask(String task) { this.task = task; }
public long getExpireTime() { return expireTime; }
public void setExpireTime(long expireTime) { this.expireTime = expireTime; }
}
public static void main(String[] args) throws InterruptedException {
RedisDelayQueue delayQueue = new RedisDelayQueue();
// 启动服务
delayQueue.start();
// 添加测试任务
delayQueue.addTask("任务1:超时关闭订单", 5, TimeUnit.SECONDS);
delayQueue.addTask("任务2:超时确认收货", 10, TimeUnit.SECONDS);
delayQueue.addTask("任务3:定时发送提醒", 15, TimeUnit.SECONDS);
// 运行一段时间
Thread.sleep(20000);
delayQueue.stop();
}
}
使用 ScheduledExecutorService 实现延迟任务
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
/**
* 基于 ScheduledExecutorService 的延迟任务实现
*/
public class ScheduledDelayQueue {
private ScheduledExecutorService executorService;
private int threadPoolSize;
public ScheduledDelayQueue(int threadPoolSize) {
this.threadPoolSize = threadPoolSize;
this.executorService = Executors.newScheduledThreadPool(threadPoolSize);
}
/**
* 添加一次性延迟任务
*/
public void addOneTimeTask(Runnable task, long delay, TimeUnit unit) {
System.out.println("添加一次性延迟任务: " + delay + " " + unit);
executorService.schedule(task, delay, unit);
}
/**
* 添加周期性任务(固定频率)
*/
public void addFixedRateTask(Runnable task, long initialDelay,
long period, TimeUnit unit) {
System.out.println("添加固定频率任务");
executorService.scheduleAtFixedRate(task, initialDelay, period, unit);
}
/**
* 添加周期性任务(固定延迟)
*/
public void addFixedDelayTask(Runnable task, long initialDelay,
long delay, TimeUnit unit) {
System.out.println("添加固定延迟任务");
executorService.scheduleWithFixedDelay(task, initialDelay, delay, unit);
}
/**
* 停止服务
*/
public void shutdown() {
executorService.shutdown();
try {
if (!executorService.awaitTermination(60, TimeUnit.SECONDS)) {
executorService.shutdownNow();
}
} catch (InterruptedException e) {
executorService.shutdownNow();
Thread.currentThread().interrupt();
}
System.out.println("调度服务已停止");
}
public static void main(String[] args) throws InterruptedException {
ScheduledDelayQueue queue = new ScheduledDelayQueue(4);
// 一次性延迟任务
queue.addOneTimeTask(() -> {
System.out.println("3秒后执行一次性延迟任务: " +
System.currentTimeMillis());
}, 3, TimeUnit.SECONDS);
// 固定频率任务(每2秒执行一次)
queue.addFixedRateTask(() -> {
System.out.println("固定频率任务,每2秒执行: " +
System.currentTimeMillis());
}, 0, 2, TimeUnit.SECONDS);
// 固定延迟任务(上一次执行完后延迟1秒)
queue.addFixedDelayTask(() -> {
System.out.println("固定延迟任务: " + System.currentTimeMillis());
try {
Thread.sleep(1000); // 模拟任务执行时间
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, 0, 1, TimeUnit.SECONDS);
// 运行10秒后停止
Thread.sleep(10000);
queue.shutdown();
}
}
应用场景示例:订单超时处理
import java.util.concurrent.DelayQueue;
import java.util.concurrent.TimeUnit;
/**
* 订单超时处理示例
*/
public class OrderTimeoutHandler {
private static final DelayQueue<DelayMessage> orderDelayQueue = new DelayQueue<>();
// 订单信息类
static class OrderInfo {
String orderId;
String userId;
double amount;
long createTime;
public OrderInfo(String orderId, String userId, double amount, long createTime) {
this.orderId = orderId;
this.userId = userId;
this.amount = amount;
this.createTime = createTime;
}
@Override
public String toString() {
return "订单{" +
"订单号='" + orderId + '\'' +
", 用户ID='" + userId + '\'' +
", 金额=" + amount +
", 创建时间=" + createTime +
'}';
}
}
// 启动延迟队列消费者
public static void startConsumer() {
Thread consumer = new Thread(() -> {
while (true) {
try {
DelayMessage message = orderDelayQueue.take();
String orderId = message.getMessageId();
System.out.println("订单超时处理: " + orderId);
// 检查订单状态
checkOrderStatus(orderId);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
});
consumer.start();
}
// 创建订单并加入延迟队列
public static OrderInfo createOrder(String orderId, String userId, double amount) {
OrderInfo order = new OrderInfo(orderId, userId, amount, System.currentTimeMillis());
System.out.println("创建订单: " + order);
// 添加延迟任务,30分钟超时
DelayMessage message = new DelayMessage(orderId,
"订单超时关闭",
30 * 60 * 1000);
orderDelayQueue.put(message);
return order;
}
// 检查订单状态
private static void checkOrderStatus(String orderId) {
// 这里需要查询数据库获取订单状态
System.out.println("检查订单 " + orderId + " 的支付状态");
// 如果未支付,则关闭订单
System.out.println("订单 " + orderId + " 超时未支付,已自动关闭");
}
public static void main(String[] args) throws InterruptedException {
// 启动消费者
startConsumer();
// 模拟创建订单
createOrder("O2024001", "U1001", 200.50);
createOrder("O2024002", "U1002", 150.00);
createOrder("O2024003", "U1003", 899.00);
// 演示用,手动添加一个3秒过期的任务
DelayMessage testMessage = new DelayMessage("TEST001", "测试过期", 3000);
orderDelayQueue.put(testMessage);
System.out.println("添加测试消息,3秒后过期");
// 运行一段时间
Thread.sleep(10000);
}
}
注意事项
-
延迟队列选择:
- 小型应用、单机部署:使用 JDK DelayQueue
- 分布式系统、需要持久化:使用 Redis ZSet
- 简单定时任务:使用 ScheduledExecutorService
-
可靠性:
- DelayedQueue 不持久化,程序重启会丢失
- Redis 延迟队列需要额外处理持久化
-
性能考虑:
- 大批量任务时推荐使用 Redis
- 考虑使用时间轮算法(如 Netty 的 HashedWheelTimer)提高性能
-
生产环境建议:
- 结合消息队列(如 RabbitMQ、Kafka)的延迟消息功能
- 使用成熟的框架(如 Quartz、XXL-JOB)
这些示例代码可以直接运行,根据实际业务需求进行修改即可。