使用 Disruptor 优化电商秒杀(超详细 + 代码可落地案例)

目标:在高并发秒杀中,用 Redis 预减库存 + Disruptor 内存队列顺序化写库,把“并发写库”变成“顺序写库”,稳定扛住峰值,同时保证不超卖、幂等、防重复下单、可恢复
代码示例以 Spring Boot 2.x + MyBatis + Redis(StringRedisTemplate)+ MySQL 为主,MQ 部分给可选落地。


1. 秒杀痛点:为什么你会崩?

秒杀最常见的崩法:

  1. 锁竞争:synchronized / ReentrantLock 把 CPU 用在上下文切换上了
  2. DB 成瓶颈:并发 update/insert 把连接池打满,慢查询、死锁、行锁排队
  3. 过早写库:大量失败请求也进入 DB,浪费 IO
  4. 重复请求:用户狂点/重试/网络抖动导致多次请求
  5. 不可靠异步:队列积压/丢消息/消费失败没补偿

一句话:

秒杀不是业务难,是“在极端并发下只让少数成功请求进入慢路径”难。


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. 核心流程(非常重要)

  1. 请求进来先做轻校验:参数、活动时间、风控等
  2. Redis 原子预减库存(失败立刻返回)
  3. Redis 写去重 key(失败说明用户已秒过 → 回滚库存)
  4. 发布 Disruptor 事件(返回“排队中”或“受理成功”)
  5. 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. 常见坑(你大概率会踩)

  1. 队列满不补偿 → “假卖光”
  2. 去重 key 不加 TTL → 活动后一直重复失败
  3. handler 里做慢 IO → 整条队列卡死
  4. 不做 DB 唯一键兜底 → Redis 去重偶发失效会重复单
  5. bufferSize 不是 2 的幂 → 性能/正确性都有坑

14. 结论(最硬的一句话)

Redis 负责 快判生死,Disruptor 负责 顺序化慢路径,DB 只处理“可能成功”的请求。


Logo

电商企业物流数字化转型必备!快递鸟 API 接口,72 小时快速完成物流系统集成。全流程实战1V1指导,营造开放的API技术生态圈。

更多推荐