前言

个人学习总结,主要是为了梳理优惠券秒杀功能中各种写法之间的区别、以及不同改进方法中各自的优缺点

短信登录

因为多台Tomcat并不共享session存储空间,当请求切换到不同tomcat服务时就会导致数据丢失的问题,所有采用了Redis实现共享session登录

整体流程

当注册完成后,用户去登录会去校验用户提交的手机号和验证码,是否一致,如果一致,则根据手机号查询用户信息,不存在则新建,最后将用户数据保存到redis,并且生成随机token作为redis的key,当我们校验用户是否登录时,会去携带着token进行访问,从redis中取出token对应的value,判断是否存在这个数据,如果没有则拦截,如果存在则将其保存到ThreadLocal中,并且放行。

// 随机生成登录令牌token
String token = UUID.randomUUID().toString(true);
// 脱敏,利用工具类将User->UserDTO
UserDTO userDTO = BeanUtil.copyProperties(user, UserDTO.class);

/* 工具类,将userDTO转为map对象
 * 所有数据类型都会被转为String类型,以匹配stringRedisTemplate
*/
Map<String, Object> userMap = BeanUtil.beanToMap(userDTO, new HashMap<>(),
	CopyOptions.create().setIgnoreNullValue(true).setFieldValueEditor(
	(fieldName, fieldValue) -> fieldValue.toString()
));
// 将用户转为hash存储到redis
stringRedisTemplate.opsForHash().putAll(LOGIN_USER_KEY + token, userMap);
// 设置token有效期
stringRedisTemplate.expire(LOGIN_USER_KEY + token, LOGIN_USER_TTL, TimeUnit.MINUTES);

return Result.ok(token);

两层拦截器

第一层RefreshTokenInterceptor,会拦截所有请求,但不会做拦截操作,用于刷新token令牌以及从redis获取用户信息并保存在ThreadLocal

ThreadLocal

为每个线程提供独立的用户信息副本,避免多线程并发问题,同时可以减少重复操作,实现数据复用,不必在每次请求时都查询redis获取数据判断是否登录,只需要调用UserHolder.getUser即可,记得在afterCompletion中主动移除user对象,避免OOM

	@Override
    public boolean preHandle(HttpServletRequest request, HttpServletResponse response, Object handler) throws Exception {
        // 获取请求头中的token
        String token = request.getHeader("authorization");
        if(StrUtil.isBlank(token)){
            return true;
        }
        // 基于token获取redis中的token
        String key=LOGIN_USER_KEY + token;
        Map<Object, Object> userMap = stringRedisTemplate.opsForHash()
                .entries(key);
        if(userMap.isEmpty()){
            return true;
        }
        // 将hash对象再转为userDTO
        UserDTO userDTO = BeanUtil.fillBeanWithMap(userMap, new UserDTO(), false);
        // 存在则保存用户到ThreadLocal
        UserHolder.saveUser(userDTO);
        // 刷新token的有效期
        stringRedisTemplate.expire(key, LOGIN_USER_TTL, TimeUnit.MINUTES);
        return true;
    }

    @Override
    public void afterCompletion(HttpServletRequest request, HttpServletResponse response, Object handler, Exception ex) throws Exception {
        // 移除用户
        UserHolder.removeUser();
    }

第二层LoginInterceptor才会对用户登录状态做真正的判断,如果登录状态过期,则会拦截该请求,注意不要拦截注册/登录等无需登录的请求哦。

缓存击穿

互斥锁(setnx)

核心思路:如果从缓存没有查询到数据,则进行互斥锁的获取,获取互斥锁后,判断是否获得到了锁,如果没有获得到,则休眠,过一会再进行尝试,直到获取到锁为止,才能进行查询

