【Redis|实战篇5】黑马点评|秒杀优化+消息队列
文章目录
5.秒杀优化
5.1异步秒杀思路

之前,整个流程同步执行,还需要频繁调用数据库,耗费时间性能较低。可以将同步转为异步,一部分工作让Redis执行,并且开启一个新线程去异步操作数据库

优化思路:
- 新增秒杀优惠券的同时,将优惠券信息保存到Redis中
- 基于Lua脚本,判断秒杀库存、一人一单,决定用户是否抢购成功
- 如果抢购成功将优惠券id和用户id封装后存入阻塞队列
- 开启线程任务,不断从阻塞队列中获取信息,实现异步下单功能

存秒杀库存用String类型,存一人一单用Set类型
5.2基于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:" + voucher.getId(),voucher.getStock().toString());
}
Lua脚本
---
--- Generated by EmmyLua(https://github.com/EmmyLua)
--- Created by ghp.
--- DateTime: 2023/7/15 15:22
--- Description 判断库存是否充足 && 判断用户是否已下单
---
-- 优惠券id
local voucherId = ARGV[1];
-- 用户id
local userId = ARGV[2];
-- 库存的key
local stockKey = 'seckill:stock:' .. voucherId;
-- 订单key
local orderKey = 'seckill:order:' .. voucherId;
-- 判断库存是否充足 get stockKey > 0 ?
local stock = redis.call('GET', stockKey);
if (tonumber(stock) <= 0) then
-- 库存不足,返回1
return 1;
end
-- 库存充足,判断用户是否已经下过单 SISMEMBER orderKey userId
if (redis.call('SISMEMBER', orderKey, userId) == 1) then
-- 用户已下单,返回2
return 2;
end
-- 库存充足,没有下过单,扣库存、下单
redis.call('INCRBY', stockKey, -1);
redis.call('SADD', orderKey, userId);
-- 返回0,标识下单成功
return 0;
下单
@Override
public Result seckillVouche(Long voucherId) {
//获取用户id
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){
//3.不为0没有购买资格
return Result.fail(r == 1 ? "库存不足" : "不能重复下单");
}
//4.为0,把下单信息存入消息队列
long orderId = redisIdWorker.nextId("order");
//TODO 存入消息队列
//5.返回订单id
return Result.ok(orderId);
}
5.3基于阻塞队列实现秒杀异步下单
package com.hmdp.service.impl;
import com.hmdp.config.RedissonConfig;
import com.hmdp.dto.Result;
import com.hmdp.entity.SeckillVoucher;
import com.hmdp.entity.VoucherOrder;
import com.hmdp.mapper.VoucherOrderMapper;
import com.hmdp.service.ISeckillVoucherService;
import com.hmdp.service.IVoucherOrderService;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.hmdp.utils.RedisConstants;
import com.hmdp.utils.RedisIdWorker;
import com.hmdp.utils.SimpleRedisLock;
import com.hmdp.utils.UserHolder;
import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import org.springframework.aop.framework.AopContext;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.core.io.ClassPathResource;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.data.redis.core.script.DefaultRedisScript;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.PostConstruct;
import javax.annotation.Resource;
import java.time.LocalDateTime;
import java.util.Collection;
import java.util.Collections;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
/**
* <p>
* 服务实现类
* </p>
*
* @author 虎哥
* @since 2021-12-22
*/
@Service
public class VoucherOrderServiceImpl extends ServiceImpl<VoucherOrderMapper, VoucherOrder> implements IVoucherOrderService {
@Resource
private ISeckillVoucherService seckillVoucherService;
@Resource
private RedisIdWorker redisIdWorker;
@Autowired
private StringRedisTemplate stringRedisTemplate;
@Resource
private RedissonClient redissonClient;
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 BlockingQueue<VoucherOrder> orderTasks = new ArrayBlockingQueue<>(1024 * 1024);
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 {
//1.获取队列中的订单信息
VoucherOrder voucherOrder = orderTasks.take();
//2.创建订单
} catch (Exception e) {
log.error("处理订单异常" + e);
}
}
}
}
private void handleVoucherOrder(VoucherOrder voucherOrder){
//1.获取用户
Long userId = voucherOrder.getUserId();
//2.创建锁对象
RLock lock = redissonClient.getLock("lock:order:" + userId);
//获取锁
boolean isLock = lock.tryLock();
if (!isLock) {
// 索取锁失败,重试或者直接抛异常(这个业务是一人一单,所以直接返回失败信息)
log.error("一人只能下一单");
return;
}
try {
// 创建订单(使用代理对象调用,是为了确保事务生效)
proxy.createVoucherOrder(voucherOrder);
} finally {
//释放锁
lock.unlock();
}
}
private IVoucherOrderService proxy;
@Override
public Result seckillVouche(Long voucherId) {
//获取用户id
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){
//3.不为0没有购买资格
return Result.fail(r == 1 ? "库存不足" : "不能重复下单");
}
//4.为0,把下单信息存入消息队列
long orderId = redisIdWorker.nextId("order");
VoucherOrder voucherOrder = new VoucherOrder();
voucherOrder.setId(orderId);
voucherOrder.setUserId(userId);
voucherOrder.setVoucherId(voucherId);
//存入消息队列
orderTasks.add(voucherOrder);
//获取代理对象
IVoucherOrderService proxy = (IVoucherOrderService) AopContext.currentProxy();
//5.返回订单id
return Result.ok(orderId);
}
//// SimpleRedisLock lock = new SimpleRedisLock(stringRedisTemplate,"order:" + userId);
// RLock lock = redissonClient.getLock("lock:order:" + userId);
// boolean isLock = lock.tryLock();//Redisson不用手动设置有效期
// if (!isLock) {
// //获取锁失败
// return Result.fail("一人只能下一单");
// }
// try {//获取代理对象
// IVoucherOrderService proxy = (IVoucherOrderService) AopContext.currentProxy();
// return proxy.createVoucherOrder(voucherId);
// } finally {
// lock.unlock();
// }
//
// }
@Transactional
public void createVoucherOrder(VoucherOrder voucherOrder) {
//4.5一人一单
Long userId = UserHolder.getUser().getId();
int count = query().eq("user_id", userId).eq("voucher_id", voucherOrder.getVoucherId()).count();//查询订单
if (count > 0) {
log.error("用户已经购买过一次了");
return;
}
//5.扣减库存
boolean success = seckillVoucherService.update()
.setSql("stock = stock - 1")
.eq("voucher_id", voucherOrder.getVoucherId()).gt("stock", 0)
.update();
if (!success) {
log.error("库存不足");
return;
}
//6.创建订单
save(voucherOrder);
}
}
- 队满时,put被阻塞
- 队空时,take被阻塞
用阻塞队列的缺点:
-
内存限制问题
-
数据安全问题
阻塞队列基于JVM内存,不具有持久化能力
- 服务器宕机—>数据丢失
- 并发处理—>数据错乱
6.Redis消息队列
6.1认识消息队列
消息队列字面意思就是存放消息的队列
简单的消息队列包括3个角色:
- 消息队列:存储和管理消息
- 生产者:发送消息到消息队列
- 消费者:从消息队列获取消息并处理消息

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

