《redis-day04-优惠券秒杀》

0. 今日总结

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

1. 全局唯一ID

1.1 全局ID生成器

  • 问题

image-20251219084908719

当用户抢购时,就会生成订单并保存到tb_voucherorder这张表中,而订单表如果使用数据库自增ID就存在一些问题:

  1. id的规律性太明显
  2. 受单表数据量的限制
  • 全局ID生成器

    全局ID生成器,是一种在分布式系统下用来生成全局唯一ID的工具,一般要满足下列特性:

    外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

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

    image-20251219085835740

    ID的组成部分:

    1. 符号位:1bit,永远为0
    2. 时间戳:31bit,以秒为单位,可以使用69年
    3. 序列号:32bit,秒内的计数器,支持每秒产生2^32个不同ID

1.2 代码开发

  • RedisIdWorker

    外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

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

    外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

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

1.3 小结

image-20251219092608670

2. 实现优惠券秒杀下单

2.1 添加秒杀卷

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

表关系如下:

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

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

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

image-20251219100831916

2.2 实现优惠券秒杀下单

下单时需要判断两点:

  1. 秒杀是否开始或结束,如果尚未开始或已经结束则无法下单

  2. 库存是否充足,不足则无法下单

    image-20251219101341081

  • 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);
        }
    
    1. 查询优惠券信息,注入seckillVoucherService,使用Mybatis-plus的getById方法查询,voucherId和seckillVoucherId相同,因此可以通过voucherId在seckillVoucher表里查

    2. 判断秒杀是否开始

    3. 判断秒杀是否结束

    4. 如果处于秒杀时间段内,开始进行业务逻辑

      1. 判断库存是否充足,如果不足,返回错误信息

      2. 如果充足,使用Mybatis-plus对库存-1

        image-20251219105029707

        image-20251219105012451

      3. 如果扣减失败,返回库存不足

    5. 如果扣减成功,表明秒杀下单成功,创建订单

      1. 构建订单id,使用自己编写的redisIdWorker构建器
      2. 构建用户id,通过thredLocal获取当前线程的用户id
      3. 构建代金券id
      4. 调用mybatis-plus的sava方法保存订单
    6. 将订单id返回给前端

3. 超卖问题

3.1 问题

  • 正常情况

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

  • 异常情况:线程交叉

    image-20251219134206492

3.2 解决方案:加锁

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

  • 悲观锁实现较为简单

  • 乐观锁的重点是如何判断之前查询得到的数据是否有被修改过,常见方式有两种

    1. 版本号法

      image-20251219140510358

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

      image-20251219140643649

3.3 实现乐观锁CAS法

  • 实现

image-20251219141526671

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

  • 问题

    外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

    实际上200个线程只卖出了41张票

  • 分析

    这是因为乐观锁认为,只要发生了修改,就要检查一致性问题,这样会导致票数在远远大于0时,只要出现了线程交叉的情况,就会购票失败,但是实际上此时票数还很多,完全可以购票成功。

3.4 乐观锁CAS法改进

  • 不判断库存是否相同,而是直接判断库存是否大于0,只有大于0才能修改

    image-20251219142051281

    按上述方式改进后能够解决超卖问题

4. 一人一单

4.1 需求

需求:修改秒杀业务,要求同一个优惠券,一个用户只能下一单

image-20251219143518620

4.2 代码实现

image-20251219145456263

在扣减库存前先判断该用户是否已经下过单

  1. 通过ThreadLocal查询用户Id
  2. 通过mybatis-plus根据用户id和优惠券id查询该用户和优惠券在数据库中的记录数
  3. 如果记录数大于0,说明用户已经购买过,则返回错误信息
  4. 如果记录为0,说明用户尚未下单,可以正常执行之后的逻辑
  • 问题

    和超卖问题类似,如果有大量线程涌入步骤5,此时都会判断出count=0从而进入之后的扣减库存、创建订单等一系列操作

  • 解决方案

    由于之前的超卖问题是修改count,对于修改类问题可以通过加乐观锁,而此时的问题是插入问题,不涉及修改,因此不能加乐观锁(乐观锁CAS法是在更新时添加一致性判断,而此时不存在数据,也就没法进行一致性判断)。

    因此只能使用悲观锁解决