如果获取到了锁的线程,再去进行查询,查询后将数据写入redis,再释放锁,返回数据,利用互斥锁就能保证只有一个线程去执行操作数据库的逻辑,防止缓存击穿

    // 互斥锁解决缓存击穿问题
    public Shop queryWithMutex(Long id) {
        String key = CACHE_SHOP_KEY + id;
        String shopJson = stringRedisTemplate.opsForValue().get(key);
        if (StrUtil.isNotBlank(shopJson)) {
            // 存在
            return JSONUtil.toBean(shopJson, Shop.class);
        }
        // 判断命中的是否是空值
        if (shopJson != null) {
            // 返回错误信息
            return null;
        }
        // 不存在,实现缓存重建
        // 获取互斥锁
        String lockKey = LOCK_SHOP_KEY + id;
        Shop shop = null;
        try {
            boolean isLock = tryLock(lockKey);
            if (!isLock) {
                // 失败,休眠并重试
                Thread.sleep(50);
                return queryWithMutex(id);
            }
            shop = getById(id);
            if (shop == null) {
                // 将空值写入redis
                stringRedisTemplate.opsForValue().set(key, "", CACHE_NULL_TTL, TimeUnit.MINUTES);
                return null;
            }
            // 存在,以json格式写入redis
            stringRedisTemplate.opsForValue().set(key, JSONUtil.toJsonStr(shop), CACHE_SHOP_TTL, TimeUnit.MINUTES);
        } catch (InterruptedException e) {
            throw new RuntimeException(e);
        } finally {
            // 释放互斥锁
            unlock(lockKey);
        }
        return shop;
    }

    // 获取互斥锁
    private boolean tryLock(String key) {
        Boolean flag = stringRedisTemplate.opsForValue().setIfAbsent(key, "1", LOCK_SHOP_TTL, TimeUnit.SECONDS);
        return BooleanUtil.isTrue(flag);
    }

    // 释放锁
    private void unlock(String key) {
        stringRedisTemplate.delete(key);
    }

逻辑过期

核心思路:设置为永不过期,添加一个逻辑过期时间,不用去添加字段,而是封装一个新的类。

public class RedisData {
    // 过期时间
    private LocalDateTime expireTime;
    // 实际数据
    private Object data;
}

用户开始查询redis时,判断是否命中,如果没有命中则直接返回空数据(一般不会出现这种情况,因为热点信息会预热提前存入redis,保证缓存存在)

命中后,将value取出,判断value中的过期时间是否满足,如果没有过期,则直接返回redis中的数据,如果过期,则在开启独立线程后直接返回之前的数据,独立线程去重构数据,重构完成后释放互斥锁。

和互斥锁的关键区别就是:管它三七二十一,先返回一个旧的数据给用户先用着,由竞争锁成功的那个线程再去创建一个新线程去完成缓存重构操作
在这里插入图片描述

	// 线程池
    private static final ExecutorService CACHE_REBUILD_EXECUTOR = Executors.newFixedThreadPool(10);

    // 逻辑过期解决缓存击穿问题
    public Shop queryWithLogicalExpire(Long id) {
        String key = CACHE_SHOP_KEY + id;
        // 从redis查询商户缓存
        String shopJson = stringRedisTemplate.opsForValue().get(key);
        if (StrUtil.isBlank(shopJson)) {
            // 未命中
            return null;
        }
        // 命中,判断过期时间
        RedisData redisData = JSONUtil.toBean(shopJson, RedisData.class);
        Shop shop = JSONUtil.toBean((JSONObject) redisData.getData(), Shop.class);
        LocalDateTime expireTime = redisData.getExpireTime();
        if (expireTime.isAfter(LocalDateTime.now())) {
            // 未过期
            return shop;
        }
        // 过期,重建缓存
        String lockKey = LOCK_SHOP_KEY + id;
        boolean isLock = tryLock(key);
        if (isLock) {
            // 开启独立线程,实现缓存构建
            CACHE_REBUILD_EXECUTOR.submit(() -> {
                try {
                    // 重构缓存
                    this.saveShop2Redis(id, 20L);// 测试用时间,根据实际开发调整
                } catch (Exception e) {
                    throw new RuntimeException(e);
                } finally {
                    // 释放锁
                    unlock(lockKey);
                }
            });
        }
        // 竞争失败、成功都先返回旧数据(更新交由上述新线程去执行)
        return shop;
    }

    // 获取互斥锁
    private boolean tryLock(String key) {
        Boolean flag = stringRedisTemplate.opsForValue().setIfAbsent(key, "1", LOCK_SHOP_TTL, TimeUnit.SECONDS);
        return BooleanUtil.isTrue(flag);
    }

    // 释放锁
    private void unlock(String key) {
        stringRedisTemplate.delete(key);
    }

    // 1.模拟预热操作,将热点信息提前加入redis
    // 2.同样也是重构缓存操作
    public void saveShop2Redis(Long id, Long expireSeconds) {
        Shop shop = getById(id);
        RedisData redisData = new RedisData();
        redisData.setData(shop);
        redisData.setExpireTime(LocalDateTime.now().plusSeconds(expireSeconds));
        stringRedisTemplate.opsForValue().set(CACHE_SHOP_KEY + id, JSONUtil.toJsonStr(redisData));
    }
  • 从Redis取出RedisData后,redisData.getData()可以直接强转成Shop类型吗?RedisData的data字段可以使用泛型代替吗?代替后可以直接强转成Shop类型吗?

    这是我在b站发现有很多人有疑问的地方,我就去ai了一下,上面三种都不可以直接操作,都需要一定的转换,(b站弹幕我也有回答哦,绿色的)

    1. 不可以直接强转,因为Hutool的JSONUtil在反序列化时,如果字段类型为 Object,通常会将嵌套的 JSON 对象解析为 JSONObjectMap,而非具体的 Shop 类型。直接强转会导致 ClassCastException

    2. 可以使用泛型代替

    3. 也不可以直接转,因为泛型会被擦除,所以即使你明确了类型,在你通过redisData.getData()取出来时仍会遇到和上面一样的问题,所以仍然需要多一步。(当然方法不唯一)

      String shopJson = stringRedisTemplate.opsForValue().get(CACHE_SHOP_KEY + id);
      // 使用 Hutool 的 TypeUtil 构建完整泛型类型
      Type type = TypeUtil.getType(RedisData.class, Shop.class);
      RedisData<Shop> redisData = JSONUtil.toBean(shopJson, type);
      Shop shop = redisData.getData(); // 直接获取 Shop 对象
      

