秒杀优化与消息队列
# 秒杀优化与消息队列
# 5. 你是怎么实现秒杀业务的?有出现什么问题吗?怎么解决的?
# 1. 优化秒杀异步下单的实现
当用户发起请求,此时会先请求 Nginx,Nginx 反向代理到 Tomcat,而 Tomcat 中的程序会进行串行操作,分为如下几个步骤:
- 查询优惠券
- 判断秒杀库存是否足够
- 查询订单
- 校验是否一人一单
- 扣减库存
- 创建订单
在这六个步骤中,有很多操作都是要去操作数据库的,而且还是一个线程串行执行,这样就会导致我们的程序执行很慢,所以我们需要异步程序执行。
优化方案:我们将耗时较短的逻辑判断放到 Redis 中,例如:库存是否充足、是否一人一单这样的操作。只要满足这两条操作,那我们是一定可以下单成功的,不用等数据真的写进数据库,我们直接告诉用户下单成功就好了。然后后台再开一个线程,后台线程再去慢慢执行队列里的消息,这样我们就能很快地完成下单业务。

这里还存在两个难点:
- 我们怎么在 Redis 中快速校验是否一人一单,还有库存判断?
- 我们校验一人一单和将下单数据写入数据库,这是两个线程,我们怎么知道下单是否完成?
我们需要将一些信息返回给前端,同时也将这些信息丢到异步 queue 中去,后续操作中,可以通过这个 id 来查询下单逻辑是否完成。
整体思路:当用户下单之后,判断库存是否充足,只需要取 Redis 中根据 key 找对应的 value 是否大于 0 即可。如果不充足,则直接结束。如果充足,则在 Redis 中判断用户是否可以下单,如果 set 集合中没有该用户的下单数据,则可以下单,并将 userId 和优惠券存入到 Redis 中,并且返回 0。整个过程需要保证是原子性的,所以我们要用 Lua 来操作。同时由于我们需要在 Redis 中查询优惠券信息,所以在我们新增秒杀优惠券的同时,需要将优惠券信息保存到 Redis 中。
完成以上逻辑判断时,我们只需要判断当前 Redis 中的返回值是否为 0,如果是 0,则表示可以下单,将信息保存到 queue 中去,然后返回,开一个线程来异步下单。其中订单可以通过返回订单的 id 来判断是否下单成功。