6.2基于List实现消息队列
Redis的list数据结构是一个双向链表,很容易模拟出队列效果

队列入口和出口不在一边,可以使用LPUSH+RPOP来实现
注意:
当队列中没有消息是RPOP或LPOP操作会返回null,并不像JVM的阻塞队列那样会阻塞并等待消息
因此应当使用BRPOP或BLPOP来实现阻塞效果
- 优点:
- 利用Redis存储,不受限于JVM内存上限
- 基于Redis的持久化机制,数据安全性有保证
- 可以满足消息有序性
- 缺点:
- 无法避免消息丢失
- 只支持单消息者
6.3PubSub实现消息队列
PubSub(发布订阅)是Redis2.0版本引入的消息传递模型
消费者可以订阅一个或多个channel,生产者向对应channel发送消息后,所有订阅者都能收到相关消息
- SUBSCRIBLE channel [channel]:订阅一个或多个频道
- PUBLISH channel msg:向一个频道发送消息
- PSUBSCRIBLE pattern [pattern]:订阅与pattern格式匹配的所有频道

- 优点:采用发布订阅模型,支持多生产、多消费
- 缺点:
- 不支持数据持久化
- 无法避免消息丢失
- 消息堆积有上限,超出时数据丢失
6.4Stream的单消费者模式
Stream是Redis5.0引入的一种新数据类型,可以实现功能非常完善的消息队列
发送消息:用于向指定的Stream流中添加一个消息
XADD key *|ID value [value ...]
# 创建名为 users 的队列,并向其中发送一个消息,内容是:{name=jack,age=21},并且使用Redis自动生成ID

读取消息:
XREAD [COUNT count] [BLOCK milliseconds] STREAMS key [key ...] ID ID