4.3 通过加悲观锁解决一人一单线程问题

image-20251219152849190

  1. 将一人一单之后的逻辑抽象为方法,并对方法根据当前用户id加锁:

    1. 思路:可以直接在整个createVoucherOrder方法上加锁,但是没必要,因为这样的话所有用户进来都会加锁,所有线程就变成串行的了,因此可以只对用户加锁,只有同一用户的请求才会被加锁
    2. userId.toString().intern(),userId是一个Long类型的对象,因此每个userId都是不同的对象,因此要采用toString方法转化为字符串,但是toString转化后依旧是对象
      因此要再调用intern方法,intern方法功能是去字符串常量池找到与该字符串相同的值的地址,并返回,这样的话只要userId对应字符串的值是相同的,那对应常量池的地址就是相同的,就能保证只给同一个用户加锁
  2. 为什么要把锁加在整个函数外?

    1. 由于整个函数加上了@Trancsactional注解,将函数交给了spring管理可能会出现已经释放锁,但是事务还没有提交的情况,这时候如果有别的线程进入了该函数就能获取对应的锁,这样又产生了线程安全问题
    2. 因此要在整个函数外面加锁,这样只有事务被提交过后别的线程才能正常的获得锁
  3. 为什么要获取代理对象,并通过代理对象的方式调用方法?

    1. Spring 事务的原理:当在方法上添加 @Transactional注解时,Spring 并不会直接增强原始的业务对象,而是会为其创建一个代理对象。当我们从 Spring 容器中注入 Service 或调用其方法时,实际上操作的是这个代理对象。代理对象会在目标方法执行前开启事务,在方法执行后提交或回滚事务 。
    2. 内部调用导致代理失效:在代码中,synchronized代码块和 createVoucherOrder方法同属于一个类。如果在 synchronized块内直接使用 this.createVoucherOrder(voucherId),这被称为“内部调用”。此时的 this指向的是原始的、未被代理的对象本身,而不是 Spring 创建的代理对象。因此,这次调用会完全绕过代理对象,导致 @Transactional注解失效,事务自然也就不会开启 。

5. 分布式锁

5.1 一人一单的并发安全问题

image-20251222085553297

5.1.1 IDEA配置

  1. 点进Edit Configurations

    image-20251222093525353

  2. ctrl+D创建新程序

    image-20251222093619610

  3. modify opitons添加add VM options

    image-20251222093723952

    image-20251222093713438

  4. 配置新端口

    image-20251222093759435

5.1.2 nginx配置

  1. 找到confi文件夹下的nginx.conf

    image-20251222094042431

  2. 修改文件配置

    image-20251222094219817

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

  3. 重新加载nginx

    image-20251222094523728

5.1.3 测试

  1. 发送多次api/voucher/list/1请求

    image-20251222095442298

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

    image-20251222095532753

    image-20251222095541796

5.1.4 并发安全问题

  • 用PostMan发送两个请求

    image-20251222102954553

    image-20251222103002085

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

    外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

    外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

  • 原因

    1. 正常情况

      image-20251222103514650

    2. 线程交叉解决单线程

      image-20251222103551721

    3. 集群部署情况

      image-20251222103701169

      用synchronized加锁是利用JVM内部的锁监视器实现,可能遇到的问题:在集群部署中每个服务器都会有一个专门的JVM,而synchronized加锁在每个服务器内都会维护一个专门的锁监视器,因此不同的JVM容器中获取锁不会受到其他JVM容器的影响,因此并发情况下依旧会出现安全问题

5.2 分布式锁

