核心架构

Disruptor 内部可视为这几个模块的精密协作:

  1. RingBuffer (环形缓冲区):数据容器,预分配。
  2. Sequencer (序列器):协调生产者的并发安全与序号分配。
  3. Sequence (序列):原子性的进度指示器,解决伪共享。
  4. SequenceBarrier (序列屏障):消费者协调器,确定可消费范围。
  5. EventProcessor (事件处理器):消费者线程封装,实现批处理。
  6. WaitStrategy (等待策略):当无数据时,消费者的等待方式。


模块一:RingBuffer —— 内存预分配与无GC

为什么快? 杜绝了运行时的内存分配与回收。

RingBuffer 的核心构造很简单,就是一个定长的 Object 数组。


public final class RingBuffer<E> extends RingBufferFields<E> {
    // 构造时一次性分配所有槽位
    RingBuffer(EventFactory<E> eventFactory, Sequencer sequencer) {
        super(eventFactory, sequencer);
    }
}

abstract class RingBufferFields<E> extends RingBufferPad {
    // 真正存放数据的数组,final + 预填充
    private final Object[] entries;
    
    RingBufferFields(EventFactory<E> eventFactory, Sequencer sequencer) {
        this.entries = new Object[sequencer.getBufferSize() + 2 * BUFFER_PAD];
        // 预填充对象,未来只修改属性,不替换引用
        for (int i = 0; i < sequencer.getBufferSize(); i++) {
            entries[i] = eventFactory.newInstance();
        }
    }
}

关键点:

  • entries 数组在构造后大小不变,元素引用不变。
  • 使用时,直接从数组取对象,修改其字段,用完丢回。整个过程是对象复用,JVM 不再需要 GC 这些事件对象。
  • + 2 * BUFFER_PAD 在数组头尾添加填充,防止数组长度等元数据与元素数据在同一缓存行带来的伪共享。

模块二:Sequence —— 填充缓存行,消除伪共享

为什么快? 解决多核CPU下缓存行失效带来的巨大开销。


class Sequence extends RhsPadding {
    // 真正存储序列值的字段,使用Unsafe进行CAS和Volatile操作
    private static final Unsafe UNSAFE = Util.getUnsafe();
    private static final long VALUE_OFFSET;

    static {
        VALUE_OFFSET = UNSAFE.objectFieldOffset(Value.class.getDeclaredField("value"));
    }
    
    // 构造时将序列初始值设为-1
    public Sequence(final long initialValue) {
        UNSAFE.putOrderedLong(this, VALUE_OFFSET, initialValue);
        // putOrderedLong 是一个带StoreStore屏障的延迟写,比volatile写成本低
    }
}

// 继承的填充类,保证了值独享缓存行
class LhsPadding { protected long p1, p2, p3, p4, p5, p6, p7; }
class Value extends LhsPadding { protected volatile long value; }
class RhsPadding extends Value { protected long p9, p10, p11, p12, p13, p14, p15; }

内存布局与原理:
一个 Sequence 实例在内存中形如:[ p1..p7 ][ value ][ p9..p15 ]。

  • value 前后各填充 56 字节(7个long),加上 value 自身 8 字节共 64 字节,精确占满一个主流 CPU 缓存行。
  • 当生产者修改这个 Sequence 的 value 时,只使它自己的缓存行失效。消费者各自的 Sequence 因处于不同缓存行,不受影响,无需从主存重新加载。这直接避免了伪共享风暴。

模块三:SingleProducerSequencer —— 无锁化生产

为什么快? 单生产者场景下,用简单的写缓冲而非CAS操作来分配序号,极大降低延迟。


public final class SingleProducerSequencer extends AbstractSequencer {
    // 不共享,无竞争,普通字段
    long nextValue = -1;
    long cachedValue = -1; // 缓存可发布的最大值,避免频繁读cursor