优惠券秒杀

1. 普通写法

乐观锁解决超卖问题,悲观锁解决一人一单问题,乐观锁比较适合更新数据,而现在是插入数据,所以我们需要使用悲观锁操作。

此版本终极解决方案是:锁/Lua脚本+同步下单

1.1 使用synchronized对象锁
    @Override
    public Result seckillVoucher(Long voucherId) {
        // 查询优惠券
        SeckillVoucher voucher = seckillVoucherService.getById(voucherId);
        if (voucher.getBeginTime().isAfter(LocalDateTime.now())) {
            // 未开始
            return Result.fail("秒杀尚未开始");
        }
        if (voucher.getEndTime().isBefore(LocalDateTime.now())) {
            // 已结束
            return Result.fail("秒杀已经结束");
        }
        if (voucher.getStock() < 1) {
            // 库存不足
            return Result.fail("库存不足");
        }
        // 一人一单功能
        Long userId = UserHolder.getUser().getId();
        /* 
          对userId加锁,细化锁的粒度
          使用intern()从常量池拿到数据,保证是同一对象
        */
        synchronized (userId.toString().intern()) {
            // 获取代理对象(事务)
            IVoucherOrderService proxy = (IVoucherOrderService) AopContext.currentProxy();
            return proxy.createVoucherOrder(voucherId);
        }    
    }

    @Transactional
    @Override
    public Result createVoucherOrder(Long voucherId) {
        Long userId = UserHolder.getUser().getId();
        // 简单select语句,无锁
        int count = query().eq("user_id", userId).eq("voucher_id", voucherId).count();
        if (count > 0) {
            // 已经购买
            return Result.fail("已经购买过一次");
        }
        boolean success = seckillVoucherService.update()
                .setSql("stock=stock-1")
                .eq("voucher_id", voucherId).gt("stock", 0)
                .update();
        if (!success) {
            return Result.fail("库存不足");
        }
        // 创建订单
        VoucherOrder voucherOrder = new VoucherOrder();
        long orderId = redisIdGenerator.nextId("order");
        voucherOrder.setId(orderId);
        voucherOrder.setUserId(userId);
        voucherOrder.setVoucherId(voucherId);
        save(voucherOrder);
        return Result.ok(orderId);
    }
  • 为啥不能用this.的方式调用,而是通过代理

    因为Spring的事务管理是基于AOP实现,当方法被@Transactional注解修饰时,Spring会为该Bean创建代理对象,代理对象在调用目标方法时,会先开启事务,再调用原始方法

    this.是绕过了代理,直接调用原始方法

存在问题:

解决了一人一单问题,但仅适用于单体项目,在集群环境下,每台服务器都有一个jvm,synchronized 底层使用的JVM级别中的Monitor 来决定当前线程是否获得了锁,monitor内部有一个变量Owner,标识持有锁的线程,一个线程获取到锁的标志就是在monitor中设置成功了Owner,一个monitor中只能有一个Owner,但是我们有多台服务器,有多个jvm,多个monitor,也就有多个Owner,所以需要下面的分布式锁来解决

