高并发秒杀“零超卖”解决方案:Redisson + Kafka + MySQL 最终一致性实践
本文档涵盖从理论到代码的完整实现,用于解决电商秒杀场景下的库存超卖问题,保证数据最终一致性。
目录
- 背景与挑战
- 整体架构模式
- 核心流程(时序)
- 代码实现
- 4.1 Redis Lua 脚本(预加载)
- 4.2 核心下单服务(Redisson 锁 + 本地消息表)
- 4.3 Kafka 消费者(MySQL 乐观锁扣减)
- 4.4 补偿定时任务(保证消息可靠)
- 4.5 凌晨对账任务(最终一致性修复)
- 数据一致性保障机制
- 关键配置参考
- 注意事项与防坑指南
- 总结
1. 背景与挑战
在高并发秒杀场景下,核心痛点在于:
- 并发冲突:大量请求同时修改同一库存记录,导致数据库行锁竞争剧烈。
- 超卖风险:若不严格保证“检查库存”与“扣减库存”的原子性,则会出现库存为负。
- 性能与一致性权衡:强一致性(如分布式事务)性能极差,需采用最终一致性方案。
2. 整体架构模式
我们采用“缓存预扣 + 异步落库 + 补偿兜底”的经典架构,各组件职责如下:
| 组件 | 角色 | 核心作用 |
|---|---|---|
| Redisson | 分布式锁 | 防止同一用户重复提交(防重入),降低无效并发 |
| Redis + Lua | 流量闸门 | 原子扣减缓存库存,拦截大部分超卖请求,保护数据库 |
| Kafka | 异步削峰 | 将下单请求异步化,平滑流量峰值,并保证消息可靠性 |
| MySQL | 最终权威 | 使用乐观锁(version或stock >= num)作为最终裁决,保证物理库存准确 |
| 本地消息表 | 可靠性保障 | 保证 Kafka 消息不丢失,同时支持幂等消费 |
| 定时补偿 + 对账 | 兜底机制 | 处理异常情况(如消息丢失、缓存不一致),实现最终一致性 |
3. 核心流程(时序)
- 用户请求→ 获取 Redisson 分布式锁(Key =
userId:productId),防止重复点击。 - 执行 Redis Lua原子扣减缓存库存(
stock:productId)。- 若扣减失败 → 直接返回“库存不足”,释放锁。
- 若扣减成功 → 进入下一步。
- 本地事务:向 MySQL 插入订单记录和本地消息日志(状态=0 待发送)。
- 发送 Kafka(异步),若失败不阻塞,依赖后续补偿任务。
- 返回用户“下单成功,请等待支付”。
- Kafka 消费者拉取消息,开启 MySQL 事务:
- 幂等性检查(查询消息状态,若已处理则跳过)。
- 执行乐观锁 SQL 更新物理库存(
UPDATE product SET stock = stock - #{num}, version = version + 1 WHERE id = #{id} AND stock >= #{num})。 - 若更新成功 → 插入订单详情,更新消息状态为“已消费(2)”,提交事务。
- 若更新失败 → 记录失败,发送补偿消息(将 Redis 库存加回),并通知用户下单失败。
- 补偿定时任务:每分钟扫描状态为“待发送”或“已发送但未确认”的旧消息,重新发送 Kafka。
- 凌晨对账:对比 Redis 缓存库存与 MySQL 物理库存,若不一致则以 MySQL 为准修复缓存。
4. 代码实现
环境:Spring Boot 3.x + MyBatis-Plus + Redisson + Kafka(
spring-kafka)
以下代码仅展示核心逻辑,请按实际业务调整。
4.1 Redis Lua 脚本(预加载)
package com.example.seckill.script; import org.springframework.stereotype.Component; @Component public class StockLuaScript { // 扣减脚本:KEYS[1]=库存Key,ARGV[1]=购买数量 // 返回 1 成功,0 失败 public static final String DECREASE_STOCK = "if redis.call('get', KEYS[1]) >= tonumber(ARGV[1]) then " + " redis.call('decrby', KEYS[1], ARGV[1]) " + " return 1 " + "else " + " return 0 " + "end"; }4.2 核心下单服务(Redisson 锁 + 本地消息表)
package com.example.seckill.service; import com.alibaba.fastjson.JSON; import com.example.seckill.entity.LocalMessageLog; import com.example.seckill.mapper.LocalMessageLogMapper; import org.redisson.api.RLock; import org.redisson.api.RedissonClient; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.data.redis.core.script.DefaultRedisScript; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import lombok.extern.slf4j.Slf4j; import java.util.Collections; import java.util.UUID; import java.util.concurrent.TimeUnit; @Service @Slf4j public class SeckillOrderService { @Autowired private RedissonClient redissonClient; @Autowired private StringRedisTemplate redisTemplate; @Autowired private LocalMessageLogMapper logMapper; @Autowired private KafkaTemplate<String, String> kafkaTemplate; // 下单入口 public String createOrder(Long userId, Long productId, Integer num) { String lockKey = "lock:seckill:" + userId + ":" + productId; String stockKey = "stock:" + productId; RLock lock = redissonClient.getLock(lockKey); boolean locked = false; try { // 1. 尝试获取锁,最多等待0秒,持有200毫秒(防死等) locked = lock.tryLock(0, 200, TimeUnit.MILLISECONDS); if (!locked) { return "请勿重复点击,稍后再试"; } // 2. Redis Lua 原子扣减 Long result = redisTemplate.execute( new DefaultRedisScript<>(StockLuaScript.DECREASE_STOCK, Long.class), Collections.singletonList(stockKey), num.toString() ); if (result == null || result == 0) { return "库存不足,秒杀失败"; } // 3. 构造订单和本地消息日志 String orderId = UUID.randomUUID().toString(); LocalMessageLog logEntity = new LocalMessageLog(); logEntity.setOrderId(orderId); logEntity.setProductId(productId); logEntity.setUserId(userId); logEntity.setNum(num); logEntity.setStatus(0); // 0=待发送,1=已发送,2=已消费 // 4. 本地事务:保存日志(同时保存订单,此处省略订单insert) // 注意:实际中需将 insert 放在 @Transactional 方法中 saveOrderAndLog(logEntity); // 内部使用 @Transactional // 5. 发送 Kafka(异步,失败不阻塞) kafkaTemplate.send("seckill-order-topic", orderId, JSON.toJSONString(logEntity)); // 可选:异步更新消息状态为1(但依赖补偿兜底,可省略) return "下单成功,订单号:" + orderId + ",请等待支付"; } catch (Exception e) { // 本地事务失败,必须回滚 Redis 库存 log.error("本地事务异常,执行Redis回滚", e); redisTemplate.opsForValue().increment(stockKey, num); return "系统繁忙,请稍后重试"; } finally { if (locked && lock.isHeldByCurrentThread()) { lock.unlock(); } } } @Transactional(rollbackFor = Exception.class) public void saveOrderAndLog(LocalMessageLog logEntity) { // 插入订单表(略) // orderMapper.insert(order); logMapper.insert(logEntity); } }4.3 Kafka 消费者(MySQL 乐观锁扣减)
package com.example.seckill.consumer; import com.alibaba.fastjson.JSON; import com.example.seckill.entity.LocalMessageLog; import com.example.seckill.mapper.LocalMessageLogMapper; import com.example.seckill.mapper.OrderMapper; import com.example.seckill.mapper.ProductMapper; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Component; @Component @Slf4j public class SeckillOrderConsumer { @Autowired private OrderMapper orderMapper; @Autowired private ProductMapper productMapper; @Autowired private LocalMessageLogMapper logMapper; @Autowired private KafkaTemplate<String, String> kafkaTemplate; @KafkaListener(topics = "seckill-order-topic", groupId = "seckill-group") public void consume(ConsumerRecord<String, String> record, Acknowledgment ack) { String orderId = record.key(); LocalMessageLog logEntity = JSON.parseObject(record.value(), LocalMessageLog.class); // 1. 幂等性校验 LocalMessageLog existing = logMapper.selectByOrderId(orderId); if (existing == null || existing.getStatus() == 2) { ack.acknowledge(); return; // 已处理,跳过 } // 2. MySQL 乐观锁扣减 int updateRows = productMapper.decreaseStockWithOptimisticLock( logEntity.getProductId(), logEntity.getNum() ); // Mapper SQL: // UPDATE product SET stock = stock - #{num}, version = version + 1 // WHERE id = #{id} AND stock >= #{num} if (updateRows > 0) { // 扣减成功:生成订单 Order order = new Order(); order.setOrderId(orderId); order.setUserId(logEntity.getUserId()); order.setStatus(1); // 待支付 orderMapper.insert(order); // 更新消息状态为已消费 logMapper.updateStatus(orderId, 2); ack.acknowledge(); log.info("订单落库成功: {}", orderId); } else { // 扣减失败:触发补偿 log.error("物理库存不足,订单失败,触发补偿: {}", orderId); logMapper.updateStatus(orderId, -1); // 失败状态 // 发送补偿消息,将 Redis 库存加回 String compensationMsg = "{\"productId\":" + logEntity.getProductId() + ",\"num\":" + logEntity.getNum() + "}"; kafkaTemplate.send("compensation-topic", orderId, compensationMsg); ack.acknowledge(); // 推送通知用户下单失败(略) } } }4.4 补偿定时任务(保证消息可靠)
package com.example.seckill.task; import com.alibaba.fastjson.JSON; import com.example.seckill.entity.LocalMessageLog; import com.example.seckill.mapper.LocalMessageLogMapper; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.util.List; @Component @Slf4j public class MessageCompensationTask { @Autowired private LocalMessageLogMapper logMapper; @Autowired private KafkaTemplate<String, String> kafkaTemplate; // 每5分钟执行一次,扫描状态为0(待发送)或1(已发送但未确认)且创建时间超过5分钟的消息 @Scheduled(cron = "0 0/5 * * * ?") public void retryUnsentMessages() { List<LocalMessageLog> pendingList = logMapper.selectPendingMessages(); // status in (0,1) and create_time < now-5min for (LocalMessageLog log : pendingList) { try { kafkaTemplate.send("seckill-order-topic", log.getOrderId(), JSON.toJSONString(log)); // 若发送成功,更新状态为1(已发送) logMapper.updateStatus(log.getOrderId(), 1); log.info("补偿重发成功: {}", log.getOrderId()); } catch (Exception e) { log.error("补偿重发失败,待下次重试: {}", log.getOrderId(), e); } } } }4.5 凌晨对账任务(最终一致性修复)
package com.example.seckill.task; import com.example.seckill.entity.Product; import com.example.seckill.mapper.ProductMapper; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.util.List; @Component @Slf4j public class StockReconciliationTask { @Autowired private ProductMapper productMapper; @Autowired private StringRedisTemplate redisTemplate; // 凌晨3点执行 @Scheduled(cron = "0 0 3 * * ?") public void reconcile() { List<Product> products = productMapper.selectAll(); for (Product p : products) { String stockKey = "stock:" + p.getId(); String redisVal = redisTemplate.opsForValue().get(stockKey); Integer redisStock = redisVal == null ? 0 : Integer.parseInt(redisVal); if (!redisStock.equals(p.getStock())) { log.error("发现不一致!Product: {}, Redis: {}, MySQL: {}, 强制修复", p.getId(), redisStock, p.getStock()); // 以 MySQL 为准覆盖 Redis redisTemplate.opsForValue().set(stockKey, String.valueOf(p.getStock())); // 可发送告警通知人工介入 } } } }5. 数据一致性保障机制
为了在异步链路中保证最终一致性,我们采用了以下三道防线:
- 本地消息表 + 补偿重试:确保 Kafka 消息不丢失,即使发送失败也有重试机制。
- 消费幂等:通过订单号(
orderId)查询消息状态,避免重复消费导致库存多扣。 - 反向补偿:当 MySQL 扣减失败时,发送补偿消息将 Redis 库存加回,并通知用户。
- 定期对账:每日凌晨比对 Redis 与 MySQL 库存,自动修复差异,并记录告警。
这套机制保证了在极端情况下(如网络分区、服务重启),数据最终会趋于一致。
6. 关键配置参考
application.yml(部分)
spring: kafka: bootstrap-servers: localhost:9092 producer: retries: 3 acks: all consumer: group-id: seckill-group enable-auto-commit: false auto-offset-reset: latest listener: ack-mode: manual redis: host: localhost port: 6379 datasource: url: jdbc:mysql://localhost:3306/seckill?useSSL=false&allowMultiQueries=true driver-class-name: com.mysql.cj.jdbc.Driver username: root password: 123456本地消息表 DDL
CREATE TABLE `local_message_log` ( `id` bigint(20) NOT NULL AUTO_INCREMENT, `order_id` varchar(64) NOT NULL COMMENT '订单号', `user_id` bigint(20) NOT NULL, `product_id` bigint(20) NOT NULL, `num` int(11) NOT NULL, `status` tinyint(4) DEFAULT '0' COMMENT '0-待发送 1-已发送 2-已消费 -1-失败', `create_time` datetime DEFAULT CURRENT_TIMESTAMP, `update_time` datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (`id`), UNIQUE KEY `uk_order_id` (`order_id`), KEY `idx_status_create` (`status`, `create_time`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;7. 注意事项与防坑指南
- 事务边界:不要在
@Transactional中调用 Kafka 发送,避免网络抖动导致 DB 事务回滚。应先提交事务,再异步发送,失败由补偿任务处理。 - Redisson 锁续期:业务执行超过 30 秒时,Redisson 会自动续期(看门狗),无需担心锁提前释放。
- Kafka 手动提交:必须使用
Acknowledgment.acknowledge()并关闭自动提交,确保消费成功后才提交 Offset,防止消息丢失。 - 乐观锁 SQL 条件:务必加上
stock >= #{num},这是防超卖的数据库最后防线。 - Redis 回滚:若本地事务(DB)失败,务必立即将 Redis 库存加回,否则会造成缓存与 DB 不一致。
8. 总结
本方案通过Redisson 防重、Redis Lua 防超、Kafka 异步削峰、MySQL 乐观锁兜底、本地消息表保可靠、定时对账修数据,构建了一套高并发下零超卖的最终一致性体系。各层职责清晰,性能与数据安全得到平衡。
实际生产部署时,请根据自身业务调整超时参数、重试次数和监控告警,以便及时发现并处理异常。