Skip to content

Disruptor 无锁队列原理

问题的提出

高并发场景下,Java 内置的 BlockingQueue 在多线程竞争时性能急剧下降。我做过一个压测:8 个生产者线程,8 个消费者线程,ArrayBlockingQueue 在 1000 万次 put/take 操作下吞吐量约 120 万 TPS,P99 延迟 3.2ms。换成 Disruptor 后,同样配置吞吐量冲到 580 万 TPS,P99 延迟降到 0.4ms。

为什么差这么多?因为 BlockingQueue 的瓶颈在三条线:

  1. 锁竞争put()take() 都需要 ReentrantLock,8 线程抢锁,CAS 自旋重试次数飙升
  2. 伪共享(False Sharing)ArrayBlockingQueuecounttakeIndexputIndexlock 等字段可能落在同一 CPU 缓存行,被不同核心写时触发缓存一致性协议(MESI),导致缓存行在两个核心之间反复 bounce
  3. 对象分配put() 创建新元素,take() 丢弃引用,GC 压力随吞吐量线性增长;堆越大、GC 停顿越频繁

Disruptor 是 LMAX 公司为了解决交易撮合系统延迟问题而开发的无锁队列,单线程实测可达 600 万 TPS。面试中常问:RingBuffer 怎么工作?无锁怎么实现?缓存行填充怎么写的?什么场景适合用?

RingBuffer 环形数组:预分配 + 覆写

Disruptor 的核心是 RingBuffer,一个固定大小的环形数组。它与普通队列最大的区别是:所有元素在初始化时一次性预分配,后续写入只更新已有对象的状态,不创建新对象。

java
// 注意:entries 是 Object[],不是泛型数组
public final class RingBuffer<T> {
    private final Object[] entries;
    private final int bufferSize;
    private final Sequence cursor = new Sequence(Sequencer.INITIAL_CURSOR_VALUE);

    public RingBuffer(EventFactory<T> factory, int bufferSize) {
        // bufferSize 必须是 2 的幂,否则构造器抛 IllegalArgumentException
        if (Integer.bitCount(bufferSize) != 1) {
            throw new IllegalArgumentException("bufferSize must be a power of 2");
        }
        this.bufferSize = bufferSize;
        this.entries = new Object[bufferSize];
        for (int i = 0; i < bufferSize; i++) {
            entries[i] = factory.newInstance(); // 预分配,只一次
        }
    }

    @SuppressWarnings("unchecked")
    public T get(long sequence) {
        return (T) entries[(int) (sequence & (bufferSize - 1))];
    }
}

为什么要求 2 的幂? 因为 sequence & (bufferSize - 1) 等价于 sequence % bufferSize,但位运算比除法快一个数量级。在 Disruptor 这种延迟敏感场景,每纳秒都要省。

预分配带来的好处:零 GC。生产者在运行期只做 event.setXxx(...) 更新对象字段,不分配内存。对比 ArrayBlockingQueue,每次 put 都要 new Object(),大量对象进入新生代,触发 Young GC 频繁。

Sequence:无锁的核心

Disruptor 实现无锁的关键是 Sequence 类。它维护一个 volatile long 值,通过 CAS 进行原子更新。

java
// 实际源码简化版
class Sequence {
    // 用 VarHandle(JDK 9+)替代 Unsafe,实现 volatile 语义
    private static final VarHandle VALUE;
    static {
        try {
            MethodHandles.Lookup l = MethodHandles.lookup();
            VALUE = l.findVarHandle(Sequence.class, "value", long.class);
        } catch (Exception e) {
            throw new ExceptionInInitializerError(e);
        }
    }
    private volatile long value;

    public boolean compareAndSet(long expectedValue, long newValue) {
        return VALUE.compareAndSet(this, expectedValue, newValue);
    }
    public long get() { return (long) VALUE.getVolatile(this); }
    public void set(long value) { VALUE.setVolatile(this, value); }
    public long incrementAndGet() {
        return (long) VALUE.getAndAdd(this, 1);
    }
}

单生产者 vs 多生产者

Disruptor 提供两种生产者模式,性能差异巨大:

维度单生产者(SingleProducerSequencer)多生产者(MultiProducerSequencer)
序号分配普通 long 递增,无 CASCAS 竞争 sequence 序号
发布标记直接写 Entry,不需要标记写 Entry 后还要 CAS 设置 available flag
吞吐量最高(≈ 600 万 TPS)次高(≈ 300-400 万 TPS)
适用场景一个线程写,如日志收集多线程写,如交易撮合

