【redis-day04-万字长文实现秒杀业务】
《redis-day04-优惠券秒杀》
0. 今日总结
- 实现了全局Id生成器工具类,通过拼接时间戳和redis序列号创建唯一订单ID
- 实现了优惠券秒杀下单,并为每个订单生成一个唯一ID
- 通过改进的乐观锁CAS法解决了多线程情况下的订单超卖问题
- 实现了一人一单逻辑,并通过加悲观锁解决了一人一单的线程问题
- 通过redisson实现的分布式锁解决了锁不可重入、不可重试、超时释放、主从一致性问题,并通过源码跟踪详细分析了实现原理
- 通过redis实现了秒杀优化,将查询优惠券、校验一人一单和操作数据库解耦,大大减少了因数据库阻塞导致的性能问题,并通过Lua脚本保证了查询优惠券、校验一人一单的事务原子性,整个过程通过阻塞队列实现。
- 了解了redis消息队列的几种实现方式,并基于redis的Stream+消费者组实现了消息队列
1. 全局唯一ID
1.1 全局ID生成器
- 问题

当用户抢购时,就会生成订单并保存到tb_voucherorder这张表中,而订单表如果使用数据库自增ID就存在一些问题:
- id的规律性太明显
- 受单表数据量的限制
-
全局ID生成器
全局ID生成器,是一种在分布式系统下用来生成全局唯一ID的工具,一般要满足下列特性:

为了增加ID的安全性,我们可以不直接使用Redis自增的数值,而是拼接一些其它信息:

ID的组成部分:
- 符号位:1bit,永远为0
- 时间戳:31bit,以秒为单位,可以使用69年
- 序列号:32bit,秒内的计数器,支持每秒产生2^32个不同ID
1.2 代码开发
-
RedisIdWorker

- 生成时间戳
- 通过LocalDateTime.now获取当前时间
- 通过now.toEpochSecond(ZoneOffset.UTC)将当前UTC作为时区偏移常量传入toEpochSecond,获得当前时间对应的long类型秒数
- 用当前时间减去开始时间
- 生成序列号
- 获取当前日期,精确到天,用yyyy:MM:dd格式接收,这样方便redis分层管理
- 对当日订单号进行递增,统计总订单量,该订单量即为序列号,由于前缀是当天,因此每天都会有专门的key
- 将时间戳和序列号拼接,作为id返回
- 生成时间戳
-
test测试

- 创建500线程的线程池ExecutorService es
CountDownLatch是一个同步工具,初始化计数为300。它的作用是让主线程等待所有300个子线程完成任务后再继续执行,确保统计的耗时准确。- Runnable task表示每个线程执行的操作:循环100次,每次调用
redisIdWorker.nextId("order")生成一个ID并打印 - latch.countDown()表示每个线程结束后计数器减一,保证 执行300次
- 循环300次,将任务提交给线程池执行。
latch.await()是关键的同步点,主线程会在这里阻塞,直到所有300个任务都完成(即计数器变为0)。- 最后计算并输出总耗时,用以评估性能。
1.3 小结

2. 实现优惠券秒杀下单
2.1 添加秒杀卷

表关系如下:
- tb voucher:优惠券的基本信息,优惠金额、使用规则等
- tb_seckill voucher:优惠券的库存、开始抢购时间,结束抢购时间。特价优惠券才需要填写这些信息

调用写好的接口添加秒杀卷

2.2 实现优惠券秒杀下单
下单时需要判断两点:
-
秒杀是否开始或结束,如果尚未开始或已经结束则无法下单
-
库存是否充足,不足则无法下单