1.2 分布式锁-Redis
1.2.1 自定义简单锁
        // public Result seckillVoucher(Long voucherId)
		// 自定义简单分布式锁
        SimpleRedisLock lock = new SimpleRedisLock("order:" + userId, stringRedisTemplate);
        boolean isLock = lock.tryLock(10);
        if (!isLock) {
            // 返回失败
            return Result.fail("不允许重复下单");
        }
        try {
            IVoucherOrderService proxy = (IVoucherOrderService) AopContext.currentProxy();
            return proxy.createVoucherOrder(voucherId);
        } finally {
            lock.unlock();
        }
    }

// SimpleRedisLock
private static final String KEY_PREFIX="lock:";

@Override
public boolean tryLock(long timeoutSec) {
    // 获取线程标示
    long threadId = Thread.currentThread().getId();
    // 获取锁
    Boolean success = stringRedisTemplate.opsForValue()
            .setIfAbsent(KEY_PREFIX + name, threadId + "", timeoutSec, TimeUnit.SECONDS);
    return Boolean.TRUE.equals(success);
}

public void unlock() {
    //通过del删除锁
    stringRedisTemplate.delete(KEY_PREFIX + name);
}

存在问题: key:"lock:order:id",v:threadId,存在误删问题

持有锁的线程在锁的内部出现了阻塞,还未执行释放操作,这时他的锁自动释放了,此时其他线程,线程2来尝试并获得了这把锁,然后线程2在持有锁执行过程中,线程1反应过来,继续执行未完成释放操作,此时就会把本属于线程2的锁进行删除

1.2.2 添加线程标识

核心逻辑:在存入锁时,放入自己线程的标识(UUID实现),在删除锁时,判断当前这把锁的标识是不是自己存入的,如果是,则进行删除,如果不是,则不进行删除。

private static final String ID_PREFIX = UUID.randomUUID().toString(true) + "-";

@Override
public boolean tryLock(long timeoutSec) {
   // 获取线程标示
   String threadId = ID_PREFIX + Thread.currentThread().getId();
   // 获取锁
   Boolean success = stringRedisTemplate.opsForValue()
                .setIfAbsent(KEY_PREFIX + name, threadId, timeoutSec, TimeUnit.SECONDS);
   return Boolean.TRUE.equals(success);
}

public void unlock() {
    // 获取线程标示
    String threadId = ID_PREFIX + Thread.currentThread().getId();
    // 获取锁中的标示
    String id = stringRedisTemplate.opsForValue().get(KEY_PREFIX + name);
    // 判断标示是否一致
    if(threadId.equals(id)) {
        // 释放锁
        stringRedisTemplate.delete(KEY_PREFIX + name);
    }
}

存在问题: key:"lock:order:id",v:ID_PREFIX + Thread.currentThread().getId(),仍然存在误删问题

更极端情况下,在线程1已经判断标识一致后,正准备释放锁,此时线程被阻塞且锁过期,被线程2抢占,此时线程1反应过来,又会误删线程2的锁。由于redis的锁会因为阻塞等原因过期释放,导致这两个操作不能同时执行成功或不执行

究其原因,还是因为多个redis命令是非原子性

1.2.3 Lua脚本解决多条命令原子性问题
-- unlock.lua
-- 获取锁中的标识与当前线程标识
if(redis.call('get',KEYS[1]) == ARGV[1]) then
    -- 一致,释放锁
    return redis.call('del',KEYS[1])
end
-- 不一致,直接返回
return 0
	// 加载lua脚本
    private static final DefaultRedisScript<Long> UNLOCK_SCRIPT;
    static {
        UNLOCK_SCRIPT=new DefaultRedisScript<>();
        UNLOCK_SCRIPT.setLocation(new ClassPathResource("unlock.lua"));
        UNLOCK_SCRIPT.setResultType(Long.class);
    }

	@Override
    public void unlock() {
        // 调用lua脚本
        stringRedisTemplate.execute(
                UNLOCK_SCRIPT,
                Collections.singletonList(KEY_PREFIX + name),
                ID_PREFIX + Thread.currentThread().getId()
        );
    }

存在问题: 基于setnx实现的分布式锁存在不可重入不可重试超时释放主从一致性等问题,方便起见,使用其他封装好的框架或就好,比如Redission。

1.3 分布式锁-Redission
		// 使用Redisson框架
        RLock lock = redissonClient.getLock("lock:order:" + userId);
        boolean isLock = lock.tryLock();// 无参,获取锁失败直接返回
        if (!isLock) {
            // 返回失败
            return Result.fail("不允许重复下单");
        }
        try {
            IVoucherOrderService proxy = (IVoucherOrderService) AopContext.currentProxy();
            return proxy.createVoucherOrder(voucherId);
        } finally {
            lock.unlock();
        }

