Disruptor优化电商秒杀
·
使用 Disruptor 优化电商秒杀(超详细 + 代码可落地案例)
目标:在高并发秒杀中,用 Redis 预减库存 + Disruptor 内存队列顺序化写库,把“并发写库”变成“顺序写库”,稳定扛住峰值,同时保证不超卖、幂等、防重复下单、可恢复。
代码示例以 Spring Boot 2.x + MyBatis + Redis(StringRedisTemplate)+ MySQL 为主,MQ 部分给可选落地。
1. 秒杀痛点:为什么你会崩?
秒杀最常见的崩法:
- 锁竞争:synchronized / ReentrantLock 把 CPU 用在上下文切换上了
- DB 成瓶颈:并发 update/insert 把连接池打满,慢查询、死锁、行锁排队
- 过早写库:大量失败请求也进入 DB,浪费 IO
- 重复请求:用户狂点/重试/网络抖动导致多次请求
- 不可靠异步:队列积压/丢消息/消费失败没补偿
一句话:
秒杀不是业务难,是“在极端并发下只让少数成功请求进入慢路径”难。
2. 方案总览:Redis 预判 + Disruptor 顺序化
2.1 架构位置
Client
↓
Gateway/Nginx(限流、黑名单、验证码、签名)
↓
应用层 SeckillController
↓
Redis:
- 库存预减(原子)
- 去重(用户是否已下单)
↓
Disruptor(内存 RingBuffer,顺序化成功请求)
↓
DB:落订单、扣库存(可选做最终一致校验)
↓
(可选)MQ:通知、履约、积分等异步链路
2.2 为什么 Disruptor 有用
- RingBuffer + CAS:无锁/少锁的高性能队列
- 你可以把“并发写库”变成“少量线程顺序写库”
- 更重要:你把“热点竞争”移动到 Redis(O(1) 原子)+ 内存顺序处理
秒杀最核心:先判失败,成功请求才进入 Disruptor。
3. 数据模型(MySQL)
3.1 秒杀活动库存表(示例)
CREATE TABLE seckill_stock (
sku_id BIGINT PRIMARY KEY,
total_stock INT NOT NULL,
sold_stock INT NOT NULL DEFAULT 0,
version INT NOT NULL DEFAULT 0,
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
);
3.2 订单表(唯一约束防重复)
CREATE TABLE seckill_order (
id BIGINT PRIMARY KEY,
user_id BIGINT NOT NULL,
sku_id BIGINT NOT NULL,
status TINYINT NOT NULL, -- 0=INIT 1=SUCCESS 2=FAIL
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_user_sku (user_id, sku_id)
);
uk_user_sku 是最后一道防线:就算上层去重失效,DB 也会把重复插入挡下来。
4. Redis Key 设计(关键)
假设 skuId=1001、userId=2002:
- 库存:
seckill:stock:{skuId}→ int - 用户去重:
seckill:order:{skuId}:{userId}→ 1(带 TTL) - 订单结果查询:
seckill:result:{skuId}:{userId}→ orderId / FAIL / PENDING
TTL 建议
order去重 key:TTL = 活动结束时间 + 1~2 天(看业务)result:TTL = 活动结束时间 + 若干天(给用户查询)
5. 核心流程(非常重要)
- 请求进来先做轻校验:参数、活动时间、风控等
- Redis 原子预减库存(失败立刻返回)
- Redis 写去重 key(失败说明用户已秒过 → 回滚库存)
- 发布 Disruptor 事件(返回“排队中”或“受理成功”)
- Disruptor consumer 顺序执行:
- 创建订单(INSERT)
- (可选)更新 DB 库存(或者仅做对账)
- 写 result 到 Redis
- (可选)发送 MQ
你会发现:DB 只处理“可能成功”的请求,失败流量被 Redis 拦截。
6. 完整可落地 Demo(核心代码)
下面是可直接照抄落地的结构(省略 package),你按模块复制即可。
6.1 Maven 依赖
<dependencies>
<!-- web -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- redis -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<!-- mybatis -->
<dependency>
<groupId>org.mybatis.spring.boot</groupId>
<artifactId>mybatis-spring-boot-starter</artifactId>
<version>2.3.1</version>
</dependency>
<!-- mysql driver -->
<dependency>
<groupId>com.mysql</groupId>
<artifactId>mysql-connector-j</artifactId>
<scope>runtime</scope>
</dependency>
<!-- disruptor -->
<dependency>
<groupId>com.lmax</groupId>
<artifactId>disruptor</artifactId>
<version>3.4.4</version>
</dependency>
<!-- lombok(可选) -->
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
7. Redis Lua:原子预减 + 结果写入(可选但推荐)
不用 Lua 也能做,但 Lua 可以把“读库存 + 判断 + 扣减”合成一次原子操作,减少 RTT。
7.1 预减库存脚本(stock > 0 才扣)
-- KEYS[1] = stockKey
-- ARGV[1] = deduct (e.g. 1)
local stock = tonumber(redis.call('GET', KEYS[1]))
if not stock then
return -2 -- NO_STOCK_KEY
end
if stock <= 0 then
return -1 -- OUT_OF_STOCK
end
redis.call('DECRBY', KEYS[1], tonumber(ARGV[1]))
return stock - tonumber(ARGV[1]) -- remaining
8. Controller + Service:入口尽量轻
8.1 返回值协议(建议)
ACCEPTED:已进入队列(不代表最终成功)OUT_OF_STOCK:没库存DUPLICATE:重复秒杀BUSY:队列满/系统保护
8.2 DTO
public class SeckillAcceptResult {
private String code; // ACCEPTED / OUT_OF_STOCK / DUPLICATE / BUSY / ALREADY
private String msg;
private String data; // 可放 orderId / FAIL reason / PENDING
public static SeckillAcceptResult accepted() {
return of("ACCEPTED", "已受理,排队中", "PENDING");
}
public static SeckillAcceptResult outOfStock(String msg) {
return of("OUT_OF_STOCK", msg, null);
}
public static SeckillAcceptResult duplicate() {
return of("DUPLICATE", "重复秒杀", null);
}
public static SeckillAcceptResult systemBusy(String msg) {
return of("BUSY", msg, null);
}
public static SeckillAcceptResult already(String existed) {
return of("ALREADY", "已有结果", existed);
}
private static SeckillAcceptResult of(String code, String msg, String data) {
SeckillAcceptResult r = new SeckillAcceptResult();
r.code = code; r.msg = msg; r.data = data;
return r;
}
// getter/setter 省略
}
8.3 SeckillService(入口逻辑)
@Service
public class SeckillService {
private static final String STOCK_KEY_PREFIX = "seckill:stock:";
private static final String ORDER_KEY_PREFIX = "seckill:order:";
private static final String RESULT_KEY_PREFIX = "seckill:result:";
@Resource private StringRedisTemplate redis;
@Resource private SeckillDisruptorPublisher publisher;
private final DefaultRedisScript<Long> stockPreDeductScript;
public SeckillService() {
stockPreDeductScript = new DefaultRedisScript<>();
stockPreDeductScript.setResultType(Long.class);
stockPreDeductScript.setScriptText(
"local stock = tonumber(redis.call('GET', KEYS[1])) " +
"if not stock then return -2 end " +
"if stock <= 0 then return -1 end " +
"redis.call('DECRBY', KEYS[1], tonumber(ARGV[1])) " +
"return stock - tonumber(ARGV[1])"
);
}
public SeckillAcceptResult seckill(Long userId, Long skuId) {
String stockKey = STOCK_KEY_PREFIX + skuId;
String orderKey = ORDER_KEY_PREFIX + skuId + ":" + userId;
String resultKey = RESULT_KEY_PREFIX + skuId + ":" + userId;
// 1) 幂等查询(已有结果直接返回)
String existed = redis.opsForValue().get(resultKey);
if (existed != null) {
return SeckillAcceptResult.already(existed);
}
// 2) Redis 原子预减库存
Long remaining = redis.execute(
stockPreDeductScript,
Collections.singletonList(stockKey),
"1"
);
if (remaining == null) return SeckillAcceptResult.systemBusy("redis null");
if (remaining == -2) return SeckillAcceptResult.outOfStock("stock key missing");
if (remaining == -1) return SeckillAcceptResult.outOfStock("out of stock");
// 3) 去重:SETNX
Boolean first = redis.opsForValue().setIfAbsent(orderKey, "1", Duration.ofDays(2));
if (Boolean.FALSE.equals(first)) {
// 重复:补偿库存
redis.opsForValue().increment(stockKey, 1);
return SeckillAcceptResult.duplicate();
}
// 4) 结果先写 PENDING
redis.opsForValue().set(resultKey, "PENDING", Duration.ofDays(2));
// 5) 投递 Disruptor(满了就拒绝并补偿)
boolean ok = publisher.tryPublish(userId, skuId);
if (!ok) {
redis.opsForValue().increment(stockKey, 1);
redis.delete(orderKey);
redis.opsForValue().set(resultKey, "FAIL:BUSY", Duration.ofHours(1));
return SeckillAcceptResult.systemBusy("ring buffer full");
}
return SeckillAcceptResult.accepted();
}
public String queryResult(Long userId, Long skuId) {
return redis.opsForValue().get(RESULT_KEY_PREFIX + skuId + ":" + userId);
}
}
8.4 Controller
@RestController
@RequestMapping("/seckill")
public class SeckillController {
@Resource private SeckillService seckillService;
@PostMapping("/{skuId}")
public SeckillAcceptResult seckill(@PathVariable Long skuId,
@RequestParam Long userId) {
return seckillService.seckill(userId, skuId);
}
@GetMapping("/{skuId}/result")
public String result(@PathVariable Long skuId,
@RequestParam Long userId) {
return seckillService.queryResult(userId, skuId);
}
}
9. Disruptor:事件模型 + 发布器 + 消费者
9.1 事件对象(RingBuffer 复用对象,避免 GC)
public class SeckillEvent {
private long userId;
private long skuId;
private long orderId;
public void set(long userId, long skuId, long orderId) {
this.userId = userId;
this.skuId = skuId;
this.orderId = orderId;
}
public long getUserId() { return userId; }
public long getSkuId() { return skuId; }
public long getOrderId() { return orderId; }
}
9.2 ID 生成器(简单版,生产建议用 Segment/雪花)
@Component
public class IdGenerator {
private final AtomicLong seq = new AtomicLong(1_000_000);
public long nextId() {
return seq.getAndIncrement();
}
}
9.3 发布器:非阻塞投递(满了就失败)
@Component
public class SeckillDisruptorPublisher {
private volatile RingBuffer<SeckillEvent> ringBuffer;
@Resource private IdGenerator idGenerator;
public void init(RingBuffer<SeckillEvent> ringBuffer) {
this.ringBuffer = ringBuffer;
}
public boolean tryPublish(Long userId, Long skuId) {
long sequence;
try {
sequence = ringBuffer.tryNext();
} catch (InsufficientCapacityException e) {
return false;
}
try {
SeckillEvent event = ringBuffer.get(sequence);
event.set(userId, skuId, idGenerator.nextId());
} finally {
ringBuffer.publish(sequence);
}
return true;
}
}
9.4 Disruptor 配置(Spring Boot 初始化)
@Configuration
public class SeckillDisruptorConfig {
@Resource private SeckillEventHandler handler;
@Resource private SeckillDisruptorPublisher publisher;
@Bean(destroyMethod = "shutdown")
public Disruptor<SeckillEvent> seckillDisruptor() {
int bufferSize = 1 << 20; // 1048576
ThreadFactory tf = r -> {
Thread t = new Thread(r);
t.setName("seckill-disruptor-" + t.getId());
t.setDaemon(true);
return t;
};
Disruptor<SeckillEvent> disruptor = new Disruptor<>(
SeckillEvent::new,
bufferSize,
tf,
ProducerType.MULTI,
new YieldingWaitStrategy()
);
disruptor.handleEventsWith(handler);
disruptor.start();
publisher.init(disruptor.getRingBuffer());
return disruptor;
}
}
10. Consumer:顺序写库 + 结果落 Redis(核心)
10.1 MyBatis Mapper
@Mapper
public interface SeckillOrderMapper {
int insertOrder(@Param("id") long id,
@Param("userId") long userId,
@Param("skuId") long skuId,
@Param("status") int status);
}
@Mapper
public interface SeckillStockMapper {
int deductSold(@Param("skuId") long skuId,
@Param("deduct") int deduct);
}
XML 示例
<!-- SeckillOrderMapper.xml -->
<insert id="insertOrder">
INSERT INTO seckill_order(id, user_id, sku_id, status)
VALUES (#{id}, #{userId}, #{skuId}, #{status})
</insert>
<!-- SeckillStockMapper.xml -->
<update id="deductSold">
UPDATE seckill_stock
SET sold_stock = sold_stock + #{deduct}
WHERE sku_id = #{skuId}
AND sold_stock + #{deduct} <= total_stock
</update>
10.2 Handler(事务 + 补偿)
Disruptor consumer 线程里,建议用
TransactionTemplate控事务。
@Component
public class SeckillEventHandler implements EventHandler<SeckillEvent> {
private static final String RESULT_KEY_PREFIX = "seckill:result:";
private static final String ORDER_KEY_PREFIX = "seckill:order:";
private static final String STOCK_KEY_PREFIX = "seckill:stock:";
@Resource private TransactionTemplate transactionTemplate;
@Resource private SeckillOrderMapper orderMapper;
@Resource private SeckillStockMapper stockMapper;
@Resource private StringRedisTemplate redis;
@Override
public void onEvent(SeckillEvent event, long sequence, boolean endOfBatch) {
long userId = event.getUserId();
long skuId = event.getSkuId();
long orderId = event.getOrderId();
String resultKey = RESULT_KEY_PREFIX + skuId + ":" + userId;
String orderKey = ORDER_KEY_PREFIX + skuId + ":" + userId;
String stockKey = STOCK_KEY_PREFIX + skuId;
try {
Boolean ok = transactionTemplate.execute(status -> {
// 1) 插入订单(DB 唯一键兜底幂等)
try {
orderMapper.insertOrder(orderId, userId, skuId, 1);
} catch (Exception dup) {
// uk_user_sku 冲突:幂等成功
return true;
}
// 2) DB 库存扣减兜底(推荐生产保留)
int updated = stockMapper.deductSold(skuId, 1);
if (updated != 1) {
status.setRollbackOnly();
return false;
}
return true;
});
if (Boolean.TRUE.equals(ok)) {
redis.opsForValue().set(resultKey, String.valueOf(orderId), Duration.ofDays(2));
} else {
// DB 不够:补偿 Redis
redis.opsForValue().increment(stockKey, 1);
redis.delete(orderKey);
redis.opsForValue().set(resultKey, "FAIL:DB_STOCK", Duration.ofHours(1));
}
} catch (Exception e) {
// 异常兜底:补偿 Redis
redis.opsForValue().increment(stockKey, 1);
redis.delete(orderKey);
redis.opsForValue().set(resultKey, "FAIL:EXCEPTION", Duration.ofHours(1));
}
}
}
11. 系统保护:队列满、慢消费、雪崩怎么办?
- 发布端用
tryNext(),满了就拒绝(并补偿 Redis) - handler 内禁止调用外部服务(HTTP/远程 RPC),把那些事丢给 MQ
- bufferSize 要够大,但别无限大(内存会爆)
12. 多 SKU 提升吞吐:分片多个 Disruptor(推荐思路)
如果 SKU 很多、并发更高,单消费线程可能成为瓶颈:
- 建 N 个 disruptor(N≈CPU 核数)
idx = (int)(skuId % N)路由到固定分片- 同 SKU 永远进同一个队列 → 顺序天然保证
13. 常见坑(你大概率会踩)
- 队列满不补偿 → “假卖光”
- 去重 key 不加 TTL → 活动后一直重复失败
- handler 里做慢 IO → 整条队列卡死
- 不做 DB 唯一键兜底 → Redis 去重偶发失效会重复单
- bufferSize 不是 2 的幂 → 性能/正确性都有坑
14. 结论(最硬的一句话)
Redis 负责 快判生死,Disruptor 负责 顺序化慢路径,DB 只处理“可能成功”的请求。
更多推荐




所有评论(0)