-
service
@Autowired private ISeckillVoucherService seckillVoucherService; @Autowired private RedisIdWorker redisIdWorker; @Override @Transactional public Result seckillVoucher(Long voucherId) { //1.查询优惠券信息 //voucherId和seckillVoucherId相同,因此可以通过voucherId在seckillVoucher表里查 SeckillVoucher seckillVoucher = seckillVoucherService.getById(voucherId); //2.判断秒杀是否开始 if (seckillVoucher.getBeginTime().isAfter(LocalDateTime.now())) { return Result.fail("秒杀尚未开始!"); } //3.判断秒杀是否结束 if (seckillVoucher.getEndTime().isBefore(LocalDateTime.now())) { return Result.fail("秒杀已经结束!"); } //4.判断库存是否充足 if (seckillVoucher.getStock() < 1) { return Result.fail("库存不足!"); } //5.扣减库存 boolean success = seckillVoucherService.update() .setSql("stock = stock - 1") .eq("voucher_id", voucherId).update(); if (!success) { //扣减失败 return Result.fail("库存不足!"); } //6.创建订单 VoucherOrder voucherOrder = new VoucherOrder(); //6.1订单id long orderId = redisIdWorker.nextId("order"); voucherOrder.setId(orderId); //6.2用户id Long userId = UserHolder.getUser().getId(); voucherOrder.setUserId(userId); //6.3代金券id voucherOrder.setVoucherId(voucherId); save(voucherOrder); //7.返回订单id return Result.ok(orderId); }-
查询优惠券信息,注入seckillVoucherService,使用Mybatis-plus的getById方法查询,voucherId和seckillVoucherId相同,因此可以通过voucherId在seckillVoucher表里查
-
判断秒杀是否开始
-
判断秒杀是否结束
-
如果处于秒杀时间段内,开始进行业务逻辑
-
判断库存是否充足,如果不足,返回错误信息
-
如果充足,使用Mybatis-plus对库存-1


-
如果扣减失败,返回库存不足
-
-
如果扣减成功,表明秒杀下单成功,创建订单
- 构建订单id,使用自己编写的redisIdWorker构建器
- 构建用户id,通过thredLocal获取当前线程的用户id
- 构建代金券id
- 调用mybatis-plus的sava方法保存订单
-
将订单id返回给前端
-
3. 超卖问题
3.1 问题
- 正常情况

-
异常情况:线程交叉

3.2 解决方案:加锁

-
悲观锁实现较为简单
-
乐观锁的重点是如何判断之前查询得到的数据是否有被修改过,常见方式有两种
-
版本号法

-
对版本号法进行优化:CAS(Compare And Set)法

-
3.3 实现乐观锁CAS法
- 实现

只需要在更新库存时多加一条判断库存是否为查询到的库存即可
-
问题

实际上200个线程只卖出了41张票
-
分析
这是因为乐观锁认为,只要发生了修改,就要检查一致性问题,这样会导致票数在远远大于0时,只要出现了线程交叉的情况,就会购票失败,但是实际上此时票数还很多,完全可以购票成功。
3.4 乐观锁CAS法改进
-
不判断库存是否相同,而是直接判断库存是否大于0,只有大于0才能修改

按上述方式改进后能够解决超卖问题
4. 一人一单
4.1 需求
需求:修改秒杀业务,要求同一个优惠券,一个用户只能下一单

4.2 代码实现

在扣减库存前先判断该用户是否已经下过单
- 通过ThreadLocal查询用户Id
- 通过mybatis-plus根据用户id和优惠券id查询该用户和优惠券在数据库中的记录数
- 如果记录数大于0,说明用户已经购买过,则返回错误信息
- 如果记录为0,说明用户尚未下单,可以正常执行之后的逻辑
-
问题
和超卖问题类似,如果有大量线程涌入步骤5,此时都会判断出count=0从而进入之后的扣减库存、创建订单等一系列操作
-
解决方案
由于之前的超卖问题是修改count,对于修改类问题可以通过加乐观锁,而此时的问题是插入问题,不涉及修改,因此不能加乐观锁(乐观锁CAS法是在更新时添加一致性判断,而此时不存在数据,也就没法进行一致性判断)。
因此只能使用悲观锁解决
4.3 通过加悲观锁解决一人一单线程问题

-
将一人一单之后的逻辑抽象为方法,并对方法根据当前用户id加锁:
- 思路:可以直接在整个createVoucherOrder方法上加锁,但是没必要,因为这样的话所有用户进来都会加锁,所有线程就变成串行的了,因此可以只对用户加锁,只有同一用户的请求才会被加锁
- userId.toString().intern(),userId是一个Long类型的对象,因此每个userId都是不同的对象,因此要采用toString方法转化为字符串,但是toString转化后依旧是对象
因此要再调用intern方法,intern方法功能是去字符串常量池找到与该字符串相同的值的地址,并返回,这样的话只要userId对应字符串的值是相同的,那对应常量池的地址就是相同的,就能保证只给同一个用户加锁
-
为什么要把锁加在整个函数外?
- 由于整个函数加上了@Trancsactional注解,将函数交给了spring管理可能会出现已经释放锁,但是事务还没有提交的情况,这时候如果有别的线程进入了该函数就能获取对应的锁,这样又产生了线程安全问题
- 因此要在整个函数外面加锁,这样只有事务被提交过后别的线程才能正常的获得锁
-
为什么要获取代理对象,并通过代理对象的方式调用方法?
- Spring 事务的原理:当在方法上添加
@Transactional注解时,Spring 并不会直接增强原始的业务对象,而是会为其创建一个代理对象。当我们从 Spring 容器中注入 Service 或调用其方法时,实际上操作的是这个代理对象。代理对象会在目标方法执行前开启事务,在方法执行后提交或回滚事务 。 - 内部调用导致代理失效:在代码中,
synchronized代码块和createVoucherOrder方法同属于一个类。如果在synchronized块内直接使用this.createVoucherOrder(voucherId),这被称为“内部调用”。此时的this指向的是原始的、未被代理的对象本身,而不是 Spring 创建的代理对象。因此,这次调用会完全绕过代理对象,导致@Transactional注解失效,事务自然也就不会开启 。
- Spring 事务的原理:当在方法上添加
5. 分布式锁
5.1 一人一单的并发安全问题

5.1.1 IDEA配置
-
点进Edit Configurations

-
ctrl+D创建新程序

-
modify opitons添加add VM options


-
配置新端口

5.1.2 nginx配置
-
找到confi文件夹下的nginx.conf

-
修改文件配置

会监听api端口,然后反向代理到backend,backend进行负载均衡,有两个端口,分别是8081和8082
-
重新加载nginx

5.1.3 测试
-
发送多次api/voucher/list/1请求

-
两个服务器都收到了请求,实现了负载均衡


5.1.4 并发安全问题
-
用PostMan发送两个请求


-
库存减2,出现并发安全问题


-
原因
-
正常情况

-
线程交叉解决单线程

-
集群部署情况

用synchronized加锁是利用JVM内部的锁监视器实现,可能遇到的问题:在集群部署中每个服务器都会有一个专门的JVM,而synchronized加锁在每个服务器内都会维护一个专门的锁监视器,因此不同的JVM容器中获取锁不会受到其他JVM容器的影响,因此并发情况下依旧会出现安全问题
-
5.2 分布式锁
5.2.1 功能原理
-
关键:满足分布式系统或集群模式下多进程可见并且互斥的锁

5.2.2 不同分布式锁


5.2.3 基于Redis的分布式锁

为了保证获取锁和添加过期时间这两个操作的原子性(即要么都成功,要么都失败),可以用Set命令+各种参数实现,例如set lock therad1 ex 1000 nx,即设置过期时间,又实现了setnx的功能(只有不存在才能设置)
-
改进

5.2.4 实现分布式锁
-
SimpleRedisLock类

- 添加属性name:具体业务的名字,stringRedisTemplate:redis模板类,KEY_PREFIX:锁的前缀(避免硬编码)
- 通过构造方法为name和stringRedisTemplate初始化
- 编写tryLock,获取锁方法
- 通过
Thread.currentThread().getId()获取当前线程id - 通过
stringRedisTemplate.opsForValue().setIfAbsent(KEY_PREFIX + name, threadId + "", timeoutSec, TimeUnit.SECONDS);实现两个功能:对不存在的key设置值为:thredId,并设置过期时间,然后返回true。对已经存在的key,直接返回false
- 通过
- 编写unlock方法,解锁
-
ShopServiceImpl业务类实现过程

- 创建锁对象SimpleRedisLock,由于是对order进行相关操作,因此将业务名传入为"order:"再拼接当前用户Id,明确有一人一单业务
- 调用SimpleRedisLock的tryLock,传入过期时间,尝试获取锁
- 如果获取失败,说明已经有其他业务上锁,返回获取失败
- 如果获取锁成功
- 获取当前事务代理对象
- 调用当前代理对象的createVoucherOrder方法,该方法在第4.2节已经实现
5.3 redis锁可能存在业务阻塞导致的误删问题
5.3.1 问题和解决思路
问题

- 如果线程1在获取锁之后,执行业务时发生的阻塞,阻塞时间超过了锁的声明周期,则锁会自动释放
- 此时,线程2就能获取锁了,线程2 获取锁,执行业务
- 此时,线程1的业务执行完成了,出发了释放锁操作,但是此时线程2在获取这个锁,因此线程1会把线程2的锁释放,而线程2的业务还未结束
- 那么此时线程3就可以获取锁了,如果线程3获取锁,执行业务,那么线程2和线程3又出现了同时执行业务的情况,发生了并发安全问题
解决思路

- 在释放锁前根据存入锁的标志判断一下是否一致,只有一致才能释放锁
5.3.2 改进redis分布式锁
需求:修改之前的分布式锁实现,满足:
- 获取锁时存入线程标示(可以用UUID表示)
- 在释放锁时先获取锁中的线程标示,判断是否与当前线程标志一致
- 如果一致则释放锁
- 如果不一致则不释放锁

5.3.3 实现改进的redis分布式锁

- 使用UUID拼线程ID的方式作为redis中的锁对应的value
- 在释放锁时,获取当前线程的UUID和线程ID,并从redis中根据锁查找对应的value,判断两者是否相同
- 如果相同,说明该锁是自己加的,可以删除
- 如果不同,说明该锁是别人加的,不删除
5.4 redis锁可能存在的释放锁阻塞导致的误删问题
5.4.1 问题和解决思路
- 问题

- 线程1执行业务,判断锁表示OK,正当准备释放锁时,阻塞了(可能时JVM自动运行垃圾回收机制或其他原因)
- 线程1一致被阻塞,直到超时自动释放锁,此时线程2尝试获取锁成功
- 线程2执行业务时,线程1阻塞完毕,执行释放锁业务,由于已经判断OK,因此会释放当前被线程2获取的锁,而线程2的业务还未结束
- 此时线程3获取锁成功,就出现了线程2和线程3的并行问题
- 解决思路:保证判断和删除操作的一致性
5.4.2 Lua脚本
Redis提供了Lua脚本功能,在一个脚本中编写多条Redis命令,确保多条命令执行时的原子性。Lua是一种编程语言,它的基本语法大家可以参考网站:https://www.runoob.com/lua/lua-tutorial.html



5.4.3 用Lua脚本改进redis分布式锁

5.4.4 在java中调用Lua脚本
- 通过将判断和删除写到lua脚本里,再通过java一次性调用整个脚本实现删除和判断操作的原子性

-
在resources中编写lua脚本

- 调用redis.call,执行get keys[1]方法,得到key[1]中存储的value
- 将获取到的value即UUID+userId与传入的参数ARGV[1]进行比较
- 如果相同,释放锁
- 否则不执行任何操作(不释放锁)


- 声明脚本对象UNLOCK_SCRIPT
- 使用静态代码块对UNLOCK_SCRIPT赋初值
- 创建UNLOCK_SCRIPT对象
- 设置脚本位置为ClassPathResource(class路径下找)下的"unlock.lua"
- 设置返回值类型为Long
- 编写unlock方法,调用stringRedisTemplate的execute方法,执行脚本
- 第一个参数为脚本本身,传入脚本对象UNLOCK_SCRIPT
- 第二个参数为存储key的集合,由于此处为锁的key,唯一,因此用Collections的singletonList创建集合,值为KEY_PREFIX+name
- 传入值,即为UUID+userId
5.5 分布式锁存在的问题及Redisson工具包
5.5.1 问题

5.5.2 Redisson介绍
Redisson是一个在Redis的基础上实现的]ava驻内存数据网格(In-MemoryData Grid)。它不仅提供了一系列的分布式的Java常用对象,还提供了许多分布式服务,其中就包含了各种分布式锁的实现。