    @Override
    public long next(int n) {
        long nextValue = this.nextValue;
        long nextSequence = nextValue + n;
        // wrapPoint = 本次申请的最后一个序号 - RingBuffer大小
        long wrapPoint = nextSequence - bufferSize;
        
        // 利用缓存,减少对volatile的读取
        long cachedGatingSequence = this.cachedValue;
        
        // 检查是否需要环绕:如果上次缓存的可消费序号不够,才真正去读消费者的进度
        if (wrapPoint > cachedGatingSequence || cachedGatingSequence > nextValue) {
            // 原子性地记录生产者当前已发布的最大序号(这个读Volatile是必须的)
            long gatingSequence = Util.getMinSequences(gatingSequences, nextValue);
            
            // 如果申请的slot会覆盖未被消费的数据,必须自旋等待
            if (wrapPoint > gatingSequence) {
                LockSupport.parkNanos(1); // 自旋等待
                // 实际是循环重试,这里简化
            }
            // 更新缓存
            this.cachedValue = gatingSequence;
        }
        
        // 无锁、无CAS,直接赋值,这就是单生产者极速的秘密
        this.nextValue = nextSequence;
        return nextSequence;
    }
}

关键逻辑:

  • 因为没有竞争,nextValue 可以用普通写,而 MultiProducerSequencer 必须用 CAS 循环。
  • cachedValue 缓存消费者的最小序号。大多数情况,wrapPoint <= cachedGatingSequence 这个判断直接通过,完全避免读取消费者端的 volatile 变量,将读屏障开销降到最低。

模块四:MultiProducerSequencer —— 乐观自旋分配

为什么快? 高并发下,它比锁的上下文切换开销小得多。


public final class MultiProducerSequencer extends AbstractSequencer {
    // 用来追踪每个slot的生产者写入状态,初始可用
    private final int[] availableBuffer;
    
    @Override
    public long next(int n) {
        long current, next;
        do {
            current = cursor.get(); // Volatile读
            next = current + n;
            // ... 环绕检查 ...
        } while (!cursor.compareAndSet(current, next)); // CAS循环分配序号
        return next;
    }
}
  • 生产者通过 CAS 竞争 cursor 这个序号序列器,拿到一段独占的连续序号。
  • 这是无锁算法,避免了互斥锁导致的上下文切换和线程挂起。冲突时,线程在用户态自旋,成本远低于内核态调度。

模块五:BatchEventProcessor —— 批处理消费

为什么快? 将多次消费合并为一次,大幅摊薄延迟和CPU缓存交换成本。


public final class BatchEventProcessor<T> implements EventProcessor {
    private final Sequence sequence = new Sequence(-1); // 消费者进度
    private final SequenceBarrier barrier;              // 获取可用序列号
    private final EventHandler<T> eventHandler;         // 用户逻辑

    @Override
    public void run() {
        long nextSequence = sequence.get() + 1;
        while (true) {
            // 阻塞等待直到有可用序号
            final long availableSequence = barrier.waitFor(nextSequence);
            
            // 批处理核心:拿到一批数据就不断处理,直到追上生产者
            while (nextSequence <= availableSequence) {
                T event = dataProvider.get(nextSequence);
                eventHandler.onEvent(event, nextSequence, nextSequence == availableSequence);
                nextSequence++;
            }
            // 批量更新消费者进度,只需一次Volatile写
            sequence.set(availableSequence);
        }
    }
}

关键逻辑:

  • barrier.waitFor(nextSequence) 可能等待第一批数据到来,但一旦拿到 availableSequence,它会在这个 while 循环里把积压的所有事件一次性、不间断地处理掉。
  • 只在批处理结束时才更新自己的 sequence 序号。这最小化了 volatile 写操作,并使得生产者需要探查的消费者最小序号变化频率极低,大大减少了生产者的环绕检查成本。


模块六:SequenceBarrier —— 依赖跟踪

为什么快? 它能够精准确定“哪些数据我已经可以安全消费了”,从而支持复杂的消费者间依赖(菱形、六边形),而这一切都通过无锁的 Sequence 操作完成。

逻辑简述:
ProcessingSequenceBarrier 内部持有所有“前置依赖”的序列引用。waitFor 实现中,它会循环获取这些序列的最小值,并与 cursor(生产进度)比较,确定可消费的上界。这个过程全是基于 volatile 读的无锁协调。


总结

优化维度核心实现模块为何快?
内存RingBuffer预分配、无GC,数据连续,缓存友好。
缓存Sequence填充缓存行,消除伪共享,多核并行度是真线性。
并发Sequencer无锁CAS或单写,无互斥锁,无上下文切换开销。
批处理BatchEventProcessor读/写屏障操作均摊,最大限度地利用CPU缓存。
指令级Sequence中的putOrderedLong用更廉价的StoreStore屏障代替昂贵的StoreLoad屏障。

这些模块环环相扣,将硬件性能压榨到了极致,从而造就了 Disruptor 在高性能并发队列中难以撼动的地位。