5.2.1 功能原理

  • 关键:满足分布式系统或集群模式下多进程可见并且互斥的锁

    image-20251222104209322

5.2.2 不同分布式锁

image-20251222104331697

image-20251222105735388

5.2.3 基于Redis的分布式锁

image-20251222110702793

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

  • 改进

    image-20251222112808679

5.2.4 实现分布式锁

  • SimpleRedisLock类

    image-20251222131438633

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

    image-20251222132709494

    1. 创建锁对象SimpleRedisLock,由于是对order进行相关操作,因此将业务名传入为"order:"再拼接当前用户Id,明确有一人一单业务
    2. 调用SimpleRedisLock的tryLock,传入过期时间,尝试获取锁
    3. 如果获取失败,说明已经有其他业务上锁,返回获取失败
    4. 如果获取锁成功
      1. 获取当前事务代理对象
      2. 调用当前代理对象的createVoucherOrder方法,该方法在第4.2节已经实现

5.3 redis锁可能存在业务阻塞导致的误删问题

5.3.1 问题和解决思路

问题

image-20251222133131485

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

解决思路

image-20251222133300795

  1. 在释放锁前根据存入锁的标志判断一下是否一致,只有一致才能释放锁

5.3.2 改进redis分布式锁

需求:修改之前的分布式锁实现,满足:

  1. 获取锁时存入线程标示(可以用UUID表示)
  2. 在释放锁时先获取锁中的线程标示,判断是否与当前线程标志一致
    1. 如果一致则释放锁
    2. 如果不一致则不释放锁

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

5.3.3 实现改进的redis分布式锁

image-20251222135158474

  1. 使用UUID拼线程ID的方式作为redis中的锁对应的value
  2. 在释放锁时,获取当前线程的UUID和线程ID,并从redis中根据锁查找对应的value,判断两者是否相同
    1. 如果相同,说明该锁是自己加的,可以删除
    2. 如果不同,说明该锁是别人加的,不删除

5.4 redis锁可能存在的释放锁阻塞导致的误删问题

5.4.1 问题和解决思路

  • 问题

image-20251222135642918

  1. 线程1执行业务,判断锁表示OK,正当准备释放锁时,阻塞了(可能时JVM自动运行垃圾回收机制或其他原因)
  2. 线程1一致被阻塞,直到超时自动释放锁,此时线程2尝试获取锁成功
  3. 线程2执行业务时,线程1阻塞完毕,执行释放锁业务,由于已经判断OK,因此会释放当前被线程2获取的锁,而线程2的业务还未结束
  4. 此时线程3获取锁成功,就出现了线程2和线程3的并行问题
  • 解决思路:保证判断和删除操作的一致性

5.4.2 Lua脚本

Redis提供了Lua脚本功能,在一个脚本中编写多条Redis命令,确保多条命令执行时的原子性。Lua是一种编程语言,它的基本语法大家可以参考网站:https://www.runoob.com/lua/lua-tutorial.html

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

image-20251222140606887

image-20251222140805301

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

image-20251222144357087

5.4.4 在java中调用Lua脚本

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

image-20251222141801161

  • 在resources中编写lua脚本

    image-20251222144642981

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

image-20251222142722184

image-20251222142707088

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

5.5 分布式锁存在的问题及Redisson工具包

5.5.1 问题

image-20251222150045462

5.5.2 Redisson介绍

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

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

5.5.3 Redisson入门

image-20251222151657610

image-20251222151706900

  • 配置类

    image-20251222161436985

  • 改写业务

    image-20251222152606126

    image-20251222152635846

    将获取锁的部分用RedissonClient对象中的方法实现

5.6 Redisson原理

