【黑马点评】Redis Stream 消息队列:异步秒杀的终极优化
1. 优化的演进之路
在秒杀业务的优化过程中,我们经历了三个阶段,每一步都是为了解决上一步的痛点:
1.1 第一阶段:同步下单(原始版本)
所有逻辑(校验、扣库存、写订单)都在主线程同步执行。
- 痛点:数据库是瓶颈,并发能力极差,响应时间长。
1.2 第二阶段:异步下单(基于 JVM 阻塞队列)
这是我们上一步刚完成的优化。我们将业务拆分:
- 主线程:利用 Lua 脚本快速判断秒杀资格,将订单放入 JVM 的
BlockingQueue阻塞队列。 - 子线程:异步从队列取出订单,慢慢写入数据库。
- 优点:极大提升了吞吐量,响应速度达到微秒级。
- 痛点(致命缺陷):
-
- 内存溢出风险 (OOM):JVM 内存有限,如果瞬时流量过大,队列可能撑爆内存。
- 数据丢失风险:
BlockingQueue是内存队列,如果服务重启或宕机,队列里的订单全没了!
1.3 第三阶段:异步下单(基于 Redis Stream)
为了解决 JVM 内存队列的数据丢失问题,我们需要引入外部的、可持久化的消息队列。Redis 刚好提供了现成的 MQ 功能。
2. Redis 消息队列的三种方案对比
在决定使用 Stream 之前,我们对比了 Redis 提供的三种实现消息队列的方式:
2.1 基于 List 结构
Redis 的 List 是双向链表,可以通过 LPUSH 和 RPOP(或 BRPOP 阻塞读)来实现队列。
- 优点:存储在 Redis 中,不受 JVM 内存限制;数据可持久化。
- 缺点:
-
- 无法避免消息丢失:当消费者
POP取出消息后,Redis 就把消息删除了。如果消费者处理时宕机,这条消息就永久丢失了。 - 功能单一:只支持单消费者,不支持多组消费。
- 无法避免消息丢失:当消费者
2.2 基于 PubSub (发布订阅)
一种广播模式,生产者发布消息,所有订阅者都能收到。
- 优点:支持多消费者消费。
- 缺点:
-
- 完全无持久化:消息“即发即失”。如果发消息时消费者下线,消息直接丢弃。
- 数据不可靠:完全无法用于秒杀这种涉及资金的业务。
2.3 基于 Stream (最终选择)
Redis 5.0 引入的 Stream 是专门为消息队列设计的数据结构,参考了 Kafka 的设计。
- 优点:
-
- 持久化:消息也会写入 RDB/AOF,重启不丢失。
- 消费者组:支持多消费者协作,加快处理速度。
- ACK 确认机制:消费者处理完必须手动确认,保证业务完成才移除消息。
- Pending List (待处理列表):如果消费者宕机,未确认的消息会保留在 Pending List 中,重启后可重新读取处理(完美解决了 List 的丢消息问题)。
结论:Stream 完美解决了 JVM 阻塞队列的“内存限制”和“数据丢失”问题。
3. 代码落地:基于 Stream 的异步秒杀
3.1 核心架构
- Lua 脚本:完成资格校验,直接将订单信息
XADD写入 Stream。 - Java 消费者:使用
XREADGROUP读取消息 -> 落库 ->ACK确认。 - 异常恢复:处理 Pending List 中的积压消息。
3.2 第一步:改造 Lua 脚本 (seckill.lua)
不再返回结果给 Java 去入队,而是脚本直接通过 XADD 命令将消息推送到 Redis Stream。
Lua
-- seckill.lua
local voucherId = ARGV[1]
local userId = ARGV[2]
local orderId = ARGV[3]
-- 1. 校验库存
local stockKey = 'seckill:stock:' .. voucherId
if(tonumber(redis.call('get', stockKey)) <= 0) then
return 1 -- 库存不足
end
-- 2. 校验一人一单
local orderKey = 'seckill:order:' .. voucherId
if(redis.call('sismember', orderKey, userId) == 1) then
return 2 -- 重复下单
end
-- 3. 扣库存 & 记录用户
redis.call('incrby', stockKey, -1)
redis.call('sadd', orderKey, userId)
-- 4. 发送消息到 Stream 🌟 (核心变化)
-- XADD stream.orders * userId ... voucherId ... id ...
redis.call('xadd', 'stream.orders', '*', 'userId', userId, 'voucherId', voucherId, 'id', orderId)
return 0
3.3 第二步:Java 消费者逻辑 (VoucherOrderServiceImpl)
这是最关键的部分,我们需要处理事务代理、消息确认以及异常消息恢复。
Java
@Service
@Slf4j
public class VoucherOrderServiceImpl extends ServiceImpl<VoucherOrderMapper, VoucherOrder> implements IVoucherOrderService {
@Resource
private RedisIdWorker redisIdWorker;
@Autowired
private StringRedisTemplate stringRedisTemplate;
@Resource
private RedissonClient redissonClient;
/**
* 🌟 技巧:注入自身代理对象
* 作用:确保子线程调用 createVoucherOrder 时事务注解 @Transactional 能生效
* 加 @Lazy 防止循环依赖报错
*/
@Resource
@Lazy
private IVoucherOrderService self;
// 线程池
private static final ExecutorService SECKILL_ORDER_EXECUTOR = Executors.newSingleThreadExecutor();
// 初始化:创建消费者组 & 启动线程
@PostConstruct
private void init() {
createGroupIfNotExist("stream.orders", "g1");
SECKILL_ORDER_EXECUTOR.submit(new VoucherOrderHandler());
}
// 健壮性代码:防止 Stream 或 Group 不存在报错
private void createGroupIfNotExist(String streamKey, String groupName) {
try {
stringRedisTemplate.opsForStream().createGroup(streamKey, groupName);
log.info("消费者组 {} 创建成功", groupName);
} catch (Exception e) {
log.warn("消费者组 {} 已存在", groupName);
}
}
// 🌟 核心任务:消费者逻辑
private class VoucherOrderHandler implements Runnable {
String queueName = "stream.orders";
String groupName = "g1";
String consumerName = "c1";
@Override
public void run() {
while (true) {
try {
// 1. 读取消息队列 (阻塞读 2秒)
// XREADGROUP GROUP g1 c1 COUNT 1 BLOCK 2000 STREAMS stream.orders >
List<MapRecord<String, Object, Object>> list = stringRedisTemplate.opsForStream().read(
Consumer.from(groupName, consumerName),
StreamReadOptions.empty().count(1).block(Duration.ofSeconds(2)),
StreamOffset.create(queueName, ReadOffset.lastConsumed())
);
// 2. 没读到,继续循环
if (list == null || list.isEmpty()) {
continue;
}
// 3. 解析消息 & 执行业务
MapRecord<String, Object, Object> record = list.get(0);
Map<Object, Object> value = record.getValue();
VoucherOrder voucherOrder = BeanUtil.fillBeanWithMap(value, new VoucherOrder(), true);
handleVoucherOrder(voucherOrder);
// 4. ACK 确认 (从 Pending List 移除)
stringRedisTemplate.opsForStream().acknowledge(queueName, groupName, record.getId());
} catch (Exception e) {
log.error("处理订单异常", e);
// 🌟 异常兜底:去处理 Pending List 里的旧消息
handlePendingList();
}
}
}
// 🌟 恢复机制:处理“已读但未确认”的死信
private void handlePendingList() {
while (true) {
try {
// 1. 读取 Pending List (ReadOffset.from("0"))
List<MapRecord<String, Object, Object>> list = stringRedisTemplate.opsForStream().read(
Consumer.from(groupName, consumerName),
StreamReadOptions.empty().count(1),
StreamOffset.create(queueName, ReadOffset.from("0"))
);
if (list == null || list.isEmpty()) {
break; // 处理完了,回到主循环
}
// 2. 解析 & 重试业务
MapRecord<String, Object, Object> record = list.get(0);
Map<Object, Object> value = record.getValue();
VoucherOrder voucherOrder = BeanUtil.fillBeanWithMap(value, new VoucherOrder(), true);
handleVoucherOrder(voucherOrder);
// 3. ACK
stringRedisTemplate.opsForStream().acknowledge(queueName, groupName, record.getId());
} catch (Exception e) {
log.error("Pending-List 处理异常", e);
try { Thread.sleep(20); } catch (InterruptedException ex) {}
}
}
}
private void handleVoucherOrder(VoucherOrder voucherOrder) {
Long userId = voucherOrder.getUserId();
RLock redisLock = redissonClient.getLock("lock:order:" + userId);
boolean isLock = redisLock.tryLock();
if (!isLock) {
log.error("不允许重复下单");
return;
}
try {
// 必须用代理对象调用事务方法
self.createVoucherOrder(voucherOrder);
} finally {
redisLock.unlock();
}
}
}
// 主线程:只负责调 Lua,非常快
@Override
public Result seckillVoucher(Long voucherId) {
Long userId = UserHolder.getUser().getId();
Long orderId = redisIdWorker.nextId("order");
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 ? "库存不足" : "不能重复下单");
}
return Result.ok(orderId);
}
@Transactional
public void createVoucherOrder(VoucherOrder voucherOrder) {
// 具体落库逻辑(一人一单校验 + 扣数据库库存 + 写入订单)
save(voucherOrder);
}
}
4.深度剖析:为什么 Redis Stream 是秒杀场景的“终极杀器”?
在代码跑通之后,我们需要跳出具体的语法细节,从架构层面思考:为什么仅仅换了一个队列,系统的稳定性就有了质的飞跃? 这里有三个核心维度的架构升级:
1. 从“有状态”到“无状态”的架构演进
- JVM 阻塞队列版本:此时的 Java 应用是有状态的(Stateful)。订单数据存在 JVM 内存中,这意味着我们不能随意重启服务,也不能随意进行水平扩容(Scale Out),因为一旦某个节点挂掉,内存里的订单就全丢了。
- Redis Stream 版本:我们将状态(订单数据)完全外置到了 Redis 中。此时的 Java 应用变成了无状态的(Stateless)。
-
- 优势:你可以随时加减 Java 服务的节点,甚至在秒杀进行中重启某个节点,都不会导致数据丢失。这是云原生架构设计的核心原则之一。
2. 真正实现了“原子性闭环”
在旧版本中,Redis 扣减库存和 Java 队列入队是两个独立的动作:
- Lua 脚本扣减 Redis 库存。
- Java 代码将订单加入阻塞队列。
风险点:如果在第 1 步成功后,Java 服务突然发生 GC 停顿或网络抖动,导致第 2 步失败,就会出现 “Redis 库存扣了,但订单没生成” 的少卖现象(即数据不一致)。
而在 Stream 版本 中,我们将 XADD(消息入队)直接写进了 Lua 脚本。
- 原子性:Redis 保证 Lua 脚本执行的原子性。“扣减库存”和“写入队列”要么同时成功,要么同时失败。这在源头上保证了数据的一致性。
3. 消费语义的升级:At-Least-Once(至少一次)
消息队列有三种投递语义:
- At-Most-Once(至多一次):像
Pub/Sub或未做处理的List。消息发出去就不管了,丢了就丢了。 - At-Least-Once(至少一次):像 Stream + ACK 机制。只要消费者没确认(ACK),消息就永远在 Pending List 里。消费者宕机重启后,依然能看见这条消息并重新处理。
- Exactly-Once(精确一次):通常需要配合幂等性设计。
Stream 的 Pending List 机制让我们在 Redis 层面实现了 At-Least-Once。配合我们在数据库层面的唯一索引(user_id + voucher_id) 实现的幂等性,我们最终构建了一个坚不可摧的可靠消费链路。
5. 总结
通过这次重构,我们完成了从 JVM 内存队列 到 Redis 分布式消息队列 的跨越:
- 安全性提升:利用 Stream 的持久化和 ACK 机制,彻底解决了服务宕机导致订单丢失的问题。
- 健壮性提升:通过
Pending List处理机制,即使消费者发生异常,也能在重启后自动恢复处理积压消息。 - 高并发:主线程全程无数据库操作,纯 Redis 交互,抗并发能力拉满。
下一篇文章,我将复盘黑马点评 Redis 高并发秒杀业务下的所有技术演进和细节。
更多推荐




所有评论(0)