存在问题: 利用现成的框架似乎解决了所有问题,但是还有优化的空间,同步操作就好像一个人既当前台小妹,又要当厨师,实在忙不过来啊,所以可以使用异步操作,前台小妹只负责点单,把客人骗进来再说,做饭就交给厨师来,慢慢来嘛

2. Lua脚本+异步下单

基于Lua脚本,判断秒杀库存、一人一单,决定用户是否抢购成功,并对redis中库存进行预扣减,保存订单信息,如果抢购成功,原线程直接返回订单id然后结束,其他交由新线程,不断从阻塞队列/消息队列中获取信息,实现异步下单功能,更新库存并保存订单到数据库。

此版本终极解决方案是:Lua脚本+Stream消息队列实现异步下单

2.1 阻塞队列

VoucherServiceImpl

    // 新增秒杀优惠券的同时,将优惠券库存信息保存到Redis中
	@Override
    @Transactional
    public void addSeckillVoucher(Voucher voucher) {
        // 保存优惠券
        save(voucher);
        // 保存秒杀信息
        SeckillVoucher seckillVoucher = new SeckillVoucher();
        seckillVoucher.setVoucherId(voucher.getId());
        seckillVoucher.setStock(voucher.getStock());
        seckillVoucher.setBeginTime(voucher.getBeginTime());
        seckillVoucher.setEndTime(voucher.getEndTime());
        seckillVoucherService.save(seckillVoucher);
        // 保存秒杀库存到redis中
        stringRedisTemplate.opsForValue().set(SECKILL_STOCK_KEY + voucher.getId(), voucher.getStock().toString());
    }

seckill.lua

-- 判断秒杀库存、一人一单,决定用户是否抢购成功,向消息队列中添加消息
-- 先确保redis中已正确添加消费者组,消息队列,库存
-- 优惠券id
local voucherId = ARGV[1]
-- 用户id
local userId = ARGV[2]
-- 订单id
-- local orderId = ARGV[3] 消息队列参数

-- 库存key
local stockKey = 'seckill:stock:' .. voucherId
-- 订单key
local orderKey = 'seckill:order:' .. voucherId

-- 判断库存是否充足 get stockKey,tonumber--字符串转数字
if(tonumber(redis.call('get', stockKey)) <= 0) then
    -- 库存不足
    return 1
end
-- 判断用户是否下单 SISMEMBER orderKey userId,如果当前用户没下过单,返回0
if(redis.call('sismember', orderKey, userId) == 1) then
    -- 存在,说明是重复下单,返回2
    return 2
end
-- 扣库存 incrby stockKey -1
redis.call('incrby', stockKey, -1)
-- 下单(保存用户)sadd orderKey userId
redis.call('sadd', orderKey, userId)
-- 发送消息到队列中, XADD stream.orders * k1 v1 k2 v2 ...
-- redis.call('xadd', 'stream.orders', '*', 'userId', userId, 'voucherId', voucherId, 'id', orderId) 消息队列命令
return 0

VoucherOrderServiceImpl

	// 代理对象
    private IVoucherOrderService proxy;

    // 加载lua脚本
    private static final DefaultRedisScript<Long> SECKILL_SCRIPT;

