Redis ZSet 延迟队列实战:MySQL 订单超时关单的可靠实现与重试机制
订单超时自动关闭是电商系统里最典型的延迟任务场景:用户下单后如果 30 分钟未支付,系统需要自动取消订单并释放库存。这个需求看似简单,但实现起来要考虑可靠性、幂等性、重试机制和最终一致性。本文以 Spring Boot 3 + MySQL 8 + Redis 7 为技术栈,手写一个基于 Redis ZSet 的延迟队列,完成订单超时关单的完整闭环。
一、为什么不用现成的消息队列?
延迟任务实现方案有不少,我们列举最常见的三种:
| 方案 |
实现复杂度 |
可靠性 |
延迟精度 |
适用规模 |
备注 |
| Redis ZSet 轮询 |
低,自己写轮询 |
中,需要自行处理 ACK/重试 |
秒级 |
中小规模 |
无额外中间件,运维成本低 |
| Redisson DelayedQueue |
低,封装好 |
高,自带 ACK/重试 |
毫秒级 |
中大规模 |
依赖 Redisson,侵入业务代码 |
| RabbitMQ 延迟插件 |
中,需要安装插件 |
高,依赖 MQ 本身 |
秒级 |
大规模 |
需要运维 RabbitMQ,增加系统复杂度 |
对于创业初期或中小型项目,直接引入 RabbitMQ 可能属于过度设计。Redis 是几乎所有后端系统都会用到的组件,基于 ZSet 实现延迟队列可以在不增加运维负担的前提下,快速解决业务痛点。这也是当初我在团队中推动这个方案的原因——技术选型不是选择最酷的,而是选择当前阶段最合适的。等业务量起来后,再平滑迁移到 RabbitMQ 或 RocketMQ 也不迟。
二、ZSet 延迟队列的设计思路
Redis 的 ZSet 是一个有序集合,每个成员(member)关联一个分数(score)。我们可以把任务的执行时间戳作为 score,把任务内容序列化成 JSON 字符串作为 member。定时扫描任务时,只需要查询 score 小于当前时间戳的成员即可。
核心流程如下:
1 2 3 4 5 6 7 8 9 10 11 12 13
| 生产者 Redis ZSet 消费者 | | | |--- ZADD delay:order:close ----->| | | score = 当前时间 + 延迟秒数 | | | |<--- 定时轮询 ZRANGEBYSCORE ---| | | score <= now | | |--- 返回到期任务列表 -------->| | | | | |<--- 原子弹出到期任务 (Lua) ---| | | | | | 执行业务操作 | | | | |<--- 成功则 ACK,失败则重新入队 --|
|
这里有两个关键点:
- 原子弹出:多个消费者实例同时轮询可能会取到同一个任务。必须用 Lua 脚本保证「查询到期任务 + 删除任务」的原子性,否则会重复消费。
- 重试机制:业务处理失败时,需要把任务重新放回队列,并增加延迟时间和重试次数。超过最大重试次数则标记为失败,进入人工处理或死信队列。
三、项目准备:完整的 Maven 依赖与配置
3.1 创建 Spring Boot 项目
先给出完整的 pom.xml,确保依赖正确:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68
| <?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.2.5</version> <relativePath/> </parent>
<groupId>com.example</groupId> <artifactId>redis-delay-queue-demo</artifactId> <version>1.0.0</version>
<properties> <java.version>17</java.version> </properties>
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-jpa</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency> <dependency> <groupId>com.mysql</groupId> <artifactId>mysql-connector-j</artifactId> <scope>runtime</scope> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> </dependencies>
<build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> <configuration> <excludes> <exclude> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> </exclude> </excludes> </configuration> </plugin> </plugins> </build> </project>
|
验证环境:JDK 17.0.10、Spring Boot 3.2.5、MySQL 8.0.36、Redis 7.2.4。
3.2 application.yml 配置
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26
| server: port: 8080
spring: datasource: url: jdbc:mysql://localhost:3306/order_db?useSSL=false&serverTimezone=Asia/Shanghai&characterEncoding=utf8 username: root password: root driver-class-name: com.mysql.cj.jdbc.Driver jpa: hibernate: ddl-auto: update show-sql: true properties: hibernate: format_sql: true data: redis: host: localhost port: 6379 timeout: 3000ms lettuce: pool: max-active: 8 max-idle: 8 min-idle: 0
|
3.3 初始化 MySQL 表结构
使用 JPA 自动建表,但为了演示方便,给出对应的建表 SQL:
1 2 3 4 5 6 7 8 9 10 11 12 13
| CREATE DATABASE IF NOT EXISTS order_db DEFAULT CHARACTER SET utf8mb4;
USE order_db;
CREATE TABLE `orders` ( `id` bigint NOT NULL AUTO_INCREMENT, `order_no` varchar(64) NOT NULL COMMENT '订单号', `status` varchar(20) NOT NULL COMMENT '状态:WAIT_PAY-待支付 PAID-已支付 CLOSED-已关闭', `create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (`id`), UNIQUE KEY `uk_order_no` (`order_no`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单表';
|
验证环境:MySQL 8.0.36。
四、核心代码实现
4.1 订单实体与状态枚举
先定义订单状态枚举,这里直接使用字符串常量,便于和数据库字段对应:
1 2 3 4 5 6 7
| package com.example.delayqueue.enums;
public enum OrderStatus { WAIT_PAY, PAID, CLOSED }
|
订单实体类:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51
| package com.example.delayqueue.entity;
import com.example.delayqueue.enums.OrderStatus; import jakarta.persistence.*; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor;
import java.time.LocalDateTime;
@Data @Builder @NoArgsConstructor @AllArgsConstructor @Entity @Table(name = "orders") public class Order {
@Id @GeneratedValue(strategy = GenerationType.IDENTITY) private Long id;
@Column(name = "order_no", nullable = false, unique = true, length = 64) private String orderNo;
@Enumerated(EnumType.STRING) @Column(nullable = false, length = 20) private OrderStatus status;
@Column(name = "create_time", nullable = false, updatable = false) private LocalDateTime createTime;
@Column(name = "update_time", nullable = false) private LocalDateTime updateTime;
@PrePersist public void prePersist() { if (createTime == null) { createTime = LocalDateTime.now(); } if (updateTime == null) { updateTime = LocalDateTime.now(); } }
@PreUpdate public void preUpdate() { updateTime = LocalDateTime.now(); } }
|
数据访问层:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28
| package com.example.delayqueue.repository;
import com.example.delayqueue.entity.Order; import com.example.delayqueue.enums.OrderStatus; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.data.jpa.repository.Modifying; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.query.Param; import org.springframework.transaction.annotation.Transactional;
import java.util.Optional;
public interface OrderRepository extends JpaRepository<Order, Long> {
Optional<Order> findByOrderNo(String orderNo);
@Modifying @Transactional @Query("UPDATE Order o SET o.status = :newStatus WHERE o.orderNo = :orderNo AND o.status = :oldStatus") int updateStatusIfInStatus(@Param("orderNo") String orderNo, @Param("oldStatus") OrderStatus oldStatus, @Param("newStatus") OrderStatus newStatus); }
|
4.2 延迟任务模型
延迟任务对象 DelayTask 会被序列化为 JSON 字符串存到 ZSet 的 member 中。它包含任务 ID、业务类型、业务数据(这里是订单号)以及重试信息。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47
| package com.example.delayqueue.model;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor;
import java.time.LocalDateTime;
@Data @Builder @NoArgsConstructor @AllArgsConstructor @JsonIgnoreProperties(ignoreUnknown = true) public class DelayTask {
private String taskId;
private String type;
private String payload;
private int retryCount;
private int maxRetry;
private LocalDateTime createTime; }
|
4.3 Redis 延迟队列封装
这是整个方案的心脏部分。RedisDelayQueue 封装了任务入队和原子弹到期任务的操作。原子弹出使用 Lua 脚本,避免并发下重复消费。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104
| package com.example.delayqueue.core;
import com.example.delayqueue.model.DelayTask; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.data.redis.core.script.DefaultRedisScript; import org.springframework.stereotype.Component;
import java.time.Instant; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.UUID;
@Slf4j @Component @RequiredArgsConstructor public class RedisDelayQueue {
private final StringRedisTemplate stringRedisTemplate; private final ObjectMapper objectMapper;
private static final DefaultRedisScript<List> POP_READY_SCRIPT = new DefaultRedisScript<>( "local items = redis.call('ZRANGEBYSCORE', KEYS[1], 0, ARGV[1]) " + "if #items > 0 then " + " redis.call('ZREM', KEYS[1], unpack(items)) " + "end " + "return items", List.class );
public void addTask(String queueKey, DelayTask task, long delaySeconds) { if (task.getTaskId() == null || task.getTaskId().isBlank()) { task.setTaskId(UUID.randomUUID().toString()); } if (task.getCreateTime() == null) { task.setCreateTime(java.time.LocalDateTime.now()); } long executeAt = System.currentTimeMillis() + delaySeconds * 1000; try { String member = objectMapper.writeValueAsString(task); stringRedisTemplate.opsForZSet().add(queueKey, member, executeAt); log.info("延迟任务已入队: queue={}, taskId={}, delay={}s", queueKey, task.getTaskId(), delaySeconds); } catch (JsonProcessingException e) { throw new RuntimeException("序列化任务失败", e); } }
public List<DelayTask> pollReadyTasks(String queueKey) { long now = System.currentTimeMillis(); List<String> members = stringRedisTemplate.execute( POP_READY_SCRIPT, Collections.singletonList(queueKey), String.valueOf(now) ); if (members == null || members.isEmpty()) { return Collections.emptyList(); } List<DelayTask> tasks = new ArrayList<>(members.size()); for (String member : members) { try { tasks.add(objectMapper.readValue(member, DelayTask.class)); } catch (JsonProcessingException e) { log.error("反序列化延迟任务失败: {}", member, e); } } log.info("弹出到期任务数量: {}", tasks.size()); return tasks; }
public void removeTask(String queueKey, DelayTask task) { try { stringRedisTemplate.opsForZSet().remove(queueKey, objectMapper.writeValueAsString(task)); } catch (JsonProcessingException e) { log.error("移除任务失败", e); } } }
|
说明:这里使用的 Lua 脚本兼容 Redis 2.6+,没有使用 Redis 5.0 的 ZPOPMIN,因为很多生产环境还在用老版本。ZRANGEBYSCORE + ZREM 的组合在 Lua 脚本中具有原子性,可以保证多个消费者同时轮询时任务不会被重复弹出。
4.4 订单超时关单处理器
该处理器负责真正执行关单业务:更新订单状态。更新操作使用带 WHERE status = 'WAIT_PAY' 的条件更新,天然支持幂等。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64
| package com.example.delayqueue.handler;
import com.example.delayqueue.entity.Order; import com.example.delayqueue.enums.OrderStatus; import com.example.delayqueue.model.DelayTask; import com.example.delayqueue.repository.OrderRepository; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.springframework.transaction.annotation.Transactional;
import java.util.Optional;
@Slf4j @Component @RequiredArgsConstructor public class OrderCloseTaskHandler {
private final OrderRepository orderRepository;
@Transactional public boolean handle(DelayTask task) { String orderNo = task.getPayload(); log.info("开始处理关单任务: taskId={}, orderNo={}", task.getTaskId(), orderNo);
Optional<Order> orderOpt = orderRepository.findByOrderNo(orderNo); if (orderOpt.isEmpty()) { log.error("订单不存在,无法关闭: orderNo={}", orderNo); return true; }
Order order = orderOpt.get(); if (order.getStatus() == OrderStatus.CLOSED) { log.info("订单已经是关闭状态,无需重复处理: orderNo={}", orderNo); return true; } if (order.getStatus() == OrderStatus.PAID) { log.info("订单已支付,无法关闭: orderNo={}", orderNo); return true; }
int affected = orderRepository.updateStatusIfInStatus( orderNo, OrderStatus.WAIT_PAY, OrderStatus.CLOSED );
if (affected > 0) { log.info("订单关闭成功: orderNo={}", orderNo); return true; } else { log.warn("订单状态已变化,关闭失败(可能已被支付或并发关闭): orderNo={}", orderNo); return true; } } }
|
这里有一个细节:如果订单状态已变为 PAID 或 CLOSED,我们直接返回 true(成功),而不是 false(重试)。因为这种情况重试是无效的,而且可能造成任务永远无法完成。这也是幂等性设计的一部分——业务上的“已处理”和“处理成功”需要区分对待。
4.5 延迟队列调度器与重试机制
调度器负责周期性扫描到期任务,提交到线程池执行。线程池是必需的,否则轮询线程会被阻塞在业务处理上。重试逻辑也在这里统一处理。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96
| package com.example.delayqueue.scheduler;
import com.example.delayqueue.core.RedisDelayQueue; import com.example.delayqueue.handler.OrderCloseTaskHandler; import com.example.delayqueue.model.DelayTask; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component;
import java.util.List; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit;
@Slf4j @Component @RequiredArgsConstructor public class DelayQueueScheduler {
private static final String ORDER_CLOSE_QUEUE = "delay:queue:order_close"; private static final int MAX_RETRY = 3; private static final long BASE_RETRY_DELAY_SECONDS = 30;
private final RedisDelayQueue redisDelayQueue; private final OrderCloseTaskHandler orderCloseTaskHandler;
private ExecutorService executorService;
@PostConstruct public void init() { executorService = Executors.newFixedThreadPool(4); }
@PreDestroy public void destroy() { if (executorService != null) { executorService.shutdown(); try { if (!executorService.awaitTermination(10, TimeUnit.SECONDS)) { executorService.shutdownNow(); } } catch (InterruptedException e) { executorService.shutdownNow(); Thread.currentThread().interrupt(); } } }
@Scheduled(fixedDelay = 1000, initialDelay = 3000) public void scanAndProcess() { List<DelayTask> tasks = redisDelayQueue.pollReadyTasks(ORDER_CLOSE_QUEUE); if (tasks.isEmpty()) { return; } for (DelayTask task : tasks) { executorService.submit(() -> processTask(task)); } }
private void processTask(DelayTask task) { try { boolean success = orderCloseTaskHandler.handle(task); if (success) { log.info("任务处理成功,已 ACK: taskId={}", task.getTaskId()); } else { retryTask(task); } } catch (Exception e) { log.error("任务处理异常,将重试: taskId={}", task.getTaskId(), e); retryTask(task); } }
private void retryTask(DelayTask task) { int retryCount = task.getRetryCount() + 1; if (retryCount > MAX_RETRY) { log.error("任务重试次数超过上限,标记为失败并丢弃: taskId={}, payload={}", task.getTaskId(), task.getPayload()); return; } task.setRetryCount(retryCount); long delaySeconds = BASE_RETRY_DELAY_SECONDS * retryCount; log.warn("任务处理失败,第 {} 次重试,延迟 {} 秒: taskId={}", retryCount, delaySeconds, task.getTaskId()); redisDelayQueue.addTask(ORDER_CLOSE_QUEUE, task, delaySeconds); } }
|
注意:这里我们使用 fixedDelay = 1000 表示上一次扫描完成后等待 1 秒再开始下一次。如果业务量大,可以适当调小扫描间隔,但要注意 Redis 的查询压力。通常 1 秒已经足够满足秒级精度。
4.6 订单服务与测试接口
最后提供创建订单和手动触发关单的 REST 接口,方便验证整个流程。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60
| package com.example.delayqueue.service;
import com.example.delayqueue.core.RedisDelayQueue; import com.example.delayqueue.entity.Order; import com.example.delayqueue.enums.OrderStatus; import com.example.delayqueue.model.DelayTask; import com.example.delayqueue.repository.OrderRepository; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional;
import java.time.LocalDateTime; import java.util.UUID;
@Slf4j @Service @RequiredArgsConstructor public class OrderService {
private static final String ORDER_CLOSE_QUEUE = "delay:queue:order_close"; private static final long CLOSE_DELAY_SECONDS = 30;
private final OrderRepository orderRepository; private final RedisDelayQueue redisDelayQueue;
@Transactional public Order createOrder() { Order order = Order.builder() .orderNo(UUID.randomUUID().toString().replace("-", "")) .status(OrderStatus.WAIT_PAY) .build(); order = orderRepository.save(order); log.info("订单创建成功: orderNo={}", order.getOrderNo());
DelayTask task = DelayTask.builder() .type("order_close") .payload(order.getOrderNo()) .retryCount(0) .maxRetry(3) .build(); redisDelayQueue.addTask(ORDER_CLOSE_QUEUE, task, CLOSE_DELAY_SECONDS); log.info("订单关单任务已投递,将在 {} 秒后执行", CLOSE_DELAY_SECONDS); return order; }
@Transactional public void payOrder(String orderNo) { Order order = orderRepository.findByOrderNo(orderNo) .orElseThrow(() -> new RuntimeException("订单不存在")); if (order.getStatus() == OrderStatus.WAIT_PAY) { order.setStatus(OrderStatus.PAID); orderRepository.save(order); log.info("订单支付成功: orderNo={}", orderNo); } else { log.warn("订单状态不允许支付: orderNo={}, status={}", orderNo, order.getStatus()); } } }
|
控制器:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27
| package com.example.delayqueue.controller;
import com.example.delayqueue.entity.Order; import com.example.delayqueue.service.OrderService; import lombok.RequiredArgsConstructor; import org.springframework.http.ResponseEntity; import org.springframework.web.bind.annotation.*;
@RestController @RequestMapping("/orders") @RequiredArgsConstructor public class OrderController {
private final OrderService orderService;
@PostMapping public ResponseEntity<Order> createOrder() { Order order = orderService.createOrder(); return ResponseEntity.ok(order); }
@PostMapping("/{orderNo}/pay") public ResponseEntity<Void> payOrder(@PathVariable String orderNo) { orderService.payOrder(orderNo); return ResponseEntity.noContent().build(); } }
|
五、运行验证
5.1 启动环境
确保 MySQL 和 Redis 已启动,并已创建 order_db 数据库和 orders 表。然后启动 Spring Boot 应用。
5.2 测试流程
- 创建订单:
1
| curl -X POST http://localhost:8080/orders
|
返回类似:
{
"id": 1,
"orderNo": "a1b2c3d4e5f6",
"status": "WAIT