5.6.1 可重入锁原理

  • 原本基于lua脚本编写的代码的问题

    image-20251222153339171

    • 以左边的代码为例:发生死锁
      1. 方法1获取锁成功后会调用方法2
      2. 方法2会尝试获取锁,但是由于方法1还没有释放锁,因此无法获取
      3. 而方法1由于没法执行方法2,又迟迟无法结束,从而导致了死锁
  • 实现原理

    为了实现通过redis实现的分布式锁的可重入性,采用的方法是采用hash而非string类型存储,field存储值,而value存储重入次数,每次获取锁都+1,每次释放锁都-1,当减为0时才真正的释放锁。

    redisson在加锁解锁时正式根据下述流程图严谨实现的

    image-20251222155457339

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();
        }
    }
}
  1. 调用RLock对象lock的tryLock方法,传入等待时间和时间单位

    image-20251225092240109

  2. tryLock逻辑

    image-20251225092809147

    1. 时间单位是毫秒
    2. 获取当前时间
    3. 获取当前线程ID
    4. 进入获取锁tryAcquire逻辑
  3. tryAcquire逻辑

    image-20251225092924645

    1. 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,如果不是-1,进入tryLockInnerAsync逻辑,传入的leaseTime即为设定的时间。如果是-1,也进入tryLockInnerAsync逻辑只不过leaseTime位置的参数通过this.commandExecutor.getConnectionManager().getCfg().getLockWatchdogTimeout()获取,是一个默认值image-20251225093709272默认超时时间为30s

      2. 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));
            }
        
        1. 将锁释放时间记录到本地成员变量internalLockLeaseTime
        2. 执行Lua脚本
        3. 如果获取锁成功,返回nil
        4. 如果获取锁失败,返回剩余有效期,单位是毫秒
      3. RFuture<Long> ttlRemainingFuture接受脚本返回的剩余有效期

        1. 由于tryLockInnerAsync是一个异步函数,因此执行完那没拿到值不一定,因此用RFuture接收
      4. 将 ttlRemainingFuture返回

  4. tryLock中的ttl接受返回的剩余有效期

  5. 判断有效期是否为null

    image-20251225100747385

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

    image-20251225100720802

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

      image-20251225101056090

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

        image-20251225101251495

      2. 如果剩余等待时间大于0

        外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

        1. 记录当前时间

        2. 对线程id进行订阅:订阅别人释放锁的信号。publish发布消息通知

          外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

        3. 等待信号,等待时间为time,如果等到最大时间还没等到就不等了,返回false,加上!变成true,那么进入if逻辑

          1. 执行unsubscribe,释放锁(因为已经等待超时)
        4. 如果没进if,说明还没到最大等待时间

          image-20251225105147814

          1. 再计算一次剩余等待时间

            image-20251225111754199

          2. 如果剩余时间大于0,进入循环,开始尝试重新获取

            1. 获取锁,返回剩余有效期,如果为空表示获取成功

              image-20251225111814683

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

              image-20251225111821343

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

                image-20251225111827225

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

                image-20251225111849324

            3. 再次计算剩余等待时间

              image-20251225112540292

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

              image-20251225112620749

5.6.3 超时续约实现

  • 为什么要实现超时续约?

    redis为了解决服务器崩溃等以外原因导致的缓存未被释放,会给缓存设一个过期时间,但是如果业务执行时间过长,超过了这个过期时间,可能会出现并发问题。

    因此要在业务执行时通过看门狗机制实现超时续约,不断更新锁的有效期,直到业务完成,但是锁本身又存在有效期,所以也能应对服务器崩溃导致的内存不释放。

  • 获取锁和释放锁逻辑