// 阻塞队列
    private BlockingQueue<VoucherOrder> orderTasks = new ArrayBlockingQueue<>(1024 * 1024);

    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 class VoucherOrderHandler implements Runnable {
        @Override
        public void run() {
            while (true) {
                try {
                    // 获取队列中的订单信息
                    VoucherOrder voucherOrder = orderTasks.take();
                    // 创建订单
                    handleVoucherOrder(voucherOrder);
                } catch (Exception e) {
                    log.error("处理订单异常", e);
                }
            }
        }
    }

    private void handleVoucherOrder(VoucherOrder voucherOrder) {
        Long userId = voucherOrder.getUserId();
        // 兜底方案,不获取锁也可以,前面已经做过并发问题的判断,这里将redis库存与订单信息同步到mysql即可
        RLock lock = redissonClient.getLock("lock:order:" + userId);
        boolean isLock = lock.tryLock();
        if (!isLock) {
            // 获取锁失败
            log.error("不允许重复下单");
            return;
        }
        try {
            proxy.createVoucherOrder(voucherOrder);
        } finally {
            lock.unlock();
        }
    }
    /**
     * 阻塞队列+lua脚本实现异步下单
     *
     * @param voucherId 优惠券id
     * @return 订单id
     */

    @Override
    public Result seckillVoucher(Long voucherId) {
        // 获取用户
        Long userId = UserHolder.getUser().getId();
        // 执行lua脚本
        Long result = stringRedisTemplate.execute(
                SECKILL_SCRIPT,
                Collections.emptyList(),
                voucherId.toString(), userId.toString()
        );
        int r = result.intValue();
        if (r != 0) {
            // 没有购买资格
            return Result.fail(r == 1 ? "库存不足" : "不能重复下单");
        }
        // 创建订单
        VoucherOrder voucherOrder = new VoucherOrder();
        long orderId = redisIdGenerator.nextId("order");
        voucherOrder.setId(orderId);
        voucherOrder.setUserId(userId);
        voucherOrder.setVoucherId(voucherId);
        // 保存到阻塞队列
        orderTasks.add(voucherOrder);
        // 获取代理对象
        proxy = (IVoucherOrderService) AopContext.currentProxy();

        return Result.ok(orderId);
    }

    @Transactional // 同步库存,保存订单信息到数据库
    @Override
    public void createVoucherOrder(VoucherOrder voucherOrder) {
        Long userId = voucherOrder.getId();
        Long voucherId = voucherOrder.getVoucherId();
        int count = query().eq("user_id", userId).eq("voucher_id", voucherId).count();
        if (count > 0) {
            // 已经购买
            log.error("重复下单");
            return;
        }
        boolean success = seckillVoucherService.update()
                .setSql("stock=stock-1")
                .eq("voucher_id", voucherId).gt("stock", 0)
                .update();
        if (!success) {
            log.error("库存不足");
            return;
        }
        save(voucherOrder);
    }

存在问题: 内存限制问题和数据安全问题

2.2 Redis Stream

VoucherOrderServiceImpl

@Slf4j
@Service
public class VoucherOrderServiceImpl extends ServiceImpl<VoucherOrderMapper, VoucherOrder> implements IVoucherOrderService {

    @Resource
    private ISeckillVoucherService seckillVoucherService;

    @Resource
    private RedisIdGenerator redisIdGenerator;

    @Resource
    private StringRedisTemplate stringRedisTemplate;

    // 代理对象
    private IVoucherOrderService proxy;