5.5.3 Redisson入门


-
配置类

-
改写业务


将获取锁的部分用RedissonClient对象中的方法实现
5.6 Redisson原理
5.6.1 可重入锁原理
-
原本基于lua脚本编写的代码的问题

- 以左边的代码为例:发生死锁
- 方法1获取锁成功后会调用方法2
- 方法2会尝试获取锁,但是由于方法1还没有释放锁,因此无法获取
- 而方法1由于没法执行方法2,又迟迟无法结束,从而导致了死锁
- 以左边的代码为例:发生死锁
-
实现原理
为了实现通过redis实现的分布式锁的可重入性,采用的方法是采用hash而非string类型存储,field存储值,而value存储重入次数,每次获取锁都+1,每次释放锁都-1,当减为0时才真正的释放锁。
redisson在加锁解锁时正式根据下述流程图严谨实现的

5.6.2 可重试原理
package com.hmdp;
import lombok.extern.slf4j.Slf4j;
import org.junit.jupiter.api.Test;
import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import java.util.concurrent.TimeUnit;
@Slf4j
@SpringBootTest
public class RedissonTest {
@Autowired
private RedissonClient redissonClient;
private RLock lock;
void setUp(){
lock = redissonClient.getLock("order");
}
@Test
void method1() throws InterruptedException {
boolean isLock = lock.tryLock(1L, TimeUnit.SECONDS);
if (!isLock) {
log.error("获取锁失败 ... 1");
return;
}
try {
log.info("获取锁成功 ... 1");
method2();
log.info("开始执行业务 ... 1");
} finally {
log.warn("准备释放锁 ... 1");
lock.unlock();
}
}
void method2(){
boolean isLock = lock.tryLock();
if (!isLock) {
log.error("获取锁失败 ... 2");
return;
}
try {
log.info("获取锁成功 ... 2");
method2();
log.info("开始执行业务 ... 2");
} finally {
log.warn("准备释放锁 ... 2");
lock.unlock();
}
}
}
-
调用RLock对象lock的tryLock方法,传入等待时间和时间单位