注意:当我们指定起始ID为$时代表读取最后一条消息(读取最新的消息)ID为0时代表读最开始的一条消息(读取最旧的消息),如果我们处理一条消息的过程中,又有超过1条以上的消息到达队列,则下次获取时也只能获取到最新的一条,会出现漏读消息的问题
-
优点:
- 消息可回溯
- 一个消息可以被多个消费者消费
- 可以阻塞读取
-
缺点:有消息漏读的风险
6.5Stream的消费者组模式
消费者组:将多个消费者划分到一个组中,监听同一个队列
特点:

创建消费者组:
XGROUP CREATE key groupName ID
- key:队列名称
- groupName:消费者组名称
- ID:起始ID标识,$代表队列中最后一个消息,0代表队列中的第一个消息
- MKSTREAM:队列不存在时自动创建队列
其他常见命令:
# 删除指定的消费者组
XGROUP DESTORY key groupName
# 给指定的消费者组添加消费者
XGROUP CREATECONSUMER key groupName consumerName
# 删除消费者组中指定消费者
XGROUP DELCONSUMER key groupName consumerName
从消费者组读取信息:
XREADGROUP GROUP

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

Stream类型消息队列的XREADGROUP命令特点:
- 消息可回溯
- 可以多消费者争抢消息,加快消费速度
- 可以阻塞读取
- 没有消息漏读的风险
- 有消息确认机制,保证消息至少被消费一次
6.6基于Stream消息队列实现异步秒杀