image-20251226091322899

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

  • tryAcquireAsync逻辑实现

    image-20251225113458564

    1. 回调成功,获取剩余有效期之后,进入上面的逻辑

      1. 如果异常不为空,表示抛出了异常

      2. 如果异常为空,表示正常执行

        1. 如果剩余有效期为null,表示获取成功

        2. 执行scheduleExpirationRenewal表示更新过期时间

          image-20251225113817988

          1. 通过putIfAbsent进行赋值,key为当前enteyName每个实例(order相关锁是一个key,voucher相关锁是一个key)独一份。如果是第一次则将一个全新的entry放入对应的锁,此时没有oldEntry,返回null。如果不是第一次重入,则不会写入,而是返回oldEntry

            1. 如果oldEntry不为空(不是第一次),则把当前线程Id加入oldEntry。由于不同线程不会获得同一把锁,因此这的逻辑一定是同一线程重入
            2. 如果oldEntry为空(是第一次),加入,并执行renewExpiration刷新缓存时间
          2. renewExpiration逻辑

            image-20251226090408135

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

              image-20251226090610552

              image-20251226090643425

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

            image-20251226090709886

            1. 定时任务负责重置有效期
            2. 递归调用本身,直到捕获释放锁信号

5.6.4 主从一致性原理(Multi lock)

  • 问题

    image-20251226094026262

    如果主节点子节点还没同步,此时主节点宕机,那么子节点没有保存在主节点保存的锁,就出现了主从不一致

  • 解决方案:Multi lock

image-20251226093857375

5.6.5 Redisson总结

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

6. Redis优化秒杀

6.1 问题和优化方案

  • 问题:整个业务流程串行,并且频繁操作数据库导致性能较差

image-20251226095707700

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

image-20251226100419360

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

    image-20251226101351836

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

    image-20251226102138933

  • 业务流程

image-20251226102358886

6.2 基于Redis实现异步秒杀

image-20251226102636712

6.2.1 添加秒杀券

image-20251226110842486

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

image-20251226111739754

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);
    }
  1. 初始化Script
  2. 调用execute方法执行lua脚本
    1. 第一个参数是脚本
    2. 第二个参数是空列表,因为lua里不需要用到Key
    3. 第三个表示voucherId
  3. 判断结果是否为0,不是0表示没有购买资格,根据返回值返回对应的错误信息
  4. 为0表示有购买资格,调用自己写的redisIdWorker.nextId,给他一个唯一的orderId

6.3.4 将秒杀卷添加到阻塞队列

  • 说明:由于VoucherOrderServiceImpl被加上了@service注解,因此其被注册为Spring容器的一个@bean对象,如果线程1,2,3都注入了该对象,其实注入的是同一个对象。儿orderTasks是该对象的一个变量,因此不同线程实际上是在操作同一个阻塞队列,而阻塞队列本身保证了线程安全,因此将voucherOrder添加到阻塞队列时,不会产生线程安全问题

image-20251229090248521

image-20251229090238105

  1. 创建阻塞队列
  2. 设置秒杀卷信息
  3. 将秒杀卷添加到阻塞队列

6.3.5 开启线程,从阻塞队列获取信息

image-20251229093048321

  1. 创建线程池,该线程池为单线程:线程池内只有一个线程

  2. init()方法由于加了PostConstrut注解,在spring boot启动时就会执行,提交一个VoucherOrderHanlder任务,因此在spring boot启动时,就会开始执行VoucherOrderHanlder里的run

  3. run方法内,会循环的执行以下操作

    1. 调用oderTasks.take获取消息队列中的数据(6.3.4已经分析过,不同用户(不同线程)操作的是同一个阻塞队列,并且不会出现线程问题,因为阻塞队列本身已经实现了线程安全)
    2. 获取完阻塞队列中的数据后,调用handleVoucherOrder方法处理订单数据
  4. handleVoucherOrder方法会获取用户Id,然后根据该id尝试获取锁对象,如果获取成功,通过代理对象proxy来调用创建订单方法,并且将代理对象创建在VoucherOrderServiceImpl的成员变量部分,然后在将数据添加到阻塞队列的过程中对代理对象进行初始化,这样就能在handleVoucherOrder中使用代理对象调用createVoucherOrder方法从而避免通过this.createVoucherOrder调用createVoucherOrder方法会出现的跳过@Transactional注解的问题

    image-20251229100700917

    image-20251229100711261

6.3 几个关键问题和解答

