基于BlockingQueue的阻塞队列
当前的业务代码包含多个读写操作和分布式锁,读写操作都使用数据库,总体耗时较长。
每个线程都需要进行长时间操作,会导致同一时间的最大线程数,因此应该有更多的线程在进行秒杀判断而不是写数据库来提高并发。
实际上只要判断完库存就可以事实上说明抢到了,所以主线程只需要判断库存,写操作另外开启线程进行。
public Result createSeckillVoucherOrder(long voucherId) {
// 1.查询优惠券
SeckillVoucher seckillVoucher = seckillVoucherService.getById(voucherId);
// 2.判断是否存在
if (seckillVoucher == null) {
return Result.fail("优惠券不存在!");
}
// 3.判断是否在抢购时间内
if (seckillVoucher.getBeginTime().isAfter(LocalDateTime.now()) ||
seckillVoucher.getEndTime().isBefore(LocalDateTime.now())) {
return Result.fail("不在抢购时间内!");
}
// 4.判断库存是否充足
if (seckillVoucher.getStock() <= 0) {
return Result.fail("库存不足!");
}
// 5.获取分布式锁
long userId = UserHolder.getUser().getId();
String key = "order:" + userId;
RLock lock = redissonClient.getLock(key);
if (!lock.tryLock()) {
return Result.fail("不允许重复下单!");
}
// 6.创建优惠券订单
try {
IVoucherOrderService proxy = (IVoucherOrderService) AopContext.currentProxy();
return proxy.createVoucherOrder(voucherId);
} finally {
lock.unlock();
}
}因此,我们可以将判断库存和扣减的逻辑放到Redis中。
首先要进行预热,将库存提前保存到Redis中。
stringRedisTemplate.opsForValue()
.set(RedisConstants.SECKILL_STOCK_KEY + voucher.getId(), voucher.getStock().toString());使用一个string表示库存,一个set表示已经购买过的用户集合,利用Redis的单线程和Lua的原子性,就能做到扣减库存等实现。
local voucherId = ARGV[1] -- 秒杀的优惠券ID
local userId = ARGV[2] -- 秒杀的用户ID
local stockKey = 'seckill:stock:' .. voucherId -- 库存key
local userKey = 'seckill:user:' .. voucherId -- 用户key
local stock = redis.call('get', stockKey) -- 获取库存
if(stock == nil or tonumber(stock) <= 0) then
return 1
end
if(redis.call('sismember', userKey, userId) == 1) then
return 2
end
redis.call('incrby', stockKey, -1) -- 库存减1
redis.call('sadd', userKey, userId) -- 将用户添加到已购买的集合中
return 0然后执行Lua脚本
private final DefaultRedisScript<Long> seckillScript = new DefaultRedisScript<>();
@PostConstruct
public void init() {
seckillScript.setLocation(new org.springframework.core.io.ClassPathResource("lua/seckill.lua"));
seckillScript.setResultType(Long.class);
}
@Override
public Result createSeckillVoucherOrder(long voucherId) {
// 1.执行Lua脚本
long userId = UserHolder.getUser().getId();
long result = stringRedisTemplate.execute(seckillScript, Collections.emptyList(),
String.valueOf(voucherId), String.valueOf(userId));
// 2.判断是否成功
if (result != 0) {
if (result == 1) {
return Result.fail("库存不足!");
} else {
return Result.fail("不允许重复下单!");
}
}
}那么在库存扣减后,数据操作都在Redis中,而我们最终需要将这些订单存入数据库,因此先用一个阻塞队列(BlockingQueue)存起来。
private final BlockingQueue<VoucherOrder> orderTasks = new ArrayBlockingQueue<>(1024 * 1024); // 最大任务数
public Result createSeckillVoucherOrder(long voucherId) {
// 1.执行Lua脚本
long userId = UserHolder.getUser().getId();
long result = stringRedisTemplate.execute(seckillScript, Collections.emptyList(),
String.valueOf(voucherId), String.valueOf(userId));
// 2.判断是否成功
if (result != 0) {
if (result == 1) {
return Result.fail("库存不足!");
} else {
return Result.fail("不允许重复下单!");
}
}
// 3.保存到阻塞队列
// 3.1 创建订单
long orderId = idGenerator.generate("order");
VoucherOrder voucherOrder = new VoucherOrder();
voucherOrder.setId(orderId);
voucherOrder.setVoucherId(voucherId);
voucherOrder.setUserId(UserHolder.getUser().getId());
// 3.2 将订单加入阻塞队列
orderTasks.add(voucherOrder);
// 4.返回订单id
return Result.ok(orderId);
}存进阻塞队列后,就需要有线程来处理消息,因此创建一个线程来处理这些消息。
private final BlockingQueue<VoucherOrder> orderTasks = new ArrayBlockingQueue<>(1024 * 1024);
private final ExecutorService orderCreator = Executors.newSingleThreadExecutor();
@Autowired
private IVoucherOrderService proxy;
@PostConstruct
public void init() {
seckillScript.setLocation(new ClassPathResource("lua/seckill.lua"));
seckillScript.setResultType(Long.class);
orderCreator.submit(() -> {
while (true) {
try {
// 1.获取订单信息
VoucherOrder voucherOrder = orderTasks.take(); // 阻塞队列的take()在队列为空会让线程进入等待状态,所以可以while(true)
// 2.处理订单
handleVoucherOrder(voucherOrder);
} catch (Exception e) {
log.error(e.getMessage());
}
}
});
}
public void handleVoucherOrder(VoucherOrder voucherOrder) {
// 1.获取分布式锁
long userId = voucherOrder.getUserId();
String key = "order:" + userId;
RLock lock = redissonClient.getLock(key);
if (!lock.tryLock()) {
log.error("不允许重复下单!");
return;
}
// 2.更新库存,保存优惠券订单
try {
proxy.createVoucherOrder(voucherOrder); //自调用事务,需要注入自己或者通过AopContext.currentProxy()获取当前对象的代理对象
} finally {
lock.unlock();
}
}
@Transactional
public void createVoucherOrder(VoucherOrder voucherOrder) {
long voucherId = voucherOrder.getVoucherId();
// 1.更新库存
boolean success = seckillVoucherService.update()
.setSql("stock = stock - 1")
.eq("voucher_id", voucherId)
.gt("stock", 0)
.update();
if (!success) {
log.error("库存不足!");
return;
}
save(voucherOrder);
}基于Redis的消息队列
基于阻塞队列的异步秒杀实现了抢购和下单的解耦,但是存在几个问题:
- JVM内存有限
- 消息可能丢失(没有重试)
- 无法持久化
- 只有单一消费者
Redis5.0新加入的Stream类型可以实现轻量的消息队列,并且具有以下特点: - 多消费者组
- 消息回溯
- 消息确认
- 消息可持久化
创建消费者组
XGROUP CREATE <stream_key> <group_name> <ID> [MKSTREAM]
XGROUP CREATE stream.orders g1 0 MKSTREAM <stream_key>: Stream 的键名。<group_name>: 要创建的消费者组名称。<ID>: 指定从哪个消息开始读取。0: 从第一条消息开始(包含 Stream 中已有的历史消息)。$: 从最新消息开始(只读取命令执行后新进入 Stream 的消息)。ID: 从指定的特定消息 ID 之后开始。
MKSTREAM(可选): 如果指定的 Stream 键不存在,自动创建一个空 Stream。如果不加此参数且 Stream 不存在,命令会报错。
发布消息
因为使用Redis的Stream消息队列,所以在扣减库存的同时可以直接发布消息。
XADD <key> [NOMKSTREAM] [<MAXLEN | MINID> [~|=] <threshold>] <ID | *> <field> <value> [<field> <value> ...]
<key>: Stream 的键名。<ID | *>: 消息的 ID。- 使用
*:Redis 会自动生成一个唯一的 ID(推荐方式)。 - 手动指定:格式必须是毫秒时间戳-序列号(如
1620000000000-0),且必须比当前 Stream 中最大的 ID 还要大。
- 使用
<field> <value>: 消息的具体内容,以键值对的形式存储(类似 Hash 结构)。
local voucherId = ARGV[1] -- 秒杀的优惠券ID
local userId = ARGV[2] -- 秒杀的用户ID
local orderId = ARGV[3] -- 订单ID
local stockKey = 'seckill:stock:' .. voucherId -- 库存key
local userKey = 'seckill:user:' .. voucherId -- 用户key
local stock = redis.call('get', stockKey) -- 获取库存
local stock_num = tonumber(stock)
if (stock_num == nil or stock_num <= 0) then
return 1
end
if(redis.call('sismember', userKey, userId) == 1) then
return 2
end
redis.call('incrby', stockKey, -1) -- 库存减1
redis.call('sadd', userKey, userId) -- 将用户添加到已购买的集合中
-- 发布消息到订单队列
redis.call('xadd', 'stream.orders', '*', 'userId', userId, 'voucherId', voucherId, 'id', orderId)
return 0消费消息
将之前的阻塞队列消费改造成Redis消息队列消费即可,不同之处在于失败后要继续处理未确认的消息。
被XREADGROUP读取的消息会放入PEL(Pending Entries List,待处理条目列表)中,直到被XACK确认才会从PEL中移除。
XREADGROUP GROUP <group> <consumer> [COUNT <count>] [BLOCK <milliseconds>] [NOACK] STREAMS <key> [<key> ...] <ID> [<ID> ...]GROUP <group>: 指定所属的消费者组名。<consumer>: 消费者的名称(如果不存在,Redis 会在读取时自动创建该消费者)。COUNT <count>: 本次读取的消息上限。BLOCK <milliseconds>: 阻塞等待新消息的时间。NOACK: 如果设置此项,消息在被读取时会立即自动确认(不推荐,因为失去可靠性保证)。<ID>:>: 读取从未分发给任何消费者的新消息。这是最常用的模式。- 任意其他 ID(如
0): 读取已经分发给当前消费者、但尚未确认(ACK)的消息。这用于消费者崩溃重启后的“故障恢复”。
XACK <key> <group> <ID> [<ID> ...]<key>: Stream 的键名。<group>: 消费者组的名称。<ID>: 消息的 ID(可以一次性确认多个 ID)。
@PostConstruct
public void init() {
orderCreator.submit(() -> {
while (true) {
try {
// 1.阻塞消费消息
List<MapRecord<String, Object, Object>> list = stringRedisTemplate.opsForStream().read(
Consumer.from("g1", "c1"),
StreamReadOptions.empty().count(1).block(Duration.ofSeconds(5)),
StreamOffset.create(RedisConstants.ORDER_STREAM_KEY,
ReadOffset.lastConsumed())
);
if (list == null || list.isEmpty()) {
continue;
}
// 2.处理订单消息
MapRecord<String, Object, Object> record = list.get(0);
Map<Object, Object> map = record.getValue();
VoucherOrder voucherOrder = new VoucherOrder();
BeanUtil.fillBeanWithMap(map, voucherOrder, true);
handleVoucherOrder(voucherOrder);
// 3.确认消息
stringRedisTemplate.opsForStream()
.acknowledge(RedisConstants.ORDER_STREAM_KEY, "g1", record.getId());
} catch (Exception e) {
log.error("处理订单异常", e);
consumePendingMessage();
}
}
});
}
private void consumePendingMessage() {
while (true) {
try {
// 1.阻塞消费消息
List<MapRecord<String, Object, Object>> list = stringRedisTemplate.opsForStream().read(
Consumer.from("g1", "c1"),
StreamReadOptions.empty().count(1).block(Duration.ZERO),
StreamOffset.create(RedisConstants.ORDER_STREAM_KEY,
ReadOffset.from("0")) // 处理
);
if (list == null || list.isEmpty()) {
break;
}
// 2.处理订单消息
MapRecord<String, Object, Object> record = list.get(0);
Map<Object, Object> map = record.getValue();
VoucherOrder voucherOrder = new VoucherOrder();
BeanUtil.fillBeanWithMap(map, voucherOrder, true);
handleVoucherOrder(voucherOrder);
// 3.确认消息
stringRedisTemplate.opsForStream()
.acknowledge(RedisConstants.ORDER_STREAM_KEY, "g1", record.getId());
} catch (Exception e) {
log.error("处理订单异常", e);
}
}
} 基于RabbitMQ的消息队列
Redis的Stream虽然可以实现轻量级的消息队列,但相较于专业的消息队列仍有不足,这里使用RabbitMQ进行改造。
安装
首先在Docker中部署单节点RabbitMQ,然后安装SpringBoot官方的ampqstarter:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
配置
修改application.yml:
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest创建配置类,新建一个队列:
@Configuration
public class RabbitMQConfig {
@Bean
public Queue queue() {
return new Queue("seckill.queue", true);
}
}发布消息
然后在Service中注入RabbitTemplate:
@Autowired
private RabbitTemplate rabbitTemplate;
修改原有业务逻辑:
@Override
public Result createSeckillVoucherOrder(long voucherId) {
// 1.执行Lua脚本
long userId = UserHolder.getUser().getId();
long result = stringRedisTemplate.execute(seckillScript, Collections.emptyList(),
String.valueOf(voucherId), String.valueOf(userId));
// 2.判断是否成功
if (result != 0) {
if (result == 1) {
return Result.fail("库存不足!");
} else {
return Result.fail("不允许重复下单!");
}
}
// 3.创建订单
VoucherOrder voucherOrder = new VoucherOrder();
long orderId = idGenerator.generate("order");
voucherOrder.setId(orderId);
voucherOrder.setUserId(userId);
voucherOrder.setVoucherId(voucherId);
// 4.发送订单消息
try {
rabbitTemplate.convertAndSend("seckill.queue", voucherOrder);
} catch (Exception e) {
stringRedisTemplate.execute(seckillFailScript, Collections.emptyList(),
String.valueOf(voucherId), String.valueOf(userId));
return Result.fail("下单失败!");
}
// 5.返回订单id
return Result.ok(orderId);
}其中seckillFailScript:
local voucherId = ARGV[1] -- 秒杀的优惠券ID
local userId = ARGV[2] -- 秒杀的用户ID
local stockKey = 'seckill:stock:' .. voucherId -- 库存key
local userKey = 'seckill:user:' .. voucherId -- 用户key
redis.call('incrby', stockKey, 1) -- 库存加1
redis.call('srem', userKey, userId) -- 将用户从已购买的集合中移除
return 0消费消息
创建一个Listener类并使用@RabbitListener:
@Component
public class VoucherOrderListener{
@Autowired
private IVoucherOrderService voucherOrderService;
@RabbitListener(queues = "seckill.queue")
public void onMessage(VoucherOrder voucherOrder) {
voucherOrderService.handleVoucherOrder(voucherOrder);
}
}