步骤:
- 新增秒杀优惠券的同时,将优惠券信息保存到 Redis 中
- 基于 Lua 脚本,判断秒杀库存、一人一单,决定用户是否秒杀成功
# 2. 基于阻塞队列实现秒杀优化
修改下单的操作,我们在下单时,是通过 Lua 表达式去原子执行判断逻辑。如果判断结果不为 0,返回错误信息;如果判断结果为 0,则将下单的逻辑保存到队列中去,然后异步执行。
需求:
- 如果秒杀成功,则将优惠券 id 和用户 id 封装后存入阻塞队列
- 开启线程任务,不断从阻塞队列中获取信息,实现异步下单功能
@Service
@Slf4j
public class VoucherOrderServiceImpl extends ServiceImpl<VoucherOrderMapper, VoucherOrder> implements IVoucherOrderService {
@Autowired
private ISeckillVoucherService seckillVoucherService;
@Autowired
private RedisIdWorker redisIdWorker;
@Resource
private StringRedisTemplate stringRedisTemplate;
@Resource
private RedissonClient redissonClient;
private IVoucherOrderService proxy;
private static final DefaultRedisScript<Long> SECKILL_SCRIPT;
static {
SECKILL_SCRIPT = new DefaultRedisScript();
SECKILL_SCRIPT.setLocation(new ClassPathResource("seckill.lua"));
SECKILL_SCRIPT.setResultType(Long.class);
}
private static final ExecutorService SECKILL_ORDER_EXECUTOR = Executors.newSingleThreadExecutor();
@PostConstruct
private void init() {
SECKILL_ORDER_EXECUTOR.submit(new VoucherOrderHandler());
}
private final BlockingQueue<VoucherOrder> orderTasks = new ArrayBlockingQueue<>(1024 * 1024);
private void handleVoucherOrder(VoucherOrder voucherOrder) {
// 1. 获取用户
Long userId = voucherOrder.getUserId();
// 2. 创建锁对象,作为兜底方案
RLock redisLock = redissonClient.getLock("order:" + userId);
// 3. 获取锁
boolean isLock = redisLock.tryLock();
// 4. 判断是否获取锁成功(理论上必成功,redis已经帮我们判断了)
if (!isLock) {
log.error("不允许重复下单!");
return;
}
try {
// 5. 使用代理对象,由于这里是另外一个线程
proxy.createVoucherOrder(voucherOrder);
} finally {
redisLock.unlock();
}
}
private class VoucherOrderHandler implements Runnable {
@Override
public void run() {
while (true) {
try {
// 1. 获取队列中的订单信息
VoucherOrder voucherOrder = orderTasks.take();
// 2. 创建订单
handleVoucherOrder(voucherOrder);
} catch (Exception e) {
log.error("订单处理异常", e);
}
}
}
}
@Override
public Result seckillVoucher(Long voucherId) {
Long result = stringRedisTemplate.execute(SECKILL_SCRIPT,
Collections.emptyList(), voucherId.toString(),
UserHolder.getUser().getId().toString());
if (result.intValue() != 0) {
return Result.fail(result.intValue() == 1 ? "库存不足" : "不能重复下单");
}
long orderId = redisIdWorker.nextId("order");
// 封装到voucherOrder中
VoucherOrder voucherOrder = new VoucherOrder();
voucherOrder.setVoucherId(voucherId);
voucherOrder.setUserId(UserHolder.getUser().getId());
voucherOrder.setId(orderId);
// 加入到阻塞队列
orderTasks.add(voucherOrder);
// 主线程获取代理对象
proxy = (IVoucherOrderService) AopContext.currentProxy();
return Result.ok(orderId);
}
@Transactional
public void createVoucherOrder(VoucherOrder voucherOrder) {
// 一人一单逻辑
Long userId = voucherOrder.getUserId();
Long voucherId = voucherOrder.getVoucherId();
synchronized (userId.toString().intern()) {
int count = query().eq("voucher_id", voucherId).eq("user_id", userId).count();
if (count > 0) {
log.error("你已经抢过优惠券了哦");
return;
}
// 5. 扣减库存
boolean success = seckillVoucherService.update()
.setSql("stock = stock - 1")
.eq("voucher_id", voucherId)
.gt("stock", 0)
.update();
if (!success) {
log.error("库存不足");
}
// 7. 将订单数据保存到表中
save(voucherOrder);
}
}
}
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
105
106
107
108
109
110
111
112
113
114
115
116
117
在优惠券秒杀系统中,当用户发出秒杀请求后,系统会将订单信息封装到一个 VoucherOrder 对象中,并放入阻塞队列 orderTasks 中。
阻塞队列是指一个内部长度固定的队列,在队列满时,新加入的元素被阻塞,直到队列中有元素被取出才能加入。这种队列通常用于并发编程中,用于保持线程安全的通信。
在这段代码中,使用 BlockingQueue 实现阻塞队列,并定义它的长度为 1024 * 1024,即最多可以同时处理 1024 * 1024 个订单请求。当订单请求被加入阻塞队列中后,另外一个线程会不断地从中取出请求,并进行处理。这样,订单请求的处理和商品库存的扣减是在不同的线程中进行的,避免了时间冲突的问题,提高了并发处理能力。
具体来说,在阻塞队列中,当有新的订单请求被加入时,处理该请求的线程会从队列中取出订单信息,并将它传递给 handleVoucherOrder() 方法进行处理。该方法会先获取用户 id,然后创建一个锁对象,锁定指定用户的订单。这是为了防止同一个用户多次下单。获得锁成功后,该方法会调用代理对象的 createVoucherOrder() 方法,在其中处理扣减库存和保存订单等操作。当订单处理完成后,线程将释放锁,并尝试去取下一个订单进行处理。
可以看出,阻塞队列的作用是将秒杀请求和订单处理隔离开来,保证了订单的正常处理,也增强了系统的并发运作能力。
小结:
- 秒杀业务的优化思路是什么?
- 先利用 Redis 完成库存容量、一人一单的判断,完成抢单业务
- 再将下单业务放入阻塞队列,利用独立线程异步下单
- 基于阻塞队列的异步秒杀存在哪些问题?
- 内存限制问题:我们现在使用的是 JDK 里的阻塞队列,它使用的是 JVM 的内存。如果在高并发的条件下,无数的订单都会放在阻塞队列里,可能就会造成内存溢出。所以我们在创建阻塞队列时,设置了一个长度,但是如果真的存满了,再有新的订单来往里塞,那就塞不进去了,存在内存限制问题。
- 数据安全问题:经典服务器宕机了,用户明明下单了,但是数据库里没看到。
# 6. 认识消息队列
# 消息队列概念
什么是消息队列? 字面意思就是存放消息的队列,最简单的消息队列模型包括 3 个角色:
- 消息队列:存储和管理消息,也被称为消息代理(Message Broker)
- 生产者:发送消息到消息队列
- 消费者:从消息队列获取消息并处理消息
使用队列的好处在于解耦:举个例子,快递员(生产者)把快递放到驿站/快递柜里去(Message Queue)去,我们(消费者)从快递柜/驿站去拿快递,这就是一个异步。如果耦合,那么快递员必须亲自上楼把快递递到你手里,服务当然好,但是万一我不在家,快递员就得一直等我,浪费了快递员的时间。所以解耦还是非常有必要的。
那么在这种场景下我们的秒杀就变成了:在我们下单之后,利用 Redis 去进行校验下单的结果,然后再通过队列把消息发送出去,然后启动一个线程去拿到这个消息,完成解耦,同时也加快我们的响应速度。
这里我们可以直接使用一些现成的(MQ)消息队列,如 Kafka、RabbitMQ 等,但是如果没有安装 MQ,我们也可以使用 Redis 提供的 MQ 方案。
# Redis 实现消息队列的三种方案
| List | PubSub | Stream | |
|---|---|---|---|
| 消息持久化 | 支持 | 不支持 | 支持 |
| 阻塞读取 | 支持 | 支持 | 支持 |
| 消息堆积处理 | 受限于内存空间,可以利用多消费者加快处理 | 受限于消费者缓冲区 | 受限于队列长度,可以利用消费者组提高消费速度,减少堆积 |
| 消息确认机制 | 不支持 | 不支持 | 支持 |
| 消息回溯 | 不支持 | 不支持 | 支持 |
# Stream 消息队列实现异步秒杀下单
步骤:
- 创建一个 Stream 类型的消息队列,名为
stream.orders - 修改之前的秒杀下单 Lua 脚本,在认定有抢购资格后,直接向
stream.orders中添加消息,内容包含 voucherId、userId、orderId - 项目启动时,开启一个线程任务,尝试获取
stream.orders中的消息,完成下单
具体实现步骤如下:
- 使用
RedisTemplate的opsForStream方法从队列中读取一条消息。 - 判断读取的消息是否为空,若为空,则继续循环等待下一条消息。
- 将读取的消息转换为
VoucherOrder对象。 - 执行下单逻辑,并将数据保存到数据库中。
- 手动 ACK,确认当前处理的消息已经被处理完成。