6.3.1 为不同用户开启的线程和代码中创建的线程有什么关系和区别?

  1. 代码中创建的是一个单线程的线程池,并在添加了@PostConsrtuct注解的init方法中提交,因此,Spring boot在启动时就会将传入了new VoucherOrderHandler()参数的SECKILL_ORDER_EXECUTOR线程池中的单线程启动。
  2. 用户开启的线程和上述单线程没有关系,例如用户1,2,3同时发起请求,系统为其分配了三个线程,这三个线程和上述线程独立,上述线程是随着spring boot启动而被自动开启的

6.3.2 为何不同的用户能操作同一个阻塞队列?

  1. 该阻塞队列是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调用秒杀接口

  1. 写入阻塞队列操作是在seckillVoucher方法中,只要用户调用了秒杀接口的该方法,就会将数据写入阻塞队列。该过程是线程1,2,3分别执行的,由于阻塞队列本身实现了线程安全,因此不会冲突
  2. 获取数据的线程就如6.3.1所说,是随着spring boot启动开启的独立线程,不断地,循环的从阻塞队列中获取数据,并尝试创建订单

6.3.4 为什么要用代理对象调用创建订单的方法?

  1. 为了实现数据库操作的事务一致性,要加上@Transactional注解,该注解的实现原理是spring Aop,对加上Transactional注解的方法进行增强,如果不用代理对象调用创建订单的方法,而是用this调用,则是内部调用,指向的是原始的,未被代理的创建订单方法本身,则Transactional注解失效。Spring在创建Bean的过程中,会动态地生成一个代理对象来“包裹”您的原始对象,因此调用该动态对象的创建订单方法,则能识别到Transactional注解

6.3.5 代理对象是怎么初始化的?

  1. 和6.3.2类似,proxy首先被创建为VoucherOrderServiceImpl的一个成员变量但未初始化,由于VoucherOrderServiceImpl被注册为了唯一bean,因此不同的用户能访问同一个proxy,只要用户调用了秒杀接口seckillVoucher,就会执行初始化proxy的操作。
  2. 后一个请求的执行可能会覆盖前一个请求设置的proxy值。但在上述业务场景下,通常不会造成问题,因为代理对象本质上是同一个。

6.3.6 会不会出现代理对象还没被初始化就通过其调用createVocherOrder方法的情况?

  1. 不会,因为线程池的唯一线程在执行时,会首先尝试从阻塞队列中获取数据,获取失败会一直等待,而只有至少有一个用户执行了将数据写到阻塞队列操作之后,阻塞队列才会有数据,而只要有用户执行了这个操作,那代理对象就会被初始化

6.3.7 为什么不能在创建代理对象时就初始化?

  1. 因为spring bean注册过程有严格的顺序

    1. 属性填充

    2. 前置处理

    3. @PostConstruct方法执行

    4. 后置处理

    5. 生成代理对象

      如果在第一步就尝试初始化,此时代理对象尚未生成,无法成功初始化。

7. Redis消息队列实现异步秒杀

7.1 Java内置阻塞队列的问题

  1. JVM内存有限,全放阻塞队列可能会导致内存不足
  2. 由于JVM不是独立的服务,因此服务器如果宕机,阻塞队列中的数据可能丢失

7.2 消息队列

消息队列(Message Queue),字面意思就是存放消息的队列。最简单的消息队列模型包括3个角色:

  1. 消息队列:存储和管理消息,也被称为消息代理(Message Broker)
  2. 生产者:发送消息到消息队列
  3. 消费者:从消息队列获取消息并处理消息

image-20251229111402807

Redis提供了三种不同的方式来实现消息队列:

  1. list结构:基于List结构模拟消息队列
  2. PubSub:基本的点对点消息模型
  3. Stream:比较完善的消息队列模型

7.2.1 基于List结构模拟消息队列

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

image-20251229112956476

7.2.1 基于PubSub的消息队列