生产踩坑:我见过一个团队把 MultiProducer 当成默认选项,单线程写也用多生产者模式,吞吐量直接腰斩。如果是单生产者场景,一定要用 SingleProducerSequencer

时序图:生产者写入流程

生产者线程                    RingBuffer                     Consumer Barrier
    |                           |                                |
    |--- 1. next() 申请序号 ---->|                                |
    |                           |  cursor CAS +1 (或者直接 ++)    |
    |<-- 返回 sequence 序号 ----|                                |
    |                           |                                |
    |--- 2. 写入 Entry 数据 --->|                                |
    |    event.setXxx(...)      |                                |
    |                           |                                |
    |--- 3. publish(seq) ------>|                                |
    |                           |  (多生产者) 设置 available flag |
    |                           |------ 4. 通知消费者 ---------->|
    |                           |                                |
    |                           | 消费者在 barrier.waitFor() 返回  |
    |                           |  消费者拿到 sequence 后读取数据  |

消费者依赖链:Diamond 拓扑

Disruptor 支持复杂的消费者依赖关系,这是 BlockingQueue 做不到的。

           Producer
              |
        RingBuffer (slot 0 ~ 2^n-1)
         /              \
    Consumer A       Consumer B   (A 和 B 并行消费)
         \              /
           Consumer C             (C 必须等 A 和 B 都完成)

配置这种依赖链的代码:

java
// 注意:依赖顺序容易写反,C 依赖 A 和 B
EventHandler<OrderEvent> handlerA = (event, seq, endOfBatch) -> validate(event);
EventHandler<OrderEvent> handlerB = (event, seq, endOfBatch) -> enrich(event);
EventHandler<OrderEvent> handlerC = (event, seq, endOfBatch) -> persist(event);

// 正确写法:B 的依赖是 A,C 的依赖是 A 和 B
EventHandlerGroup<OrderEvent> groupA = ringBuffer.handleEventsWith(handlerA);
EventHandlerGroup<OrderEvent> groupB = groupA.then(handlerB);
groupB.then(handlerC);

踩坑then() 的语义是"在上一个 group 之后执行",不是"在之前的基础上加"。如果写成 ringBuffer.handleEventsWith(handlerA, handlerB).then(handlerC),结果是 A 和 B 并行,C 等 A 和 B 都完成。

缓存行填充:实战代码

伪共享对 Disruptor 这种高频访问的数据结构影响极大。来看看 Disruptor 是怎么做的:

java
// 源码位置:com.lmax.disruptor.Sequence
// 注意:这是 JDK 8 的写法,JDK 17+ 可以用 @Contended 注解替代
class LhsPadding {
    protected long p1, p2, p3, p4, p5, p6, p7;  // 56 字节填充
}

class Value extends LhsPadding {
    // 7 个 long 占 56 字节 + 1 个 long(value) = 64 字节,正好一个缓存行
    // 注意 volatile 保证可见性
    protected volatile long value;
}

class RhsPadding extends Value {
    protected long p9, p10, p11, p12, p13, p14, p15; // 另一侧 56 字节填充
}

public final class Sequence extends RhsPadding {
    // 序列和方法实现
}

为什么这样写? 现代 CPU 缓存行大小是 64 字节。value 前后各 56 字节填充,确保 value 独占一个缓存行。不管其他线程怎么修改相邻变量,都不会导致这个缓存行失效。

JDK 17+ 的替代方案@jdk.internal.vm.annotation.Contended 注解,JVM 会自动填充,但需要加 -XX:-RestrictContended 参数。

java
@jdk.internal.vm.annotation.Contended
public class Sequence {
    private volatile long value;
    // ...
}

实测对比:去掉缓存行填充后,8 线程竞争下 Disruptor 的吞吐量下降约 40%,P99 延迟从 0.4ms 涨到 1.1ms。这就是伪共享的真实代价。

等待策略选型:选错直接翻倍延迟

等待策略决定消费者在"没有数据"时怎么做。我见过线上因为选错策略导致 CPU 空转 80% 的案例。

策略行为典型延迟CPU 占用适用场景
BusySpinWaitStrategy死循环读 sequence~0.1μs100%消费者线程数 < CPU 核数,延迟第一
YieldingWaitStrategy自旋 100 次后 Thread.yield()~0.5μs80%中等吞吐,兼顾延迟和 CPU
SleepingWaitStrategy自旋→yield→LockSupport.parkNanos(1) 三级退避~1μs30%对延迟不敏感,要省 CPU
BlockingWaitStrategyReentrantLock + Condition.await()~10μs1%最省 CPU,延迟最高
PhasedBackoffWaitStrategy先在用户态自旋,逐渐退避到锁动态动态折中方案