    // 加载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);
    }

    // 单线程化线程池,异步处理慢慢搞就好
    private static final ExecutorService SECKILL_ORDER_EXECUTOR = Executors.newSingleThreadExecutor();

    @PostConstruct // 当前类初始化完成后执行
    private void init() {
        SECKILL_ORDER_EXECUTOR.submit(new VoucherOrderHandler());
    }

    // 线程任务,在秒杀前执行
    private class VoucherOrderHandler implements Runnable {
        String mqName = "stream.orders";

        @Override
        public void run() {
            while (true) {
                try {
// 获取队列中的订单信息 XREADGROUP GROUP g1 c1 COUNT 1 BLOCK 2000 STREAMS stream.orders >
                    List<MapRecord<String, Object, Object>> list = stringRedisTemplate.opsForStream().read(
                            Consumer.from("g1", "c1"),
                            StreamReadOptions.empty().count(1).block(Duration.ofSeconds(2)),
                            StreamOffset.create(mqName, ReadOffset.lastConsumed())
                    );
                    // 判断消息获取是否成功
                    if (list == null || list.isEmpty()) {
                        continue;
                    }
                    // 解析消息中的订单消息
                    MapRecord<String, Object, Object> record = list.get(0);
                    Map<Object, Object> values = record.getValue();
                    VoucherOrder voucherOrder = BeanUtil.fillBeanWithMap(values, new VoucherOrder(), true);
                    // 获取成功,创建订单
                    handleVoucherOrder(voucherOrder);
                    // ACK确认 SACK stream.orders g1 id
                    stringRedisTemplate.opsForStream().acknowledge(mqName, "g1", record.getId());
                } catch (Exception e) {
                    log.error("处理订单异常", e);
                    handlePendingList();
                }
            }
        }

        // 处理pending_list中已消费但未确认的消息
        private void handlePendingList() {
            while (true) {
                try {
// 获取pending_list中的订单信息 XREADGROUP GROUP g1 c1 COUNT 1 STREAMS stream.orders 0
                    List<MapRecord<String, Object, Object>> list = stringRedisTemplate.opsForStream().read(
                            Consumer.from("g1", "c1"),
                            StreamReadOptions.empty().count(1).block(Duration.ofSeconds(2)),
                            StreamOffset.create(mqName, ReadOffset.from("0"))
                    );
                    // 判断消息获取是否成功
                    if (list == null || list.isEmpty()) {
                        // 说明pending_list中没有异常消息
                        break;
                    }
                    // 解析消息中的订单消息
                    MapRecord<String, Object, Object> record = list.get(0);
                    Map<Object, Object> values = record.getValue();
                    VoucherOrder voucherOrder = BeanUtil.fillBeanWithMap(values, new VoucherOrder(), true);
                    // 获取成功,创建订单
                    handleVoucherOrder(voucherOrder);
                    // ACK确认 SACK stream.orders g1 id
                    stringRedisTemplate.opsForStream().acknowledge(mqName, "g1", record.getId());
                } catch (Exception e) {
                    log.error("处理pending_list异常", e);
                    try {
                        Thread.sleep(20);
                    } catch (InterruptedException interruptedException) {
                        interruptedException.printStackTrace();
                    }
                }
            }
        }
    }

    private void handleVoucherOrder(VoucherOrder voucherOrder) {
        Long userId = voucherOrder.getUserId();
        // 兜底方案,不获取锁也可以,前面已经做过并发问题的判断,这里将redis库存与订单信息同步到mysql即可
        RLock lock = redissonClient.getLock("lock:order:" + userId);
        boolean isLock = lock.tryLock();
        if (!isLock) {
            // 获取锁失败
            log.error("不允许重复下单");
            return;
        }
        try {
            proxy.createVoucherOrder(voucherOrder);
        } finally {
            lock.unlock();
        }
    }

    /**
     * stream消息队列+lua脚本异步下单
     *
     * @param voucherId 优惠券id
     * @return 订单id
     */
    @Override
    public Result seckillVoucher(Long voucherId) {
        // 查询优惠券
        SeckillVoucher voucher = seckillVoucherService.getById(voucherId);
        // 虽然前端做了,但还是要做验证,防止绕过前端进行请求
        if (voucher.getBeginTime().isAfter(LocalDateTime.now())) {
            // 未开始
            return Result.fail("秒杀尚未开始");
        }
        if (voucher.getEndTime().isBefore(LocalDateTime.now())) {
            // 已结束
            return Result.fail("秒杀已经结束");
        }
        // 获取用户
        Long userId = UserHolder.getUser().getId();
        // 生成订单id
        long orderId = redisIdGenerator.nextId("order");
        // 执行lua脚本(此处lua脚本应有消息队列逻辑)
        Long result = stringRedisTemplate.execute(
                SECKILL_SCRIPT,
                Collections.emptyList(),
                voucherId.toString(), userId.toString(), String.valueOf(orderId)
        );
        int r = result.intValue();
        if (r != 0) {
            // 没有购买资格
            return Result.fail(r == 1 ? "库存不足" : "不能重复下单");
        }
        // 获取代理对象
        proxy = (IVoucherOrderService) AopContext.currentProxy();
        return Result.ok(orderId);
    }

    @Transactional // 同步库存,保存订单信息到数据库
    @Override
    public void createVoucherOrder(VoucherOrder voucherOrder) {
        Long userId = voucherOrder.getId();
        Long voucherId = voucherOrder.getVoucherId();
        int count = query().eq("user_id", userId).eq("voucher_id", voucherId).count();
        if (count > 0) {
            // 已经购买
            log.error("重复下单");
            return;
        }
        boolean success = seckillVoucherService.update()
                .setSql("stock=stock-1")
                .eq("voucher_id", voucherId).gt("stock", 0)
                .update();
        if (!success) {
            log.error("库存不足");
            return;
        }
        save(voucherOrder);
    }
}

好友关注

关注推送

采用的是推模式

@Override
    public Result saveBlog(Blog blog) {
        Long userId = UserHolder.getUser().getId();
        blog.setUserId(userId);
        boolean success = save(blog);
        if (!success) {
            return Result.fail("新增blog失败");
        }
        // 查询blog作者所有的粉丝
        List<Follow> follows = followService.query().eq("follow_user_id", userId).list();
        if (follows == null || follows.isEmpty()) {
            // 没有粉丝,不用推
            return Result.ok(blog.getId());
        }
        // 推送笔记id给所有粉丝
        for (Follow follow : follows) {
            Long fanId = follow.getUserId();
            String key = FEED_KEY + fanId;
            // 推到粉丝收件箱
            stringRedisTemplate.opsForZSet().add(key, blog.getId().toString(), System.currentTimeMillis());
        }
        return Result.ok(blog.getId());
    }

Feed流的滚动分页