Pubsub(发布订阅)是Redis2.0版本引入的消息传递模型。顾名思义,消费者可以订阅一个或多个channel,生产者向对应channel发送消息后,所有订阅者都能收到相关消息。

  • SUBSCRlBE channel[channel]:订阅一个或多个频道

  • PUBLlSH channelmsg :向一个频道发送消息

  • PSUBSCRIBE pattern[pattern]:订阅与pattern格式匹配的所有频道

    外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传

    image-20251229113231381

    image-20251229113447912

7.3 基于Redis的Stream实现消息队列

7.3.1 基于Stream的消息队列

Stream 是 Redis 5.0 引入的一种新数据类型,可以实现一个功能非常完善的消息队列。

发送消息的命令:

1653577301737

例如:

1653577349691

读取消息的方式之一:XREAD

1653577445413

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

1653577643629

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

1653577659166

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

1653577689129

注意:当我们指定起始ID为$时,代表读取最新的消息,如果我们处理一条消息的过程中,又有超过1条以上的消息到达队列,则下次获取时也只能获取到最新的一条,会出现漏读消息的问题

STREAM类型消息队列的XREAD命令特点:

  • 消息可回溯
  • 一个消息可以被多个消费者读取
  • 可以阻塞读取
  • 有消息漏读的风险

7.3.2 基于Stream的消息队列-消费者组

消费者组(Consumer Group):将多个消费者划分到一个组中,监听同一个队列。具备下列特点:

1653577801668

创建消费者组:
1653577984924
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中的第一个消息开始

消费者监听消息的基本思路:

1653578211854STREAM类型消息队列的XREADGROUP命令特点:

  • 消息可回溯
  • 可以多消费者争抢消息,加快消费速度
  • 可以阻塞读取
  • 没有消息漏读的风险
  • 有消息确认机制,保证消息至少被消费一次

最后我们来个小对比

1653578560691

Redis Stream 的 消费者组(Consumer Group) 是一个强大的功能,它允许多个消费者协同处理同一个消息队列中的消息,不仅提高了处理效率,还确保了消息不会在消费过程中丢失。下面我用一个简单的类比和梳理,帮你理解它的核心机制。

🏷️ 先了解基本概念

可以把 Redis Stream 想象成一个不断增长的快递流水线,每个快递包裹就是一个消息,带有唯一的 ID(如 1703235642574-0,通常由时间戳和序列号组成)。

消费者组(Consumer Group) 就像是承包了这条流水线某个环节的一个快递团队。这个团队里有多个快递员(消费者),他们共同从流水线上取包裹去派送。

🔑 核心运行机制

消费者组之所以能可靠地工作,依赖于以下几个关键机制:

  1. 消息分发规则:竞争消费

    团队里的快递员是竞争关系。同一个包裹一旦被一个快递员取走,其他快递员就不会再拿到这个包裹了。这保证了同一条消息不会被组内的多个消费者重复处理。

  2. 进度记录:last_delivered_id

    团队里有个记录员,始终记着最后一个被取走的包裹的编号。这个编号就是 last_delivered_id。每当有快递员取走一个包裹,记录员就会把这个编号往前移动,确保大家知道下一个该处理哪个包裹。

  3. 确保消息不丢失: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脚本

    image-20260104113535240

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

    image-20260104113208057

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

    image-20260104113242427

    1. 线程不断的获取消息队列中的数据,等待时间为两秒

      1. 如果获取失败,说明消息队列中没有数据,则跳过本次循环

      2. 如果获取成功,则解析订单信息

        1. list.get(0)获取消息队列中最早的一条信息(包含消息id和消息内容)

          image-20260104130819373

        2. record.getValue()获得消息内容

        3. 调用BeanUtil将消息内容转化为VoucherOrder形式

        4. 确认下单

      3. 下单完成后ACK确认,将消息从peding list中移除

  • handlePendingList

    image-20260104113302531

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

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