Disruptor无锁环形队列详解:从数组结构到伪共享缓存优化

发布时间:2026/10/10 21:41:24
Disruptor无锁环形队列详解:从数组结构到伪共享缓存优化
队列这东西往小了说是一个先进先出的数据结构往大了说就是整个并发系统的“血管”。Disruptor 正是因为把最常见的环形队列做到了极致才能在 LMAX 的交易场景里支撑每秒数百万级的状态更新后来也被很多人用在日志采集、合并转发、流式计算等地方。这篇文章我想从数据结构层面开始逐步拆解 Disruptor 是怎么把“数组 头尾指针”的环形队列做成一套无锁、批次化、能规避缓存伪共享的高性能骨架。如果你写过阻塞队列用过ArrayBlockingQueue、LinkedBlockingQueue并且发现它们在高峰流量下锁竞争严重、延迟抖动明显那这篇文章应该对你有用。先说明一下我下面会先讲教科书里的环形队列再一层层向 Disruptor 的实现靠拢。你会发现知识脉络其实是顺下来的先理解rear和length的意义再理解为什么2 的幂容量 位运算更好用接着理解无锁序号的协同最后理解缓存行填充。把这四条线串起来Disruptor 的核心原理也就拿下了。1. 环形队列基础rear 与 length 的经典表达1.1 教科书里的环形队列定义很多教材里有这么一段话“假设以数组 q[m] 存放循环队列中的元素同时以 rear 和 length 分别指示环形队列中的队尾位置和队列长度。”这是描述环形队列最经典的写法。传统的线性队列用数组实现时出队之后队头指针后移前面的位置就浪费了需要不断搬移元素。环形队列的思路是逻辑上把数组首尾相接让队尾指针绕回数组开头继续使用。在这个定义里rear表示下一个元素要写入的位置也就是队尾指针length表示当前队列里实际有多少个元素数组总共能容纳m个元素。于是判断队列状态就变得非常简洁队空length 0队满length m入队q[rear] x; rear (rear 1) % m; length出队front (rear - length m) % m; x q[front]; length--为什么出队要这么算front因为队列里现存length个元素它们是从队头到队尾连续摆放的。用rear减去length再按数组长度取模就能倒推出最前面的那个元素位置。这在只给rear和length、不给独立front指针的场景下非常实用。我拿旋转寿司来打个比方。一条环形传送带相当于数组厨师往rear位置放新盘子客人从front位置取走盘子。传送带上正在流转的盘子总数就是length。只要知道下一个空位在哪里以及台面上有多少盘子头尾都能推算出来。用length而不是单独加标志位的好处是不会浪费一个数组槽位。经典循环队列为了避免“队空”和“队满”时front rear造成的歧义往往会强制空出一个位置实际可用容量成了m - 1。引入length之后所有槽位都能装数据这是最简单直接的空间换逻辑。1.2 环形下标与模运算的数学基础环形队列的核心操作是下标回绕rear (rear 1) % m;取模运算保证了下标始终落在[0, m)之间。模运算本身代价不低但计算机底层对“2 的幂”有天然的位运算加速所以当我们把容量设计为m 2^n时(pos 1) % m可以改写为(pos 1) (m - 1)。这是个等价的数学结论因为m是 2 的幂m - 1的低 n 位全是 1高位全是 0按位与的结果就是保留低 n 位恰好等于 mod 结果。Disruptor 正是这么做的。它要求RingBuffer的容量必须是 2 的幂如果不满足就直接抛异常。这不是什么洁癖而是两个理由推着它必须这么做位运算比取模快得多在每秒几百万次入队出队的场景下省下的 CPU 周期能直接转化为吞吐量。数组槽位分配更规整缓存行亲和性更好配合无锁序号设计时也更容易推算“环绕点”。不要小看这个细节。很多人在自己实现环形队列时随手写一个capacity然后到处用%功能上没问题但性能上已经开始落后了。高性能队列的第一课就是先把取模优化掉。1.3 Disruptor 的 RingBuffer 雏形理解了教科书环形队列再看 Disruptor 的内部结构就顺畅多了。RingBuffer内部本质上就是一个Object[]数组固定容量下标回绕方式就是index sequence (bufferSize - 1)。但它比教科书版本多了一层抽象不再直接操作rear和length而是引入了一个“全局递增序号”的概念。在 Disruptor 里每个事件在数组里都有一个唯一的序号sequence。生产者每发布一个事件公开的“游标”cursor就前进一格消费者要读某个事件也是拿着序号去找对应的槽位。因为数组容量固定序号增长到超过容量时自然环绕复用旧槽位所以环形结构依然在只是状态表达从“头尾指针”变成了“生产者游标 消费者游标”。这个转变非常关键。它把“队列位置”和“事件身份”解耦了。教科书里rear front时可能会想“空还是满”而 Disruptor 里完全不存在这种歧义生产者发布了多少序号消费者消费到多少序号一清二楚。剩下的问题只有一个当生产者试图覆盖还没被消费的槽位时怎么协调。这就是 Disruptor 高性能的核心战场。2. 无锁化改造Disruptor 如何把竞争降到最低2.1 阻塞队列到底慢在哪里先回到传统的ArrayBlockingQueue和LinkedBlockingQueue。它们都是线程安全的阻塞队列正确性没问题但高并发下性能上不去主要瓶颈在于锁。Java 里的锁在无竞争时开销很小但一旦发生竞争线程就要从运行态切换到阻塞态再切换到就绪态这个上下文切换成本往往在上百微秒量级。更麻烦的是锁竞争还会带来“惊群效应”多个生产者同时等一把锁锁释放后只能有一个线程抢到其他线程继续阻塞。吞吐量上不去尾部延迟也会很难看。Disruptor 的选择是不用锁用无锁 CAS 和序号协同。CAS 是 CPU 提供的原子原语失败就重试线程不会挂起。没有线程切换自然没有那些切换开销。2.2 Sequence 与序号协同逻辑Disruptor 里的Sequence是核心类它内部保存一个long value可以理解成一个加强版的AtomicLong。但为什么不用现成的AtomicLong有两个原因单生产者场景下生产者的cursor只被一个线程写不需要原子 CAS只需要保证“写入事件数据之后再更新游标”的存储顺序用UNSAFE.putOrderedLong这种 store-store 屏障就够了。Sequence自带缓存行填充这对性能影响很大下一节详说。序号的协同逻辑大概是这样的场景单生产者写入时next cursor.get() 1 if 环绕位置next - bufferSize大于所有消费者最小编号 说明消费者还没读走旧数据自旋等待或退避 否则 cursor 前进到 next生产者向对应槽位写入事件最后 publish多生产者写入时更复杂一点多个生产者会通过 CAS 竞争获得一段连续序号区间。比如生产者 A 拿到 100~109生产者 B 拿到 110~119然后它们各自写各自的槽位互不冲突。这里 CAS 只发生在“序号分配”这一瞬间而不是每一次数据写入都加锁。消费者侧也有类似的Sequence。每个消费者维护自己的消费进度生产者发布新事件后消费者根据自己的进度可以从上一条已消费序号 1 一直读到当前游标位置批量处理。无锁不是说“没有协调”而是把协调变成了“对序号的读取和比较”。这比内核锁轻量太多。2.3 批量申请与批量消费一次同步做多件事教科书环形队列的入队出队一次只处理一个元素。Disruptor 把这一条也改了支持批处理。生产者可以一次申请多个序号比如调用next(n)一口气拿到 n 个连续的槽位全部写入之后再统一发布。这样做的好处是序号分配的原子操作次数从 n 次降到 1 次多线程竞争窗口大幅收窄吞吐自然就上去了。消费者也一样。一个消费者在唤醒后不一定要一个一个地处理而是先看一下当前游标cursor已经发布到哪个序号了。如果有 0 到 100 的事件可读它可能一次性把这 100 个都取出来处理最后只更新一次自己的消费序号。尤其适合批处理业务比如攒一批日志再往外写。批量处理还带来一个隐藏收益减少缓存行同步频率。缓存一致性协议在多个核之间同步数据是有代价的一次同步能多干几件事平均到每一条消息上的开销就小得多。3. 缓存行与伪共享环形队列不能忽略的性能杀手3.1 什么是伪共享现代 CPU 的缓存是以“缓存行”为单位加载和同步的x86 架构下一般是 64 字节。也就是说不管程序只改一个 4 字节的 intCPU 实际上会把它所在的那一整块 64 字节内存一起缓存起来。设想一个场景两个不同的变量被安排在同一个缓存行里线程 A 频繁写变量 1线程 B 频繁写变量 2。A 写变量 1 后为了保证缓存一致性B 核上对应的缓存行要失效B 写变量 2 后A 核上的缓存行又要失效。两边互相拖累明明操作的是不同的数据却产生了类似“共享冲突”的代价这就叫伪共享。在高性能队列里生产者的游标和消费者的进度很可能会落在相邻内存位置。如果不做隔离生产者推进游标会让消费者的缓存行持续失效消费者更新进度又反过来让生产者缓存行失效。整个队列就会退化成一场“缓存行互踢”的灾难。3.2 Disruptor 的缓存行填充写法Disruptor 的解法非常直白让每个需要频繁读写的Sequence对象独占一个缓存行。一个long是 8 字节要凑满 64 字节可以在字段前后各填充若干个long。Disruptor 源码里早期版本的写法是abstract class LhsPadding { protected long p1, p2, p3, p4, p5, p6, p7; } class Value extends LhsPadding { protected volatile long value; } abstract class RhsPadding extends Value { protected long p9, p10, p11, p12, p13, p14, p15; }这样value前后各有 56 字节加上自身 8 字节整个对象核心区域占满 64 字节。当不同线程各自持有一个这样的Sequence时它们不会落在同一缓存行伪共享被隔离掉了。JDK 后来提供了Contended注解也能做类似的事但需要配合-XX:-RestrictContended参数并且在不同版本上表现不太一样。Disruptor 选择手工填充主要是为了兼容性和可控性。3.3 伪共享优化的边界与误用伪共享优化不是万灵药也不是什么地方都要填。如果两个变量确实经常被同一个线程访问那它们放在同一缓存行反而是好事因为局部性好。比如事件对象的多个字段经常被消费者一起读取你就没必要给每个字段都加 padding那样只会让内存膨胀、缓存利用率下降。常见的误用是把所有字段都填充一遍最后对象体积从几十字节变成几百字节数组遍历时反而把缓存行缓存能力消耗光了。正确做法是先分析并发访问模型找出“不同线程写、却可能落在同一缓存行”的变量再针对性地做隔离。另外还要注意数组元素本身很难做缓存行填充。Object[]里存的是对象引用数组相邻槽位的引用天然挨在一起这部分伪共享靠 JVM 的字段布局解决不了。Disruptor 的做法是让每个槽位指向一个预分配的事件对象事件对象里的字段布局自己控制但它也不会蠢到给每个事件对象都前后填满几百字节那样内存根本装不下几十万队列。所以填充只用于少数全局序号而不是每个数据元素。4. 实操搭建一个可运行的 Disruptor 风格环形队列4.1 第一步用 rear 和 length 写一个基础版先不用锁不用 Disruptor做一遍最原始的环形队列。这里我选择用synchronized保证线程安全让逻辑先跑通public class SimpleRingQueueT { private final Object[] elements; private int rear; private int length; public SimpleRingQueue(int capacity) { elements new Object[capacity]; } public synchronized boolean offer(T value) { if (length elements.length) { return false; } elements[rear] value; rear (rear 1) % elements.length; length; return true; } SuppressWarnings(unchecked) public synchronized T poll() { if (length 0) { return null; } int front (rear - length elements.length) % elements.length; T value (T) elements[front]; length--; return value; } public synchronized boolean isEmpty() { return length 0; } }这段代码对应前面公式的全部逻辑。poll()里用rear - length elements.length再取模就是为了算出队头位置。如果你是第一次写环形队列建议先把这段代码跑通再进入下一步。功能正确永远比性能先行。4.2 第二步把锁换成无锁序号基础版能用但synchronized锁在高并发下会拖垮吞吐。现在我们模仿 Disruptor 的思路把队列入队出队改成“序号分配 CAS 协同”。先定义一个简单的Sequencepublic class Sequence { private final AtomicLong value new AtomicLong(); public long get() { return value.get(); } public boolean compareAndSet(long expected, long newValue) { return value.compareAndSet(expected, newValue); } public void set(long v) { value.set(v); } }单生产者的核心逻辑可以简化成public class SingleProducerSequencer { private final int bufferSize; private final Sequence cursor new Sequence(); private final Sequence[] gatingSequences; public long next() { while (true) { long current cursor.get(); long next current 1; long wrapPoint next - bufferSize; long minSequence getMinimumGatingSequence(); if (wrapPoint minSequence) { // 消费者还没跟上不能覆盖 Thread.yield(); continue; } if (cursor.compareAndSet(current, next)) { return next; } } } public void publish(long sequence) { // 在单生产者场景下publish 需要保证 // 先把事件数据写入数组再把 cursor 游标推进到对消费者可见 cursor.set(sequence); } }这里compareAndSet的作用是防止多线程并发拿到同一个序号。拿到序号后生产者往对应槽位写数据publish时通过内存屏障让消费者看到新数据。实际 Disruptor 的 publish 还会使用UNSAFE.putOrderedLong这类操作保证“先写事件、后发布游标”的顺序不被 CPU 重排。这相当于给这个环节加了一道 store-store 屏障。消费者读取时则使用带 acquire 语义的读取防止读到过期的游标值却拿到尚未写入完成的事件数据。4.3 第三步接入真实 Disruptor API自己实现无锁序列能做出来但工程上直接用 LMAX 的 Disruptor 更省事。下面是一个最小可运行示例事件用OrderEvent表示public class OrderEvent { private long orderId; private double price; public long getOrderId() { return orderId; } public void setOrderId(long orderId) { this.orderId orderId; } public double getPrice() { return price; } public void setPrice(double price) { this.price price; } }消费者处理器public class OrderEventHandler implements EventHandlerOrderEvent { Override public void onEvent(OrderEvent event, long sequence, boolean endOfBatch) { // 在这里处理订单事件比如写入数据库、聚合统计、转发下游 System.out.println(处理订单 event.getOrderId()); } }装配方式int bufferSize 1024; DisruptorOrderEvent disruptor new Disruptor( OrderEvent::new, bufferSize, Thread::new, ProducerType.SINGLE, new YieldingWaitStrategy() ); EventHandlerOrderEvent handler new OrderEventHandler(); disruptor.handleEventsWith(handler); disruptor.start(); RingBufferOrderEvent ringBuffer disruptor.getRingBuffer(); long sequence ringBuffer.next(); try { OrderEvent event ringBuffer.get(sequence); event.setOrderId(1L); event.setPrice(123.45); } finally { ringBuffer.publish(sequence); }这段代码背后new Disruptor(OrderEvent::new, ...)会在启动时预创建bufferSize个OrderEvent对象而不是每次发布都 new 一个。事件对象被重复复用字段被覆写于是对象创建和 GC 压力被压到极低。这也是 Disruptor 能保持低延迟的重要原因之一。4.4 等待策略低延迟和低 CPU 的取舍消费者在等待新事件时不能干瞪眼需要有一个等待策略。Disruptor 提供了几种等待策略行为适用场景BlockingWaitStrategy内部用锁和条件变量消费者阻塞最省 CPU对延迟不敏感、并发量不大的场景SleepingWaitStrategy循环检查 Thread.sleep用延迟换 CPU异步日志、吞吐优先但不追求极端延迟YieldingWaitStrategy循环检查 Thread.yield延迟低较耗 CPU交易系统、实时计算CPU 资源充足BusySpinWaitStrategy纯自旋延迟最低CPU 核一直被占用专用低延迟机器绝不能用在 CPU 紧张的环境选等待策略没有银弹。我见过有人把所有等待策略挨个测一遍最后发现Yielding在 4 核以下反而不如Blocking因为自旋把核心占满了业务线程没地方跑。所以一定要结合部署机器的核数和业务负载来判断别只看参数名。5. 常见问题与排查技巧实录5.1 非 2 的幂容量会导致什么后果环形队列依赖(sequence) (bufferSize - 1)定位下标。如果容量不是 2 的幂比如 100那么bufferSize - 1 99按位与之后得到的下标永远落在 [0, 99]看似没问题但这里有个陷阱这个映射不再等价于sequence % 100会破坏数组槽位的均匀覆盖甚至可能出现一个槽位永远轮不到、另一个槽位却被重复写入的错乱。Disruptor 在创建RingBuffer时会校验容量非 2 的幂直接抛出业务异常。排查时第一反应就是看容量。5.2 消费者跟不上生产者拿到“环绕点”后自旋生产速率大于消费速率时生产者的next()会发现最老的槽位还没被消费于是一直自旋或退避。现象是某个 CPU 核占用率很高但业务线程的吞吐上不去。这个不是 bug是背压机制在起作用。不过要注意如果采用BusySpinWaitStrategy自旋会占满整个核导致消费者线程也分不到 CPU最终形成死锁式的假象。排查办法降低生产速率、增加消费者、把单条消费改批处理或者换更省 CPU 的等待策略。5.3 消费者读到的数据是旧值这种情况多半是发布了新事件但消费者读到的还是上一轮的数据。先检查发布顺序错误示范 ringBuffer.get(seq).setValue(...); // 忘记调用 ringBuffer.publish(seq); 正确示范 long seq ringBuffer.next(); try { ringBuffer.get(seq).setValue(...); } finally { ringBuffer.publish(seq); }没有publish生产者游标不会前进消费者自然看不到新事件。另一个常见坑是多个生产者同时写同一个槽位。无锁序号分配虽然能保证两个线程不会拿到同一个新序号但如果你在next()之后、publish()之前大量耗时操作而消费者进度又允许新生产者覆盖旧槽位旧事件就可能被另一个线程改掉一半。所以在多生产者场景下拿到序号后要快速写入并发布不要在中间做耗时操作。5.4 伪共享优化“没效果”怎么排查如果加了缓存行填充吞吐量却没有变化可能的原因有三个填充没有真正生效。JVM 可能对字段重新排列尤其继承结构里的字段布局不一定按顺序摆放。建议用 JOL 等工具打印对象内存布局确认。伪共享原本就不是瓶颈。先做一次火焰图或者 perf 分析看看是不是真的存在缓存一致性竞争不要为了优化而优化。填充过度导致缓存行占用过多。Sequence 对象膨胀后如果数量很多反而挤占缓存空间。这个问题在少数全局 Sequence 上不明显但要心中有数。下面把上面几个问题整理成速查表方便你对照排查现象可能原因排查方向队列容量配置后出现异常容量不是 2 的幂检查 RingBuffer 初始化参数改成 1024、2048 等生产者高 CPU吞吐不涨消费者消费太慢生产端自旋调整等待策略、增加消费者、优化消费逻辑消费者拿到旧数据publish 遗漏或时序不对确认 next/try/finally/publish 三段式写法多生产者数据错乱多个线程写了同一个槽位确认使用多生产者模式并缩短写槽时间加了缓存填充性能不变伪共享不是瓶颈或填充布局被 JVM 改变JOL 验证内存布局再用火焰图分析真实瓶颈我从实际使用中有几个体会。早期我负责过一个订单事件缓冲模块当时用的是ArrayBlockingQueue压测到一定峰值后 CPU 飙升尾部延迟开始抖动。后来换成 Disruptor同样负载下吞吐大约高了两三倍尾部延迟也明显收敛问题确实出在锁竞争和伪共享上。但那次改造也让我明白一个道理无锁队列并不是“银弹”它把问题从“锁等待”转移到了“消费进度与背压控制”上。如果你没理解消费者游标的意义就算换上了 Disruptor数据错乱和丢失的风险反而会更大。另外补一个小技巧在工程里做技术升级时我建议先把业务逻辑用普通阻塞队列完整地跑通一遍再替换成 Disruptor然后对比两项指标——吞吐量和尾部延迟。这样做的好处是如果性能提升明显你很容易判断优化来自无锁化如果性能反而下降你也能快速定位到等待策略或消费者批处理粒度的问题而不是在一片新框架的迷雾里找原因。