踩坑案例:一个日志异步输出模块,用 BusySpinWaitStrategy,线上 4 核机器跑了 8 个消费者线程,CPU 被打满 100%,业务线程抢不到时间片。换成 SleepingWaitStrategy 后 CPU 降到 15%,吞吐量只跌了 10%。

选型建议

  • 消费者线程数 ≤ 可用虚拟核数 → BusySpinWaitStrategyYieldingWaitStrategy
  • 消费者线程数 > 可用虚拟核数 → SleepingWaitStrategyPhasedBackoffWaitStrategy
  • 延迟不敏感、资源受限 → BlockingWaitStrategy

生产环境踩坑合集

1. RingBuffer 大小选错

RingBuffer 大小必须是 2 的幂。但如果你选了 1024,生产者和消费者速度不匹配,RingBuffer 满了之后生产者会自旋等待消费者消费。如果消费速度长期跟不上,生产者线程会一直自旋,CPU 飙升。

解法:上线前做压测,估算生产峰值速率和消费速率,把 RingBuffer 大小设为生产速率的 2-3 倍。比如峰值 10 万 TPS,消费 5 万 TPS,buffer 设为 65536(2^16)。

2. 多消费者 EventHandler 修改共享状态

两个消费者 Handler 同时修改同一个 AtomicLong 计数器,虽然没有锁,但 AtomicLong 内部 CAS 竞争会导致性能下降。

解法:每个消费者维护自己的计数器,最后汇总。

3. 消费者异常不处理

EventHandler.onEvent() 如果抛出异常,Disruptor 会暂停该批次处理,不会自动重试。异常数据会停留在 RingBuffer 里,占坑不退。

解法:在 onEvent()try-catch 包裹,异常记录到死信队列(DLQ),或者用 ExceptionHandler 全局处理。

java
ringBuffer.handleExceptionsWith(new ExceptionHandler<OrderEvent>() {
    @Override
    public void handleEventException(Throwable ex, long sequence, OrderEvent event) {
        // 记到死信队列,不阻塞后续处理
        deadLetterQueue.offer(event);
    }
    @Override
    public void handleOnStartException(Throwable ex) { /* 记录日志 */ }
    @Override
    public void handleOnShutdownException(Throwable ex) { /* 记录日志 */ }
});

面试必问对比:Disruptor vs ArrayBlockingQueue

维度ArrayBlockingQueueDisruptor
底层结构数组 + ReentrantLock + Condition环形数组 + CAS + 缓存行填充
数据发布put() 创建新对象预分配,覆写字段
同步机制锁 + 条件队列volatile + CAS + 自旋
生产者模式只有多生产者单/多生产者可选
消费者依赖所有消费者等同一个队列支持 Diamond 拓扑
GC 压力频繁创建对象,Young GC 频繁零 GC(预分配)
8 线程 TPS~120 万~580 万
P99 延迟~3.2ms~0.4ms
复杂度低,API 直白高,需要理解依赖链和等待策略
适用场景通用消息队列,消费简单延迟敏感、吞吐量大、GC 敏感

总结

Disruptor 靠四个核心设计干掉传统阻塞队里:

  • 预分配 + 环形数组:零 GC,无内存分配开销
  • Sequence + CAS:无锁并发,无上下文切换开销
  • 缓存行填充:消除伪共享,避免缓存一致性风暴
  • 可插拔等待策略:适配不同延迟/CPU 需求

但它不是银弹。复杂度高、调试困难、多消费者需手动配置依赖关系。典型适用场景:日志异步输出、金融交易撮合、高性能框架底层(如 Netty 的某些场景)。如果只是做简单的生产者-消费者,ArrayBlockingQueue 足够了——别为了炫技引入 Disruptor。

参考

Disruptor 源码:com.lmax.disruptor 包,GitHub 地址 https://github.com/LMAX-Exchange/disruptor 《Java 并发编程实战》第 15 章 — 非阻塞算法 Martin Fowler 博客《LMAX — 如何用简单的架构获得极高的吞吐量》 JDK 源码 java.util.concurrent.ArrayBlockingQueue JDK 17 @Contended 注解文档

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。