1.创建队列
# 创建队列(消费者组模式)
XGROUP CREATE stream.orders g1 0 MKSTREAM
2.VoucherOrderServiceImpl
@Service
public class VoucherOrderServiceImpl extends ServiceImpl<VoucherOrderMapper, VoucherOrder> implements IVoucherOrderService {
@Resource
private ISeckillVoucherService seckillVoucherService;
@Resource
private RedisIdWorker redisIdWorker;
@Resource
private StringRedisTemplate stringRedisTemplate;
@Resource
private RedissonClient redissonClient;
/**
* 当前类初始化完毕就立马执行该方法
*/
@PostConstruct
private void init() {
// 执行线程任务
SECKILL_ORDER_EXECUTOR.submit(new VoucherOrderHandler());
}
/**
* 线程池
*/
private static final ExecutorService SECKILL_ORDER_EXECUTOR = Executors.newSingleThreadExecutor();
/**
* 队列名
*/
private static final String queueName = "stream.orders";
/**
* 线程任务: 不断从消息队列中获取订单
*/
private class VoucherOrderHandler implements Runnable {
@Override
public void run() {
while (true) {
try {
// 1、从消息队列中获取订单信息 XREADGROUP GROUP g1 c1 COUNT 1 BLOCK 1000 STREAMS streams.order >
List<MapRecord<String, Object, Object>> messageList = stringRedisTemplate.opsForStream().read(
Consumer.from("g1", "c1"),
StreamReadOptions.empty().count(1).block(Duration.ofSeconds(1)),
StreamOffset.create(queueName, ReadOffset.lastConsumed())
);
// 2、判断消息获取是否成功
if (messageList == null || messageList.isEmpty()) {
// 2.1 消息获取失败,说明没有消息,进入下一次循环获取消息
continue;
}
// 3、消息获取成功,可以下单
// 将消息转成VoucherOrder对象
MapRecord<String, Object, Object> record = messageList.get(0);
Map<Object, Object> messageMap = record.getValue();
VoucherOrder voucherOrder = BeanUtil.fillBeanWithMap(messageMap, new VoucherOrder(), true);
handleVoucherOrder(voucherOrder);
// 4、ACK确认 SACK stream.orders g1 id
stringRedisTemplate.opsForStream().acknowledge(queueName, "g1", record.getId());
} catch (Exception e) {
log.error("处理订单异常", e);
// 处理异常消息
handlePendingList();
}
}
}
}
private void handlePendingList() {
while (true) {
try {
// 1、从pendingList中获取订单信息 XREADGROUP GROUP g1 c1 COUNT 1 BLOCK 1000 STREAMS streams.order 0
List<MapRecord<String, Object, Object>> messageList = stringRedisTemplate.opsForStream().read(
Consumer.from("g1", "c1"),
StreamReadOptions.empty().count(1).block(Duration.ofSeconds(1)),
StreamOffset.create(queueName, ReadOffset.from("0"))
);
// 2、判断pendingList中是否有效性
if (messageList == null || messageList.isEmpty()) {
// 2.1 pendingList中没有消息,直接结束循环
break;
}
// 3、pendingList中有消息
// 将消息转成VoucherOrder对象
MapRecord<String, Object, Object> record = messageList.get(0);
Map<Object, Object> messageMap = record.getValue();
VoucherOrder voucherOrder = BeanUtil.fillBeanWithMap(messageMap, new VoucherOrder(), true);
handleVoucherOrder(voucherOrder);
// 4、ACK确认 SACK stream.orders g1 id
stringRedisTemplate.opsForStream().acknowledge(queueName, "g1", record.getId());
} catch (Exception e) {
log.error("处理订单异常", e);
// 这里不用调自己,直接就进入下一次循环,再从pendingList中取,这里只需要休眠一下,防止获取消息太频繁
try {
Thread.sleep(20);
} catch (InterruptedException ex) {
log.error("线程休眠异常", ex);
}
}
}
}
/**
* 创建订单
*
* @param voucherOrder
*/
private void handleVoucherOrder(VoucherOrder voucherOrder) {
Long userId = voucherOrder.getUserId();
RLock lock = redissonClient.getLock(RedisConstants.LOCK_ORDER_KEY + userId);
boolean isLock = lock.tryLock();
if (!isLock) {
// 索取锁失败,重试或者直接抛异常(这个业务是一人一单,所以直接返回失败信息)
log.error("一人只能下一单");
return;
}
try {
// 创建订单(使用代理对象调用,是为了确保事务生效)
proxy.createVoucherOrder(voucherOrder);
} finally {
lock.unlock();
}
}
/**
* 加载 判断秒杀券库存是否充足 并且 判断用户是否已下单 的Lua脚本
*/
private static final DefaultRedisScript<Long> SECKILL_SCRIPT;
static {
SECKILL_SCRIPT = new DefaultRedisScript<>();
SECKILL_SCRIPT.setLocation(new ClassPathResource("lua/stream-seckill.lua"));
SECKILL_SCRIPT.setResultType(Long.class);
}
/**
* VoucherOrderServiceImpl类的代理对象
* 将代理对象的作用域进行提升,方面子线程取用
*/
private IVoucherOrderService proxy;
/**
* 抢购秒杀券
*
* @param voucherId
* @return
*/
@Transactional
@Override
public Result seckillVoucher(Long voucherId) {
Long userId = ThreadLocalUtls.getUser().getId();
long orderId = redisIdWorker.nextId(SECKILL_VOUCHER_ORDER);
// 1、执行Lua脚本,判断用户是否具有秒杀资格
Long result = null;
try {
result = stringRedisTemplate.execute(
SECKILL_SCRIPT,
Collections.emptyList(),
voucherId.toString(),
userId.toString(),
String.valueOf(orderId)
);
} catch (Exception e) {
log.error("Lua脚本执行失败");
throw new RuntimeException(e);
}
if (result != null && !result.equals(0L)) {
// result为1表示库存不足,result为2表示用户已下单
int r = result.intValue();
return Result.fail(r == 2 ? "不能重复下单" : "库存不足");
}
// 2、result为0,下单成功,直接返回ok
// 索取锁成功,创建代理对象,使用代理对象调用第三方事务方法, 防止事务失效
IVoucherOrderService proxy = (IVoucherOrderService) AopContext.currentProxy();
this.proxy = proxy;
return Result.ok();
}
/**
* 创建订单
*
* @param voucherOrder
* @return
*/
@Transactional
@Override
public void createVoucherOrder(VoucherOrder voucherOrder) {
Long userId = voucherOrder.getUserId();
Long voucherId = voucherOrder.getVoucherId();
// 1、判断当前用户是否是第一单
int count = this.count(new LambdaQueryWrapper<VoucherOrder>()
.eq(VoucherOrder::getUserId, userId));
if (count >= 1) {
// 当前用户不是第一单
log.error("当前用户不是第一单");
return;
}
// 2、用户是第一单,可以下单,秒杀券库存数量减一
boolean flag = seckillVoucherService.update(new LambdaUpdateWrapper<SeckillVoucher>()
.eq(SeckillVoucher::getVoucherId, voucherId)
.gt(SeckillVoucher::getStock, 0)
.setSql("stock = stock -1"));
if (!flag) {
throw new RuntimeException("秒杀券扣减失败");
}
// 3、将订单保存到数据库
flag = this.save(voucherOrder);
if (!flag) {
throw new RuntimeException("创建秒杀券订单失败");
}
}
}
更多推荐

所有评论(0)