本文目录导读:

- 基于 Redis 的延迟队列(ZSET实现)
- RabbitMQ 延迟消息(死信队列实现)
- 使用 RocketMQ 延迟消息(内置延迟等级)
- Java 本地延迟队列(ScheduledExecutorService)
- 场景化完整案例
- 技术选型建议
- 注意事项
我将为您提供几个常见的延迟消息实现案例,涵盖不同技术栈和场景。
基于 Redis 的延迟队列(ZSET实现)
场景:订单超时自动取消
@Service
public class OrderDelayService {
@Autowired
private StringRedisTemplate redisTemplate;
private static final String DELAY_QUEUE_KEY = "order:delay:queue";
/**
* 添加延迟任务
* @param orderId 订单ID
* @param delayTime 延迟时间(毫秒)
*/
public void addDelayTask(String orderId, long delayTime) {
// 到期时间 = 当前时间 + 延迟时间
long expireTime = System.currentTimeMillis() + delayTime;
redisTemplate.opsForZSet().add(DELAY_QUEUE_KEY, orderId, expireTime);
log.info("订单[{}]已加入延迟队列,将在{}秒后执行", orderId, delayTime/1000);
}
/**
* 消费者:定时扫描处理过期任务
*/
@Scheduled(fixedDelay = 1000) // 每秒执行一次
public void consumeDelayTask() {
// 获取当前时间
long currentTime = System.currentTimeMillis();
// 获取所有到期的任务(分数 <= 当前时间)
Set<String> expiredOrders = redisTemplate.opsForZSet()
.rangeByScore(DELAY_QUEUE_KEY, 0, currentTime);
if (CollectionUtils.isEmpty(expiredOrders)) {
return;
}
for (String orderId : expiredOrders) {
// 原子性移除任务,避免并发重复消费
Long removed = redisTemplate.opsForZSet()
.remove(DELAY_QUEUE_KEY, orderId);
if (removed != null && removed > 0) {
// 执行具体业务逻辑
handleExpiredOrder(orderId);
log.info("订单[{}]已超时,执行自动取消", orderId);
}
}
}
private void handleExpiredOrder(String orderId) {
// 业务处理:检查订单状态,更新为已取消等
Order order = orderMapper.selectById(orderId);
if (order != null && "待支付".equals(order.getStatus())) {
order.setStatus("已取消");
orderMapper.updateById(order);
}
}
}
使用示例
@Test
public void testDelayMessage() throws InterruptedException {
// 1. 创建订单
String orderId = "ORDER_2024001";
orderService.createOrder(orderId);
// 2. 添加延迟任务:30分钟后自动取消
orderDelayService.addDelayTask(orderId, 30 * 60 * 1000);
// 3. 等待测试
Thread.sleep(5000);
}
RabbitMQ 延迟消息(死信队列实现)
场景:支付超时提醒
@Configuration
public class RabbitDelayConfig {
// 交换机
public static final String ORDER_EXCHANGE = "order.exchange";
public static final String ORDER_DELAY_EXCHANGE = "order.delay.exchange";
// 队列
public static final String ORDER_QUEUE = "order.queue";
public static final String ORDER_DELAY_QUEUE = "order.delay.queue";
// 路由键
public static final String ORDER_ROUTING_KEY = "order.routing";
public static final String ORDER_DELAY_ROUTING_KEY = "order.delay.routing";
// 延迟时间
public static final long DELAY_TIME = 30 * 60 * 1000; // 30分钟
@Bean
public Queue delayQueue() {
Map<String, Object> args = new HashMap<>();
// 设置死信交换机
args.put("x-dead-letter-exchange", ORDER_EXCHANGE);
// 设置死信路由键
args.put("x-dead-letter-routing-key", ORDER_ROUTING_KEY);
// 设置消息过期时间
args.put("x-message-ttl", DELAY_TIME);
return QueueBuilder.durable(ORDER_DELAY_QUEUE).withArguments(args).build();
}
@Bean
public Queue orderQueue() {
return QueueBuilder.durable(ORDER_QUEUE).build();
}
@Bean
public DirectExchange exchange() {
return new DirectExchange(ORDER_EXCHANGE);
}
@Bean
public DirectExchange delayExchange() {
return new DirectExchange(ORDER_DELAY_EXCHANGE);
}
// 绑定
@Bean
public Binding delayBinding() {
return BindingBuilder.bind(delayQueue())
.to(delayExchange()).with(ORDER_DELAY_ROUTING_KEY);
}
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue())
.to(exchange()).with(ORDER_ROUTING_KEY);
}
}
生产者与消费者
@Service
public class PayTimeoutService {
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 发送延迟消息
*/
public void sendDelayMessage(String orderId) {
// 创建消息
Message message = MessageBuilder
.withBody(orderId.getBytes(StandardCharsets.UTF_8))
.setContentType(MessageProperties.CONTENT_TYPE_JSON)
.build();
// 发送到延迟队列
rabbitTemplate.send(ORDER_DELAY_EXCHANGE, ORDER_DELAY_ROUTING_KEY, message);
log.info("支付提醒延迟消息已发送,订单ID:{}", orderId);
}
/**
* 消费者:处理超时未支付的订单
*/
@RabbitListener(queues = ORDER_QUEUE)
public void handlePayTimeout(String orderId) {
log.info("收到延迟消息,订单ID:{}", orderId);
// 业务处理:检查支付状态
Order order = orderService.getById(orderId);
if (order != null && "未支付".equals(order.getStatus())) {
// 发送支付提醒短信或推送
smsService.sendPayRemind(order.getMemberId(), orderId);
log.info("已向用户发送支付提醒");
}
}
}
使用 RocketMQ 延迟消息(内置延迟等级)
场景:秒杀活动结束通知
@Service
public class SeckillNotificationService {
@Autowired
private RocketMQTemplate rocketMQTemplate;
private static final String TOPIC = "seckill-notification";
/**
* 发送延迟通知
*/
public void sendSeckillEndNotification(String activityId) {
// RocketMQ延迟等级:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
// 这里使用 30分钟 = 延迟等级 16
Message<String> message = MessageBuilder
.withPayload(activityId)
.build();
// 设置延迟等级
rocketMQTemplate.syncSend(
TOPIC,
message,
3000,
// 延迟等级,16 表示30分钟
10 // 延迟等级3 => 10秒
);
log.info("秒杀活动结束通知已发送,活动ID:{}", activityId);
}
/**
* 消费者
*/
@Service
@RocketMQMessageListener(
topic = TOPIC,
consumerGroup = "seckill-notification-group"
)
public static class NotificationConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String activityId) {
log.info("收到秒杀活动结束通知:{}", activityId);
// 批量处理参与用户
List<Long> userIds = seckillService.getParticipants(activityId);
notificationService.batchSend(userIds, "秒杀活动已结束");
}
}
}
Java 本地延迟队列(ScheduledExecutorService)
场景:定时任务调度
@Component
public class LocalDelayQueueExample {
private final ScheduledExecutorService scheduler =
Executors.newScheduledThreadPool(10);
private final DelayQueue<DelayTask> delayQueue = new DelayQueue<>();
/**
* 任务实体
*/
public static class DelayTask implements Delayed {
private final String taskId;
private final long executeTime;
private final Runnable task;
public DelayTask(String taskId, long delayMillis, Runnable task) {
this.taskId = taskId;
this.executeTime = System.currentTimeMillis() + delayMillis;
this.task = task;
}
@Override
public long getDelay(TimeUnit unit) {
return unit.convert(executeTime - System.currentTimeMillis(),
TimeUnit.MILLISECONDS);
}
@Override
public int compareTo(Delayed other) {
return Long.compare(this.executeTime,
((DelayTask) other).executeTime);
}
public void execute() {
task.run();
}
}
/**
* 添加延迟任务
*/
public void addTask(String taskId, long delayMillis, Runnable task) {
DelayTask delayTask = new DelayTask(taskId, delayMillis, task);
delayQueue.offer(delayTask);
log.info("延迟任务[{}]已添加,延迟{}毫秒", taskId, delayMillis);
}
/**
* 启动消费者线程
*/
@PostConstruct
public void startConsumer() {
scheduler.scheduleAtFixedRate(() -> {
try {
DelayTask task = delayQueue.poll(1, TimeUnit.SECONDS);
if (task != null) {
log.info("开始执行延迟任务:{}", task);
task.execute();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.error("延迟队列消费失败", e);
}
}, 0, 100, TimeUnit.MILLISECONDS);
}
/**
* 使用示例
*/
public void demo() {
// 添加不同类型的延迟任务
addTask("task-1", 5000, () -> {
System.out.println("5秒后执行的定时任务");
});
addTask("task-2", 10000, () -> {
System.out.println("10秒后执行的定时任务");
});
}
}
场景化完整案例
电商交易场景:订单状态流转
@Service
public class OrderStatusFlowService {
// 不同业务的延迟时间
private static final long PAY_TIMEOUT = 30 * 60 * 1000; // 30分钟支付超时
private static final long CONFIRM_TIMEOUT = 7 * 24 * 3600 * 1000; // 7天自动确认收货
private static final long REFUND_TIMEOUT = 24 * 3600 * 1000; // 24小时退款超时
@Autowired
private OrderDelayService orderDelayService;
/**
* 订单创建后设置支付超时
*/
public void createOrderAndSetTimeout(String orderId) {
// 创建订单
orderService.create(orderId);
// 设置30分钟支付超时
orderDelayService.addDelayTask(
orderId,
PAY_TIMEOUT,
DelayTaskType.PAY_TIMEOUT
);
}
/**
* 支付成功后设置自动确认收货
*/
public void paySuccessAndSetConfirm(String orderId) {
orderService.payMark(orderId);
// 7天后自动确认收货
delayTaskManager.addTask(
orderId,
CONFIRM_TIMEOUT,
DelayTaskType.AUTO_CONFIRM
);
}
/**
* 统一处理延迟任务
*/
@Scheduled(cron = "*/5 * * * * ?")
public void processDelayTasks() {
List<DelayTask> tasks = delayTaskManager.getExpiredTasks();
for (DelayTask task : tasks) {
switch (task.getTaskType()) {
case PAY_TIMEOUT:
handlePayTimeout(task.getBizId());
break;
case AUTO_CONFIRM:
handleAutoConfirm(task.getBizId());
break;
case REFUND_TIMEOUT:
handleRefundTimeout(task.getBizId());
break;
}
}
}
/**
* 处理支付超时
*/
private void handlePayTimeout(String orderId) {
Order order = orderService.getById(orderId);
if ("待支付".equals(order.getStatus())) {
orderService.cancelOrder(orderId, "支付超时");
inventoryService.restoreStock(order.getSkuId(), order.getQuantity());
log.info("订单[{}]支付超时,已自动取消", orderId);
}
}
}
技术选型建议
| 方案 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| Redis ZSET | 实现简单、高性能 | 需自建轮询、数据可能丢失 | 小型系统、对实时性要求不高 |
| RabbitMQ 死信 | 可靠、消息不丢失 | 需额外配置、延迟精度有限 | 金融、交易等关键业务 |
| RocketMQ | 内置延迟等级、高吞吐 | 依赖特定MQ | 大规模分布式系统 |
| Kafka | 高吞吐、可靠 | 延迟精度不足 | 日志分析、异步批处理 |
| Java DelayQueue | 轻量、简单 | 单机、无法持久化 | 单体应用、内存任务 |
注意事项
- 消息可靠性:生产环境建议使用 MQ 方案,确保消息不丢失
- 延迟精度:Redis 方案受轮询频率影响,MQ 有最小延迟限制
- 幂等性:消费端要做幂等处理,防止重复消费
- 监控告警:对延迟消息队列进行监控,设置积压告警
- 优雅关闭:提供关闭钩子,处理内存中的延迟任务
选择方案时需要结合业务场景、系统规模、技术栈等因素综合评估。