基于redis的sortedset实现,以时间戳为分数,倒序展示,需要注意offset的问题:

eg:score:8 6 6 6 5 4 3,limit 1(第一次为0),count 3,每页3个

  • 第一次:8 6 6,minTime=6

  • 第二次:从上一页最小的值开始,为6

我们希望的结果是:

  • (8 6 66 5 4,从最后一个6开始,偏移1位,展示3个

但结果是:

  • (8 66 6 5,从第一个最小值6开始了,而不是上一页最后一个数据6

所以在代码中,首先需要记录展示的这页中的最小值,再计算出下一次的偏移量

这个例子里上一页有2个重复的最小值6,因此下一次应该偏移2

记住是同一页内多个重复的最小值哦,如果当前页一共有4条数据,8 6 6 4,最小值为4,那下一次分页只要从4开始偏移一位就够了,和几个6没啥关系

BlogServiceImpl

	@Override
    public Result queryBlogOfFollow(Long max, Integer offset) {
        Long userId = UserHolder.getUser().getId();
        // 查询收件箱,接受关注博主推送的消息,逆序展示,ZREVRANGEBYSCORE key Max Min LIMIT offset count
        String key = FEED_KEY + userId;
        // 滚动分页查询,一次展示一页,如果有新增,用户在前端向上滚动时会再次调用这个方法
        Set<ZSetOperations.TypedTuple<String>> typedTuples = stringRedisTemplate.opsForZSet()
                .reverseRangeByScoreWithScores(key, 0, max, offset, 2);
        if (typedTuples == null || typedTuples.isEmpty()) {
            return Result.ok();
        }
        //、
        /* 解析获取数据
           blogId: 用于查询blog具体信息并展示
           minTime(时间戳):  记录当前页最后一条数据的值,作为下一页的最大值max
           os: 记录从最大值max开始,需要偏移多少才能到下一页,offset
           更新后minTime,os,和blogs会打包给前端,其中blogs用于展示,
           minTime,os会作为max和offset,成为再次调用该方法的参数,展示下一页
         */
        List<Long> ids = new ArrayList<>(typedTuples.size());
        long minTime = 0;
        int os = 1;
        for (ZSetOperations.TypedTuple<String> tuple : typedTuples) {
            ids.add(Long.valueOf(tuple.getValue()));
            long time = tuple.getScore().longValue();
            if (time == minTime) {
                os++;
            } else {
                minTime = time;
                os = 1;
            }
        }
        // 根据id查询blogs
        String idStr = StrUtil.join(",", ids);
        List<Blog> blogs = query().in("id", ids).last("ORDER BY FIELD (id, " + idStr + ")").list();
        for (Blog blog : blogs) {
            // 查询blog相关用户
            queryBlogUser(blog);
            // 查询blog点赞状态
            isBlogLiked(blog);
        }
        // 封装结果并返回,为的是下一次用户发起查询下一页的请求时,前端再返回给后端
        ScrollResult scrollResult = new ScrollResult();
        scrollResult.setList(blogs);
        scrollResult.setOffset(os);
        scrollResult.setMinTime(minTime);
        return Result.ok(scrollResult);
    }

简历写法

项目名称:探味志

项目介绍:该项目是适用于高并发场景下的休闲生活类点评项目,实现了宝藏店铺点评、生活经历分享、大额优惠抢购秒杀,签到等功能。

技术栈:Spring Boot + MySQL + Redis + RocketMQ + MyBatis Plus + Redisson

主要工作:

  1. 使用 Redis 解决了集群模式下Session共享问题,通过自定义两层拦截器实现用户的登录校验和权限刷新,并通过 ThreadLocal 实现用户信息跨层
  2. 基于 Cache Aside Pattern 解决数据库与缓存的一致性问题,实现双写一致,保证缓存更新策略的高一致性需求
  3. 通过缓存关键信息,降低了数据库查询压力,解决缓存穿透(空值缓存)、雪崩(随机TTL),击穿(互斥锁、逻辑过期)问 题,并封装成通用工具类
  4. 运用 Redisson 分布式锁和 Lua 脚本,以乐观锁解决集群环境下一人一单问题,规避超卖和线程安全风险
  5. 利用 RocketMQ 消息队列 实现秒杀下单的异步削峰与流程解耦,JMeter压测平均响应时间降低80%,吞吐量提升57%
  6. 使用 BitMap 位图实现用户签到功能,存储海量签到数据,节省大量空间,同时基于 ZSet 实现点赞排行榜,支持排名实时展示
Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