-
tryLock逻辑

- 时间单位是毫秒
- 获取当前时间
- 获取当前线程ID
- 进入获取锁tryAcquire逻辑
-
tryAcquire逻辑

-
tryAcquireAsync逻辑:获取锁剩余时间
private <T> RFuture<Long> tryAcquireAsync(long waitTime, long leaseTime, TimeUnit unit, long threadId) { if (leaseTime != -1L) { return this.<Long>tryLockInnerAsync(waitTime, leaseTime, unit, threadId, RedisCommands.EVAL_LONG); } else { RFuture<Long> ttlRemainingFuture = this.<Long>tryLockInnerAsync(waitTime, this.commandExecutor.getConnectionManager().getCfg().getLockWatchdogTimeout(), TimeUnit.MILLISECONDS, threadId, RedisCommands.EVAL_LONG); ttlRemainingFuture.onComplete((ttlRemaining, e) -> { if (e == null) { if (ttlRemaining == null) { this.scheduleExpirationRenewal(threadId); } } }); return ttlRemainingFuture; } }-
判断释放时间是否是-1,如果不是-1,进入tryLockInnerAsync逻辑,传入的leaseTime即为设定的时间。如果是-1,也进入tryLockInnerAsync逻辑只不过leaseTime位置的参数通过
this.commandExecutor.getConnectionManager().getCfg().getLockWatchdogTimeout()获取,是一个默认值
默认超时时间为30s -
tryLockInnerAsync逻辑
<T> RFuture<T> tryLockInnerAsync(long waitTime, long leaseTime, TimeUnit unit, long threadId, RedisStrictCommand<T> command) { this.internalLockLeaseTime = unit.toMillis(leaseTime); return this.<T>evalWriteAsync(this.getName(), LongCodec.INSTANCE, command, "if (redis.call('exists', KEYS[1]) == 0) then redis.call('hincrby', KEYS[1], ARGV[2], 1); redis.call('pexpire', KEYS[1], ARGV[1]); return nil; end; if (redis.call('hexists', KEYS[1], ARGV[2]) == 1) then redis.call('hincrby', KEYS[1], ARGV[2], 1); redis.call('pexpire', KEYS[1], ARGV[1]); return nil; end; return redis.call('pttl', KEYS[1]);", Collections.singletonList(this.getName()), this.internalLockLeaseTime, this.getLockName(threadId)); }- 将锁释放时间记录到本地成员变量
internalLockLeaseTime、 - 执行Lua脚本
- 如果获取锁成功,返回nil
- 如果获取锁失败,返回剩余有效期,单位是毫秒
- 将锁释放时间记录到本地成员变量
-
RFuture<Long> ttlRemainingFuture接受脚本返回的剩余有效期
- 由于tryLockInnerAsync是一个异步函数,因此执行完那没拿到值不一定,因此用RFuture接收
-
将 ttlRemainingFuture返回
-
-
-
tryLock中的ttl接受返回的剩余有效期
-
判断有效期是否为null

