基于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);  
    }  
}