核心架构
Disruptor 内部可视为这几个模块的精密协作:
- RingBuffer (环形缓冲区):数据容器,预分配。
- Sequencer (序列器):协调生产者的并发安全与序号分配。
- Sequence (序列):原子性的进度指示器,解决伪共享。
- SequenceBarrier (序列屏障):消费者协调器,确定可消费范围。
- EventProcessor (事件处理器):消费者线程封装,实现批处理。
- 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 在高性能并发队列中难以撼动的地位。