- 如果剩余有效期为null,说明获取锁成功,返回true
- 否则说明剩余有效期不为null,进入重试逻辑
-
重试

-
设置的等待时间减去系统当前时间再减去current

-
如果剩余等待时间小于0,表示获取锁失败,返回false

-
如果剩余等待时间大于0

-
记录当前时间
-
对线程id进行订阅:订阅别人释放锁的信号。publish发布消息通知

-
等待信号,等待时间为time,如果等到最大时间还没等到就不等了,返回false,加上!变成true,那么进入if逻辑
- 执行unsubscribe,释放锁(因为已经等待超时)
-
如果没进if,说明还没到最大等待时间

-
再计算一次剩余等待时间

-
如果剩余时间大于0,进入循环,开始尝试重新获取
-
获取锁,返回剩余有效期,如果为空表示获取成功

-
如果不为空,表示获取失败,计算剩余等待时间

-
如果剩余等待时间小于0,放弃获取,返回false

-
如果剩余时间大于0,且剩余有效期小于剩余等待时间,则等待时间传入为剩余有效期,如果剩余有效期大于剩余等待时间,那么等待时间传入为剩余等待时间。

-
-
再次计算剩余等待时间

如果剩余时间大于0,则再次尝试,否则退出循环,返回错误,表示获取锁失败

-
-
-
-
-
5.6.3 超时续约实现
-
为什么要实现超时续约?
redis为了解决服务器崩溃等以外原因导致的缓存未被释放,会给缓存设一个过期时间,但是如果业务执行时间过长,超过了这个过期时间,可能会出现并发问题。
因此要在业务执行时通过看门狗机制实现超时续约,不断更新锁的有效期,直到业务完成,但是锁本身又存在有效期,所以也能应对服务器崩溃导致的内存不释放。
-
获取锁和释放锁逻辑


-
tryAcquireAsync逻辑实现

-
回调成功,获取剩余有效期之后,进入上面的逻辑
-
如果异常不为空,表示抛出了异常
-
如果异常为空,表示正常执行
-
如果剩余有效期为null,表示获取成功
-
执行
scheduleExpirationRenewal表示更新过期时间
-
通过putIfAbsent进行赋值,key为当前enteyName每个实例(order相关锁是一个key,voucher相关锁是一个key)独一份。如果是第一次则将一个全新的entry放入对应的锁,此时没有oldEntry,返回null。如果不是第一次重入,则不会写入,而是返回oldEntry
- 如果oldEntry不为空(不是第一次),则把当前线程Id加入oldEntry。由于不同线程不会获得同一把锁,因此这的逻辑一定是同一线程重入
- 如果oldEntry为空(是第一次),加入,并执行renewExpiration刷新缓存时间
-
renewExpiration逻辑

-
先获得entry,如果entry不为空,执行定时任务Timeout task


任务在delay到期以后才执行,到期时间为10s

- 定时任务负责重置有效期
- 递归调用本身,直到捕获释放锁信号
-
-
-
-
-
5.6.4 主从一致性原理(Multi lock)
-
问题

如果主节点子节点还没同步,此时主节点宕机,那么子节点没有保存在主节点保存的锁,就出现了主从不一致
-
解决方案:Multi lock

5.6.5 Redisson总结
- **可重入:**利用hash结构记录线程id和重入次数
- **可重试:**利用信号量和PubSub功能实现等待、唤醒,获取锁失败的重试机制
- **超时续约:**利用watchDog,每隔一段时间,重置超时时间
- **主从一致性:**要求客户端必须向多个独立的 Redis 节点(而不仅是主从节点)同时获取锁,只有全部成功才算真正持有锁。
- 不可重入锁
- 原理:利用setnx的互斥性;利用ex避免死锁;释放锁时判断线程标示
- 缺陷:不可重入、无法重试、锁超时失效
- 可重入redis分布式锁
- 原理:利用hash结构,记录线程标示和重入次数;利用watchDog延续锁时间;利用信号量控制锁重试等待
- 缺陷:redis宕机引起锁失效问题
- Redisson的multiLock
- 原理:多个独立的Redis节点,必须在所有节点都获取重入锁,才算获取锁成功
- 缺陷:运维成本高、实现复杂
6. Redis优化秒杀
6.1 问题和优化方案
- 问题:整个业务流程串行,并且频繁操作数据库导致性能较差

- 优化方案:将查询优惠券和校验一人一单放到redis中实现,将id返回给用户,并创建一个独立线程进行后续对数据库的操作

-
判断秒杀库存使用的redis数据结构:String类型

-
一人一单使用的redis数据结构:set集合

-
业务流程

6.2 基于Redis实现异步秒杀

6.2.1 添加秒杀券

6.2.2 基于Lua脚本,实现一人一单

6.2.3 在java中调用lua脚本
private static final DefaultRedisScript<Long> SECKILL_SCRIPT;
static {
SECKILL_SCRIPT = new DefaultRedisScript<>();
SECKILL_SCRIPT.setLocation(new ClassPathResource("seckill.lua"));
SECKILL_SCRIPT.setResultType(Long.class);
}
@Override
public Result seckillVoucher(Long voucherId) {
//获取用户
Long userId = UserHolder.getUser().getId();
//1.执行lua脚本
Long result = stringRedisTemplate.execute(
SECKILL_SCRIPT,
Collections.emptyList(),
voucherId.toString(), userId.toString()
);
//2.判断结果是否为0
int r = result.intValue();
if (r != 0) {
//2.1 不为0,代表没有购买资格
return Result.fail(r == 1 ? "库存不足" : "不能重复下单");
}
//2.2 为0,有购买资格,吧下单信息保存到阻塞队列
long orderId = redisIdWorker.nextId("order");
VoucherOrder voucherOrder = new VoucherOrder();
//2.3 设置orderId
voucherOrder.setId(orderId);
//2.4 用户id
voucherOrder.setUserId(userId);
//2.5 代金券id
voucherOrder.setVoucherId(voucherId);
//2.6 创建阻塞队列
orderTasks.add(voucherOrder);
//3.获取代理对象
proxy = (IVoucherOrderService) AopContext.currentProxy();
//4. 返回订单id
return Result.ok(orderId);
}
- 初始化Script
- 调用execute方法执行lua脚本
- 第一个参数是脚本
- 第二个参数是空列表,因为lua里不需要用到Key
- 第三个表示voucherId
- 判断结果是否为0,不是0表示没有购买资格,根据返回值返回对应的错误信息
- 为0表示有购买资格,调用自己写的redisIdWorker.nextId,给他一个唯一的orderId
6.3.4 将秒杀卷添加到阻塞队列
- 说明:由于VoucherOrderServiceImpl被加上了@service注解,因此其被注册为Spring容器的一个@bean对象,如果线程1,2,3都注入了该对象,其实注入的是同一个对象。儿orderTasks是该对象的一个变量,因此不同线程实际上是在操作同一个阻塞队列,而阻塞队列本身保证了线程安全,因此将voucherOrder添加到阻塞队列时,不会产生线程安全问题


- 创建阻塞队列
- 设置秒杀卷信息
- 将秒杀卷添加到阻塞队列
6.3.5 开启线程,从阻塞队列获取信息

-
创建线程池,该线程池为单线程:线程池内只有一个线程
-
init()方法由于加了PostConstrut注解,在spring boot启动时就会执行,提交一个VoucherOrderHanlder任务,因此在spring boot启动时,就会开始执行VoucherOrderHanlder里的run
-
run方法内,会循环的执行以下操作
- 调用oderTasks.take获取消息队列中的数据(6.3.4已经分析过,不同用户(不同线程)操作的是同一个阻塞队列,并且不会出现线程问题,因为阻塞队列本身已经实现了线程安全)
- 获取完阻塞队列中的数据后,调用handleVoucherOrder方法处理订单数据
-
handleVoucherOrder方法会获取用户Id,然后根据该id尝试获取锁对象,如果获取成功,通过代理对象proxy来调用创建订单方法,并且将代理对象创建在VoucherOrderServiceImpl的成员变量部分,然后在将数据添加到阻塞队列的过程中对代理对象进行初始化,这样就能在handleVoucherOrder中使用代理对象调用createVoucherOrder方法从而避免通过this.createVoucherOrder调用createVoucherOrder方法会出现的跳过@Transactional注解的问题


6.3 几个关键问题和解答
6.3.1 为不同用户开启的线程和代码中创建的线程有什么关系和区别?
- 代码中创建的是一个单线程的线程池,并在添加了@PostConsrtuct注解的init方法中提交,因此,Spring boot在启动时就会将传入了new VoucherOrderHandler()参数的SECKILL_ORDER_EXECUTOR线程池中的单线程启动。
- 用户开启的线程和上述单线程没有关系,例如用户1,2,3同时发起请求,系统为其分配了三个线程,这三个线程和上述线程独立,上述线程是随着spring boot启动而被自动开启的
6.3.2 为何不同的用户能操作同一个阻塞队列?
- 该阻塞队列是VoucherOrderServiceImpl的一个成员变量,而VoucherOrderServiceImpl又加上了@Service注解被注册为了Spring bean,是容器中唯一的VoucherOrderServiceImpl bean,不同的用户(假设用户1,2,3分别分配了线程1,2,3)在注入VoucherOrderServiceImpl时,实际上注入的是同一个bean,因此其中的成员变量orderTasks也是同一个阻塞队列
6.3.3 写入阻塞队列和从阻塞队列获取数据的线程是怎么执行的?会不会冲突?
假设用户1,2,3分别分配了线程1,2,3调用秒杀接口
- 写入阻塞队列操作是在seckillVoucher方法中,只要用户调用了秒杀接口的该方法,就会将数据写入阻塞队列。该过程是线程1,2,3分别执行的,由于阻塞队列本身实现了线程安全,因此不会冲突
- 获取数据的线程就如6.3.1所说,是随着spring boot启动开启的独立线程,不断地,循环的从阻塞队列中获取数据,并尝试创建订单
6.3.4 为什么要用代理对象调用创建订单的方法?
- 为了实现数据库操作的事务一致性,要加上@Transactional注解,该注解的实现原理是spring Aop,对加上Transactional注解的方法进行增强,如果不用代理对象调用创建订单的方法,而是用this调用,则是内部调用,指向的是原始的,未被代理的创建订单方法本身,则Transactional注解失效。Spring在创建Bean的过程中,会动态地生成一个代理对象来“包裹”您的原始对象,因此调用该动态对象的创建订单方法,则能识别到Transactional注解
6.3.5 代理对象是怎么初始化的?
- 和6.3.2类似,proxy首先被创建为VoucherOrderServiceImpl的一个成员变量但未初始化,由于VoucherOrderServiceImpl被注册为了唯一bean,因此不同的用户能访问同一个proxy,只要用户调用了秒杀接口seckillVoucher,就会执行初始化proxy的操作。
- 后一个请求的执行可能会覆盖前一个请求设置的
proxy值。但在上述业务场景下,通常不会造成问题,因为代理对象本质上是同一个。
6.3.6 会不会出现代理对象还没被初始化就通过其调用createVocherOrder方法的情况?
- 不会,因为线程池的唯一线程在执行时,会首先尝试从阻塞队列中获取数据,获取失败会一直等待,而只有至少有一个用户执行了将数据写到阻塞队列操作之后,阻塞队列才会有数据,而只要有用户执行了这个操作,那代理对象就会被初始化
6.3.7 为什么不能在创建代理对象时就初始化?
-
因为spring bean注册过程有严格的顺序
-
属性填充
-
前置处理
-
@PostConstruct方法执行
-
后置处理
-
生成代理对象
如果在第一步就尝试初始化,此时代理对象尚未生成,无法成功初始化。
-
7. Redis消息队列实现异步秒杀
7.1 Java内置阻塞队列的问题
- JVM内存有限,全放阻塞队列可能会导致内存不足
- 由于JVM不是独立的服务,因此服务器如果宕机,阻塞队列中的数据可能丢失
7.2 消息队列
消息队列(Message Queue),字面意思就是存放消息的队列。最简单的消息队列模型包括3个角色:
- 消息队列:存储和管理消息,也被称为消息代理(Message Broker)
- 生产者:发送消息到消息队列
- 消费者:从消息队列获取消息并处理消息

Redis提供了三种不同的方式来实现消息队列:
- list结构:基于List结构模拟消息队列
- PubSub:基本的点对点消息模型
- Stream:比较完善的消息队列模型
7.2.1 基于List结构模拟消息队列


7.2.1 基于PubSub的消息队列
Pubsub(发布订阅)是Redis2.0版本引入的消息传递模型。顾名思义,消费者可以订阅一个或多个channel,生产者向对应channel发送消息后,所有订阅者都能收到相关消息。
-
SUBSCRlBE channel[channel]:订阅一个或多个频道
-
PUBLlSH channelmsg :向一个频道发送消息
-
PSUBSCRIBE pattern[pattern]:订阅与pattern格式匹配的所有频道



7.3 基于Redis的Stream实现消息队列
7.3.1 基于Stream的消息队列
Stream 是 Redis 5.0 引入的一种新数据类型,可以实现一个功能非常完善的消息队列。
发送消息的命令:

例如:

读取消息的方式之一:XREAD

例如,使用XREAD读取第一个消息:

XREAD阻塞方式,读取最新的消息:

在业务开发中,我们可以循环的调用XREAD阻塞方式来查询最新消息,从而实现持续监听队列的效果,伪代码如下

注意:当我们指定起始ID为$时,代表读取最新的消息,如果我们处理一条消息的过程中,又有超过1条以上的消息到达队列,则下次获取时也只能获取到最新的一条,会出现漏读消息的问题
STREAM类型消息队列的XREAD命令特点:
- 消息可回溯
- 一个消息可以被多个消费者读取
- 可以阻塞读取
- 有消息漏读的风险
7.3.2 基于Stream的消息队列-消费者组
消费者组(Consumer Group):将多个消费者划分到一个组中,监听同一个队列。具备下列特点:

创建消费者组:
key:队列名称
groupName:消费者组名称
ID:起始ID标示,$代表队列中最后一个消息,0则代表队列中第一个消息
MKSTREAM:队列不存在时自动创建队列
其它常见命令:
删除指定的消费者组
XGROUP DESTORY key groupName
给指定的消费者组添加消费者
XGROUP CREATECONSUMER key groupname consumername
删除消费者组中的指定消费者
XGROUP DELCONSUMER key groupname consumername
从消费者组读取消息:
XREADGROUP GROUP group consumer [COUNT count] [BLOCK milliseconds] [NOACK] STREAMS key [key ...] ID [ID ...]
- group:消费组名称
- consumer:消费者名称,如果消费者不存在,会自动创建一个消费者
- count:本次查询的最大数量
- BLOCK milliseconds:当没有消息时最长等待时间
- NOACK:无需手动ACK,获取到消息后自动确认
- STREAMS key:指定队列名称
- ID:获取消息的起始ID:
- “>”:从下一个未消费的消息开始
- 其它:根据指定id从pending-list中获取已消费但未确认的消息,例如0,是从pending-list中的第一个消息开始
消费者监听消息的基本思路:
STREAM类型消息队列的XREADGROUP命令特点:
- 消息可回溯
- 可以多消费者争抢消息,加快消费速度
- 可以阻塞读取
- 没有消息漏读的风险
- 有消息确认机制,保证消息至少被消费一次
最后我们来个小对比

Redis Stream 的 消费者组(Consumer Group) 是一个强大的功能,它允许多个消费者协同处理同一个消息队列中的消息,不仅提高了处理效率,还确保了消息不会在消费过程中丢失。下面我用一个简单的类比和梳理,帮你理解它的核心机制。
🏷️ 先了解基本概念
可以把 Redis Stream 想象成一个不断增长的快递流水线,每个快递包裹就是一个消息,带有唯一的 ID(如
1703235642574-0,通常由时间戳和序列号组成)。而 消费者组(Consumer Group) 就像是承包了这条流水线某个环节的一个快递团队。这个团队里有多个快递员(消费者),他们共同从流水线上取包裹去派送。
🔑 核心运行机制
消费者组之所以能可靠地工作,依赖于以下几个关键机制:
消息分发规则:竞争消费
团队里的快递员是竞争关系。同一个包裹一旦被一个快递员取走,其他快递员就不会再拿到这个包裹了。这保证了同一条消息不会被组内的多个消费者重复处理。
进度记录:
last_delivered_id团队里有个记录员,始终记着最后一个被取走的包裹的编号。这个编号就是
last_delivered_id。每当有快递员取走一个包裹,记录员就会把这个编号往前移动,确保大家知道下一个该处理哪个包裹。确保消息不丢失:Pending List (PEL) 与 ACK 机制
这是消费者组最核心的可靠性保障。当快递员拿到一个包裹后,这个包裹并不会立刻从流水线上消失,而是会进入一个“已取件未签收”的清单(Pending List, PEL)。
- 只有快递员派送成功,并明确回复“已签收”(即发送
XACK命令)后,这个包裹信息才会从 Pending List 中移除。- 如果某个快递员中途发生意外(比如程序崩溃),这个包裹会一直留在 Pending List 中。过一段时间,团队负责人可以通过
XPENDING命令查看这些“滞留”的包裹,并利用XCLAIM命令将它们重新分配给其他健康的快递员去处理。这就保证了即使消费者出故障,消息也不会丢失。
7.3.3 基于Redis的Stream结构作为消息队列,实现异步秒杀下单

-
改进:之前Lua脚本仅仅判断是否有抢购资格,再利用java代码往JVM阻塞队列中添加数据,而现在是往redis中发消息,因此可以一次性在Lua脚本中实现
-
lua脚本

- 通过redis一次性实现判断和保存消息,因此要传入订单id
- 进行校验的逻辑没有改变,依旧先判断库存是否充足以及用户是否下单
- 如果库存不足或用户已下单,都中止脚本,返回
- 如果库存充足且用户未下单,则执行之后的命令
- 口库存
- 下单
- 将订单数据保存到消息队列中
-
seckillVoucher

- 由于lua脚本需要订单id,因此在执行脚本时要传入第三个参数:订单id
- 由于不再需要将数据存到阻塞队列中,因此删除阻塞队列相关代码(通过Lua脚本将消息保存到消息队列)
- 判断是否下单成功,成功则获取代理对象
-
线程执行操作

-
线程不断的获取消息队列中的数据,等待时间为两秒
-
如果获取失败,说明消息队列中没有数据,则跳过本次循环
-
如果获取成功,则解析订单信息
-
list.get(0)获取消息队列中最早的一条信息(包含消息id和消息内容)

-
record.getValue()获得消息内容
-
调用BeanUtil将消息内容转化为VoucherOrder形式
-
确认下单
-
-
下单完成后ACK确认,将消息从peding list中移除
-
-
-
handlePendingList

- 如果抛出异常,则会处理PedingLsit
- 获取消息队列中的订单信息
ReadOffset.from("0"):这是关键区别点。与正常消费使用>不同,这里的"0"表示读取所有已领取但未确认的消息,即Pending List中的消息。count(1):一次只读取1条消息,避免一次性处理过多消息导致内存压力
- 判断是否获取成功,如果获取失败,说明pending list中没有未确认信息,直接跳出即可
- 如果有数据,则解析信息,保存到数据库,ACK确认
- 如果线程被意外中断,休眠一定时间后抛出异常
更多推荐




所有评论(0)