深入理解 Disruptor:高性能无锁事件队列的原理与实战
前言
Disruptor是英国外汇交易公司LMAX开发的一个高性能队列,研发的初衷是解决内存队列的延迟问题(在性能测试中发现竟然与I/O操作处于同样的数量级)。基于Disruptor开发的系统单线程能支撑每秒600万订单,2010年在QCon演讲后,获得了业界关注。2011年,企业应用软件专家Martin Fowler专门撰写长文介绍。同年它还获得了Oracle官方的Duke大奖。
在高并发系统中,线程之间经常需要传递数据。例如:
- 业务线程将日志事件交给日志线程;
- 网络线程将请求交给业务线程;
- 订单接收线程将订单交给风控、撮合和持久化线程;
- 数据采集线程将指标交给计算与存储线程。
Java 提供了 BlockingQueue、ConcurrentLinkedQueue 等并发队列,它们能够很好地满足大多数业务场景。但在金融交易、实时风控、行情处理、异步日志等对延迟和吞吐量极其敏感的系统中,传统并发队列可能面临以下问题:
- 锁竞争和线程上下文切换;
- 频繁创建队列节点,引发垃圾回收;
- 多线程竞争同一个队头或队尾;
- 缓存局部性较差;
- 多消费者协作关系难以表达。
为了解决这些问题,LMAX Exchange 开发了 Disruptor。
Disruptor 是一个面向进程内线程通信的高性能事件处理框架。它以环形数组为基础,通过事件预分配、序号协调、缓存友好设计和可配置等待策略,实现了低延迟、高吞吐的线程间数据交换。
需要特别说明的是:
Disruptor 虽然经常被称为“高性能队列”,但它并不是一个传统意义上的 FIFO 队列,而是一套以 RingBuffer 为数据容器、以 Sequence 为协调机制的事件处理框架。
一、Disruptor 是什么
Disruptor 的核心用途,是在同一个 Java 进程内,将事件从生产者线程高效地传递给一个或多个消费者线程。
其典型模型如下:
┌─────────────────┐
│ Producer │
│ 事件生产线程 │
└────────┬────────┘
│ publish
▼
┌─────────────────────────────┐
│ RingBuffer │
│ │
│ [0][1][2][3][4][5]...[N-1] │
└───────┬─────────┬───────────┘
│ │
┌──────▼───┐ ┌──▼──────────┐
│Consumer A│ │ Consumer B │
│ 业务处理 │ │ 日志持久化 │
└──────────┘ └─────────────┘
与普通队列不同,Disruptor 支持:
- 一个事件同时广播给多个消费者;
- 消费者之间建立依赖关系;
- 预先创建并重复使用事件对象;
- 单生产者和多生产者模式;
- 根据延迟与 CPU 使用率选择不同等待策略;
- 通过批量方式连续处理已经到达的事件。
二、Disruptor 与 BlockingQueue 的区别
BlockingQueue 和 Disruptor 都可以完成线程间数据传递,但两者的设计思想并不相同。
| 对比维度 | BlockingQueue | Disruptor |
|---|---|---|
| 数据结构 | 链表或数组队列 | 预分配的环形数组 |
| 消费模型 | 一个元素通常只被一个消费者取走 | 一个事件可以广播给多个消费者 |
| 对象创建 | 入队过程中可能创建新对象或节点 | 事件对象启动时预分配并循环复用 |
| 并发协调 | 锁、条件变量或 CAS | Sequence、内存屏障和 CAS |
| 消费者依赖 | 通常需要业务代码额外协调 | 原生支持消费者依赖图 |
| 等待方式 | 通常阻塞等待 | 阻塞、休眠、让步、自旋等 |
| 主要目标 | 通用性和易用性 | 低延迟与高吞吐 |
| 跨进程通信 | 不支持 | 不支持 |
| 消息持久化 | 不提供 | 不提供 |
普通队列常见的是“竞争消费”:
┌──────────────┐
Producer ──►│ BlockingQueue│
└──────┬───────┘
│
┌────────┴────────┐
▼ ▼
Consumer A Consumer B
同一条消息只会被其中一个消费者取走
Disruptor 默认采用“广播消费”:
┌────────────┐
Producer ─────────►│ RingBuffer │
└─────┬──────┘
│
┌───────────┼───────────┐
▼ ▼ ▼
Consumer A Consumer B Consumer C
每个消费者都会收到同一条事件
因此,Disruptor 更接近一个高性能的进程内事件总线,而不仅仅是一个队列。
三、Disruptor 的核心组件
3.1 Event
Event 是生产者和消费者之间传递的数据载体。
例如,一个订单事件可以定义为:
public class OrderEvent {
private long orderId;
private String symbol;
private long price;
private long quantity;
public long getOrderId() {
return orderId;
}
public void setOrderId(long orderId) {
this.orderId = orderId;
}
public String getSymbol() {
return symbol;
}
public void setSymbol(String symbol) {
this.symbol = symbol;
}
public long getPrice() {
return price;
}
public void setPrice(long price) {
this.price = price;
}
public long getQuantity() {
return quantity;
}
public void setQuantity(long quantity) {
this.quantity = quantity;
}
public void clear() {
orderId = 0L;
symbol = null;
price = 0L;
quantity = 0L;
}
}
需要注意,Disruptor 中的 Event 通常是一个可变对象,并且会被重复使用。
消费者不能长期持有 Event 的引用,否则当 RingBuffer 再次使用这个槽位时,消费者看到的数据可能已经被覆盖。
3.2 RingBuffer
RingBuffer 是一个固定大小的环形数组,用于保存预分配的事件对象。
假设 RingBuffer 大小为 8:
逻辑序号:
0 1 2 3 4 5 6 7 8 9 10 ...
映射槽位:
0 1 2 3 4 5 6 7 0 1 2 ...
序号会不断增长,但底层数组槽位会循环使用。
如果 RingBuffer 容量为 bufferSize,某个序号对应的数组下标可以表示为:
int index = (int) sequence & (bufferSize - 1);
例如,容量为 8:
sequence = 10
index = 10 & 7
index = 2
这也是 RingBuffer 容量必须为 2 的整数次幂的重要原因之一。
当容量为 2 的整数次幂时,可以使用位运算代替取模运算:
sequence % bufferSize
可以转换为:
sequence & (bufferSize - 1)
除了运算效率,这种设计也简化了环形缓冲区内部的索引计算。
3.3 Sequence
Sequence 是 Disruptor 中最重要的概念之一。
生产者和每个消费者都会维护自己的序号,用于表示当前处理进度。
例如:
Producer Sequence = 105
Consumer A Sequence = 103
Consumer B Sequence = 100
RingBuffer Size = 8
这里表示:
- 生产者已经发布到第 105 个事件;
- 消费者 A 已经处理到第 103 个事件;
- 消费者 B 已经处理到第 100 个事件。
生产者不能无限向前发布,因为 RingBuffer 的槽位会被循环使用。如果生产者继续发布到序号 108,就可能覆盖消费者 B 尚未处理的数据。
生产者判断是否可以继续写入时,需要考虑最慢消费者的位置:
wrapPoint = nextSequence - bufferSize
只有当:
wrapPoint <= minimumConsumerSequence
时,生产者才能安全写入。
因此,Disruptor 的背压并不是通过无限扩容实现的,而是通过消费者进度阻止生产者覆盖尚未处理的数据。
3.4 Sequencer
Sequencer 是 Disruptor 真正的并发协调核心。
它负责:
- 为生产者分配可写序号;
- 判断 RingBuffer 是否还有可用容量;
- 防止生产者覆盖未消费事件;
- 发布已经写入完成的序号;
- 管理单生产者或多生产者的并发逻辑。
Disruptor 提供两种生产者模式:
ProducerType.SINGLE
ProducerType.MULTI
单生产者模式
只有一个线程向 RingBuffer 发布事件。
优点:
- 不需要处理多个生产者之间的序号竞争;
- 并发控制更加简单;
- 通常具有更低的延迟;
- 如果业务只有一个生产线程,应优先选择。
多生产者模式
多个线程同时向 RingBuffer 发布事件。
多生产者模式需要额外的 CAS 和可用性标记,以保证不同生产者不会占用相同序号。
因此:
不要因为系统是多线程程序,就直接选择
ProducerType.MULTI。只有确实存在多个发布线程时,才应该使用多生产者模式。
3.5 SequenceBarrier
SequenceBarrier 用于控制消费者可以处理到哪个序号。
消费者不仅需要等待生产者发布事件,还可能需要等待其他消费者处理完成。
例如:
┌───────────────┐
│ RingBuffer │
└───────┬───────┘
│
┌──────────┴──────────┐
▼ ▼
JournalHandler ReplicationHandler
│ │
└──────────┬──────────┘
▼
BusinessHandler
这里:
JournalHandler负责写日志;ReplicationHandler负责数据复制;- 两者可以并行执行;
BusinessHandler必须等待前两个处理器完成后才能执行。
DSL 配置可以写成:
disruptor
.handleEventsWith(journalHandler, replicationHandler)
.then(businessHandler);
底层会为 BusinessHandler 建立一个 SequenceBarrier,使其等待两个上游消费者的序号。
3.6 EventHandler
EventHandler 是业务消费者需要实现的接口。
public class OrderEventHandler implements EventHandler<OrderEvent> {
@Override
public void onEvent(
OrderEvent event,
long sequence,
boolean endOfBatch) {
System.out.printf(
"sequence=%d, orderId=%d, symbol=%s%n",
sequence,
event.getOrderId(),
event.getSymbol()
);
}
}
三个参数分别表示:
event:当前事件对象;sequence:当前事件的逻辑序号;endOfBatch:当前事件是否为本批次最后一个事件。
endOfBatch 对批量写磁盘、批量提交数据库和批量发送网络数据非常有用。
例如:
@Override
public void onEvent(
OrderEvent event,
long sequence,
boolean endOfBatch) {
buffer.add(convert(event));
if (endOfBatch) {
flush();
}
}
相比每处理一条事件就执行一次 I/O,批量提交通常能显著降低系统调用和网络交互开销。
3.7 WaitStrategy
当消费者已经追上生产者时,需要等待新的事件到达。
WaitStrategy 决定消费者如何等待。
不同等待策略在延迟、吞吐量和 CPU 使用率之间做出了不同权衡。
| 等待策略 | 工作方式 | CPU 占用 | 延迟 | 适用场景 |
|---|---|---|---|---|
BlockingWaitStrategy |
锁和条件变量阻塞 | 低 | 较高 | 普通业务系统、CPU 紧张环境 |
SleepingWaitStrategy |
自旋、让步后短暂休眠 | 较低 | 中等 | 异步日志、后台任务 |
YieldingWaitStrategy |
自旋并调用 Thread.yield() |
高 | 较低 | 有充足逻辑核心的低延迟系统 |
BusySpinWaitStrategy |
持续忙等待 | 极高 | 最低 | CPU 独占、极低延迟系统 |
PhasedBackoffWaitStrategy |
分阶段自旋、让步和回退 | 可配置 | 可配置 | 需要平衡延迟与资源使用 |
选择原则可以简化为:
普通业务系统:
BlockingWaitStrategy
异步日志或后台处理:
SleepingWaitStrategy
对延迟敏感且 CPU 资源充足:
YieldingWaitStrategy
线程绑核、CPU 独占、极低延迟:
BusySpinWaitStrategy
BusySpinWaitStrategy 并不是“性能优化开关”。
如果消费者线程数量超过可用物理核心,或者运行环境中还有大量其他线程,忙等待会导致严重的 CPU 竞争,实际性能可能反而下降。
四、Disruptor 为什么快
4.1 事件对象预分配
普通队列中,生产者经常需要不断创建消息对象:
queue.put(new OrderEvent(...));
在高吞吐系统中,这会产生大量短生命周期对象,增加垃圾回收压力。
Disruptor 在启动时通过 EventFactory 创建 RingBuffer 中的全部 Event:
EventFactory<OrderEvent> factory = OrderEvent::new;
如果 RingBuffer 大小为 1024,就会在初始化阶段创建 1024 个 OrderEvent。
运行期间,生产者不是创建新事件,而是取得某个预分配槽位并修改字段:
long sequence = ringBuffer.next();
try {
OrderEvent event = ringBuffer.get(sequence);
event.setOrderId(orderId);
} finally {
ringBuffer.publish(sequence);
}
事件被消费者处理后,槽位还会在未来被继续复用。
需要注意:
Event 预分配不代表业务代码一定是零分配。
如果每次发布时仍然创建新的 String、集合、包装类型或临时对象,系统依然会产生垃圾。要进一步减少分配,应优先使用基本类型、固定长度结构或可复用对象。
4.2 连续数组带来更好的缓存局部性
链表队列中的节点可能分散在堆内存的不同位置:
Node A ──► Node B ──► Node C ──► Node D
CPU 访问下一个节点时,可能需要重新加载缓存行。
RingBuffer 使用固定数组保存对象引用:
[Event0][Event1][Event2][Event3][Event4]
连续访问数组通常具有更好的空间局部性,也更容易利用 CPU 缓存和硬件预取机制。
4.3 减少锁竞争
Disruptor 的生产者和消费者主要通过 Sequence、CAS 和内存屏障进行协调。
在使用非阻塞等待策略时,核心事件传递链路可以避免传统互斥锁。
不过,不能简单地说 Disruptor 在所有配置下都完全无锁。
例如,BlockingWaitStrategy 内部需要使用锁和条件变量让消费者线程进入阻塞状态。
更加准确的说法是:
Disruptor 大量使用非阻塞并发算法,并允许根据等待策略构建低锁或无锁的事件处理链路。
4.4 避免伪共享
现代 CPU 缓存通常以缓存行为单位读写数据。
假设两个线程分别频繁修改两个变量:
class Data {
volatile long producerSequence;
volatile long consumerSequence;
}
即使两个变量彼此独立,只要它们处于同一个缓存行中,两个 CPU 核心仍可能不断使对方的缓存失效。
这种现象称为 伪共享。
同一缓存行:
┌─────────────────────────────────────────┐
│ producerSequence │ consumerSequence │...│
└─────────────────────────────────────────┘
CPU 1 修改 CPU 2 修改
Disruptor 的 Sequence 在实现上针对伪共享进行了隔离处理,使不同线程高频修改的序号尽量避免落入同一个缓存行。
4.5 批量消费
消费者等待序号 100 时,生产者可能已经发布到 110。
消费者被唤醒后,不需要每处理一条事件就重新等待,而是可以连续处理:
100、101、102、103……110
这能够减少:
- 等待策略调用次数;
- 内存屏障开销;
- 线程调度开销;
- I/O 提交次数。
因此,Disruptor 的优势不仅来自“环形数组”,还来自完整的批量事件处理模型。
五、事件发布流程
Disruptor 的一次完整发布过程可以分为三步:
申请序号 → 填充事件 → 发布序号
代码如下:
long sequence = ringBuffer.next();
try {
OrderEvent event = ringBuffer.get(sequence);
event.setOrderId(orderId);
event.setSymbol(symbol);
event.setPrice(price);
event.setQuantity(quantity);
} finally {
ringBuffer.publish(sequence);
}
为什么 publish 必须放在 finally 中
调用 next() 后,生产者已经申请了一个序号。
如果填充事件过程中出现异常,却没有发布这个序号,后续消费者可能一直等待该位置,导致整个处理链停滞。
因此,应始终使用:
try {
// 填充事件
} finally {
ringBuffer.publish(sequence);
}
六、完整 Java 实战
下面实现一个简单的订单处理系统:
- 生产者发布订单;
- 校验处理器检查订单是否合法;
- 持久化处理器模拟记录订单;
- 业务处理器在前两个处理器完成后执行;
- 清理处理器释放 Event 中的对象引用。
事件处理拓扑如下:
┌──────────────────┐
│ RingBuffer │
└────────┬─────────┘
│
┌────────────┴────────────┐
▼ ▼
ValidateEventHandler JournalEventHandler
│ │
└────────────┬────────────┘
▼
BusinessEventHandler
│
▼
ClearEventHandler
6.1 Maven 依赖
<dependency>
<groupId>com.lmax</groupId>
<artifactId>disruptor</artifactId>
<version>4.0.0</version>
</dependency>
Disruptor 4.0.0 要求 Java 11 或更高版本。
6.2 定义事件
package com.example.disruptor;
public class OrderEvent {
private long orderId;
private String symbol;
private long price;
private long quantity;
private boolean valid;
public long getOrderId() {
return orderId;
}
public void setOrderId(long orderId) {
this.orderId = orderId;
}
public String getSymbol() {
return symbol;
}
public void setSymbol(String symbol) {
this.symbol = symbol;
}
public long getPrice() {
return price;
}
public void setPrice(long price) {
this.price = price;
}
public long getQuantity() {
return quantity;
}
public void setQuantity(long quantity) {
this.quantity = quantity;
}
public boolean isValid() {
return valid;
}
public void setValid(boolean valid) {
this.valid = valid;
}
public void clear() {
orderId = 0L;
symbol = null;
price = 0L;
quantity = 0L;
valid = false;
}
@Override
public String toString() {
return "OrderEvent{" +
"orderId=" + orderId +
", symbol='" + symbol + '\'' +
", price=" + price +
", quantity=" + quantity +
", valid=" + valid +
'}';
}
}
6.3 参数校验处理器
package com.example.disruptor;
import com.lmax.disruptor.EventHandler;
public class ValidateEventHandler
implements EventHandler<OrderEvent> {
@Override
public void onEvent(
OrderEvent event,
long sequence,
boolean endOfBatch) {
boolean valid =
event.getOrderId() > 0
&& event.getSymbol() != null
&& !event.getSymbol().isBlank()
&& event.getPrice() > 0
&& event.getQuantity() > 0;
event.setValid(valid);
if (!valid) {
System.err.printf(
"[validate] invalid event, sequence=%d, event=%s%n",
sequence,
event
);
}
}
}
6.4 日志持久化处理器
package com.example.disruptor;
import com.lmax.disruptor.EventHandler;
public class JournalEventHandler
implements EventHandler<OrderEvent> {
@Override
public void onEvent(
OrderEvent event,
long sequence,
boolean endOfBatch) {
// 实际项目中可以写入 WAL、文件或数据库。
System.out.printf(
"[journal] sequence=%d, orderId=%d, endOfBatch=%s%n",
sequence,
event.getOrderId(),
endOfBatch
);
}
}
6.5 业务处理器
package com.example.disruptor;
import com.lmax.disruptor.EventHandler;
public class BusinessEventHandler
implements EventHandler<OrderEvent> {
@Override
public void onEvent(
OrderEvent event,
long sequence,
boolean endOfBatch) {
if (!event.isValid()) {
return;
}
System.out.printf(
"[business] process order: %d, symbol=%s%n",
event.getOrderId(),
event.getSymbol()
);
}
}
6.6 数据清理处理器
package com.example.disruptor;
import com.lmax.disruptor.EventHandler;
public class ClearEventHandler
implements EventHandler<OrderEvent> {
@Override
public void onEvent(
OrderEvent event,
long sequence,
boolean endOfBatch) {
event.clear();
}
}
清理处理器必须位于整个消费链末尾。
否则,前面的处理器还没有使用完数据,Event 中的字段就可能被提前清空。
对于基本类型字段,清理通常不是必需的;但对于 String、集合、大对象等引用类型字段,应及时设为 null,避免 RingBuffer 长期持有对象引用。
6.7 启动 Disruptor
package com.example.disruptor;
import com.lmax.disruptor.BlockingWaitStrategy;
import com.lmax.disruptor.RingBuffer;
import com.lmax.disruptor.dsl.Disruptor;
import com.lmax.disruptor.dsl.ProducerType;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.atomic.AtomicInteger;
public class DisruptorApplication {
private static final int BUFFER_SIZE = 1024;
public static void main(String[] args) {
AtomicInteger threadNumber = new AtomicInteger();
ThreadFactory threadFactory = runnable -> {
Thread thread = new Thread(
runnable,
"order-disruptor-" + threadNumber.incrementAndGet()
);
thread.setDaemon(false);
return thread;
};
Disruptor<OrderEvent> disruptor = new Disruptor<>(
OrderEvent::new,
BUFFER_SIZE,
threadFactory,
ProducerType.SINGLE,
new BlockingWaitStrategy()
);
ValidateEventHandler validateHandler =
new ValidateEventHandler();
JournalEventHandler journalHandler =
new JournalEventHandler();
BusinessEventHandler businessHandler =
new BusinessEventHandler();
ClearEventHandler clearHandler =
new ClearEventHandler();
/*
* 校验和持久化并行执行。
* 业务处理等待二者全部完成。
* 数据清理等待业务处理完成。
*/
disruptor
.handleEventsWith(
validateHandler,
journalHandler
)
.then(businessHandler)
.then(clearHandler);
RingBuffer<OrderEvent> ringBuffer =
disruptor.start();
try {
for (long orderId = 1; orderId <= 10; orderId++) {
publish(
ringBuffer,
orderId,
"AAPL",
18_500L,
100L
);
}
} finally {
/*
* 调用 shutdown 前必须停止发布新事件。
* shutdown 会等待已经发布的事件处理完成。
*/
disruptor.shutdown();
}
}
private static void publish(
RingBuffer<OrderEvent> ringBuffer,
long orderId,
String symbol,
long price,
long quantity) {
long sequence = ringBuffer.next();
try {
OrderEvent event = ringBuffer.get(sequence);
event.setOrderId(orderId);
event.setSymbol(symbol);
event.setPrice(price);
event.setQuantity(quantity);
} finally {
ringBuffer.publish(sequence);
}
}
}
七、消费者编排模式
7.1 并行消费
多个处理器相互独立:
disruptor.handleEventsWith(
handlerA,
handlerB,
handlerC
);
拓扑结构:
┌──► Handler A
RingBuffer ────┼──► Handler B
└──► Handler C
三个处理器都会收到每一条事件,并且可以并行运行。
7.2 串行消费
后一个处理器依赖前一个处理器:
disruptor
.handleEventsWith(handlerA)
.then(handlerB)
.then(handlerC);
拓扑结构:
RingBuffer ──► Handler A ──► Handler B ──► Handler C
适合具有明确处理阶段的流水线,例如:
解析请求 → 参数校验 → 业务计算 → 结果持久化
7.3 菱形依赖
两个处理器并行执行,之后汇聚到一个处理器:
disruptor
.handleEventsWith(handlerA, handlerB)
.then(handlerC);
拓扑结构:
┌──► Handler A ──┐
RingBuffer ───┤ ├──► Handler C
└──► Handler B ──┘
Handler C 只有在 A 和 B 都处理完当前序号后才会继续执行。
7.4 独立支路
一个处理链负责核心业务,另一个消费者独立进行监控:
disruptor
.handleEventsWith(validateHandler)
.then(businessHandler);
disruptor.handleEventsWith(metricsHandler);
拓扑结构:
┌──► Validate ──► Business
RingBuffer ───┤
└──► Metrics
需要注意,多个根消费组也会影响 RingBuffer 的最慢消费进度。任何一个末端消费者过慢,都可能最终阻塞生产者。
八、RingBuffer 容量应该如何设置
RingBuffer 容量必须是 2 的整数次幂,例如:
1024
2048
4096
8192
65536
容量并不是越大越好。
容量过小:
- 突发流量下容易被写满;
- 慢消费者更容易阻塞生产者;
- 系统对短时间延迟抖动的容忍度低。
容量过大:
- 启动时需要预分配更多 Event;
- Event 体积较大时会占用大量内存;
- 数据工作集可能超出 CPU 缓存;
- 消费者严重落后时,积压时间过长。
可以根据流量和允许的最大阻塞时间进行初步估算:
RingBuffer 容量
≈ 峰值每秒事件数 × 可容忍积压秒数
例如:
峰值吞吐量:50,000 条/秒
允许积压:0.1 秒
最低容量约为:
50,000 × 0.1 = 5,000
向上取最近的 2 的整数次幂:
8192
这只是初始值,最终仍然需要结合真实事件大小、消费者速度和生产环境压测结果确定。
九、生产环境调优建议
9.1 优先使用单生产者模式
如果只有一个线程发布事件,应明确配置:
ProducerType.SINGLE
不要无条件使用多生产者模式。
9.2 避免在消费者中执行长时间阻塞操作
以下操作可能让消费者长时间停滞:
- 同步调用远程服务;
- 无超时的数据库请求;
- 文件系统同步刷盘;
- 长时间锁等待;
- 大量日志同步输出。
由于生产者必须考虑最慢消费者进度,一个消费者阻塞可能最终导致整个 RingBuffer 被写满。
对于不可避免的慢操作,可以考虑:
- 批量处理;
- 设置超时;
- 将慢操作放入独立处理阶段;
- 使用另一个队列进行速度隔离;
- 将非关键支路与核心低延迟链路拆分。
9.3 根据部署环境选择等待策略
不要只在开发机上测试 BusySpinWaitStrategy,然后直接用于生产环境。
等待策略的实际表现与以下因素密切相关:
- 物理核心数量;
- 超线程配置;
- 容器 CPU 配额;
- 虚拟化环境;
- 操作系统调度;
- 是否进行线程绑核;
- 同一机器上的其他进程。
在 Kubernetes 容器中,如果 Pod 的 CPU 限额较低,忙等待策略通常会不断消耗时间片,可能造成明显的调度抖动。
9.4 避免捕获型 Lambda
下面的发布代码捕获了外部变量:
ringBuffer.publishEvent((event, sequence) -> {
event.setOrderId(orderId);
});
在部分场景下,捕获型 Lambda 可能引入额外对象分配。
对极致低延迟场景,可以:
- 使用静态、无捕获的
EventTranslator; - 使用直接申请序号的
next/get/publish模式; - 使用 JFR、Async Profiler 或 GC 日志验证是否真的产生分配。
不要仅凭代码形式猜测,应以实际分析结果为准。
9.5 Event 中尽量减少复杂引用
为了降低分配和内存占用,Event 字段可优先使用:
long
int
double
boolean
对于固定协议数据,可以考虑:
- 使用数值 ID 代替字符串;
- 使用枚举编码;
- 使用预分配字节数组;
- 使用 off-heap 存储较大的消息体;
- Event 中仅保存外部数据区域的索引。
不过,这些优化会增加代码复杂度,应根据实际性能目标决定是否采用。
9.6 正确处理异常
消费者异常不能被忽略,否则可能出现:
- 某个事件未完成业务处理;
- 消费者线程停止;
- 序号无法正常推进;
- 生产者最终因 RingBuffer 写满而阻塞。
生产环境中应明确配置异常处理策略,并记录:
- 消费者名称;
- 事件序号;
- 事件关键标识;
- 异常堆栈;
- 是否允许跳过当前事件;
- 是否需要停止整个处理链。
异常处理不能只考虑“系统是否继续运行”,还要考虑继续运行后数据是否仍然正确。
9.7 监控消费者延迟
可以通过生产者游标与消费者 Sequence 的差值计算积压量:
backlog = producerCursor - consumerSequence
建议监控:
- RingBuffer 当前剩余容量;
- 每个消费者的处理序号;
- 最慢消费者积压量;
- 单事件处理延迟;
- 端到端延迟;
- 消费批次大小;
- 生产者等待次数;
- 消费异常次数;
- JVM GC 暂停;
- CPU 使用率和上下文切换次数。
只监控平均延迟是不够的。低延迟系统更应关注:
P99
P99.9
P99.99
最大延迟
9.8 停机前必须先停止生产者
disruptor.shutdown() 会等待已经发布的事件全部处理完成。
如果生产者仍然持续发布,shutdown() 可能一直无法返回。
推荐顺序为:
停止接收新请求
↓
停止事件生产者
↓
等待已经发布的事件消费完成
↓
关闭 Disruptor
↓
释放数据库、网络和文件资源
对于生产系统,还应使用带超时的关闭方式,避免停机流程永久阻塞。
十、常见错误与误区
10.1 把 Disruptor 当成普通线程池任务队列
Disruptor 默认是广播模型。
如果配置了三个 EventHandler:
disruptor.handleEventsWith(
handler1,
handler2,
handler3
);
不是三个消费者竞争处理任务,而是三个消费者都会处理每一条事件。
Disruptor 4.0.0 已移除了旧版的 WorkerPool、WorkProcessor 和 handleEventsWithWorkerPool API。
因此,需要竞争消费、每个任务只执行一次的场景,应重新评估是否直接使用:
ThreadPoolExecutor;ArrayBlockingQueue;- 自定义事件分片;
- 多个独立 RingBuffer;
- 其他任务调度框架。
10.2 发布后继续修改 Event
下面的代码是错误的:
long sequence = ringBuffer.next();
OrderEvent event = ringBuffer.get(sequence);
event.setOrderId(1001L);
ringBuffer.publish(sequence);
// 错误:发布后继续修改
event.setOrderId(2002L);
调用 publish() 后,消费者可能立即读取该 Event。
生产者不能再继续修改它。
10.3 消费者长期保存 Event 引用
下面的做法也存在问题:
private final List<OrderEvent> history = new ArrayList<>();
@Override
public void onEvent(
OrderEvent event,
long sequence,
boolean endOfBatch) {
history.add(event);
}
RingBuffer 会重复使用 Event。列表中保存的多个引用,最终可能指向已经被覆盖的数据。
需要保留数据时,必须复制真正需要的字段:
OrderSnapshot snapshot = new OrderSnapshot(
event.getOrderId(),
event.getSymbol(),
event.getPrice(),
event.getQuantity()
);
复制会产生额外开销,因此应明确数据生命周期和性能目标。
10.4 认为容量足够大就不会阻塞
Disruptor 是有界缓冲区。
即使容量很大,只要消费者长期低于生产者速度,RingBuffer 最终仍会写满。
设生产速度为:
100,000 条/秒
消费速度为:
80,000 条/秒
每秒积压:
20,000 条
容量为 1,000,000 时,也只能缓冲约:
1,000,000 ÷ 20,000 = 50 秒
增大容量只能吸收短期流量波动,无法解决长期生产消费速率不匹配。
10.5 认为使用 Disruptor 就一定更快
Disruptor 的优势通常出现在:
- 事件量大;
- 线程间传递开销占比较高;
- 延迟要求严格;
- 消费处理相对轻量;
- 可以预先分配对象;
- 消费者依赖关系明确;
- CPU 和线程部署可控。
如果业务处理本身需要几十毫秒,例如远程接口调用、复杂 SQL 或大文件读写,那么队列层节省的微秒级开销可能不会显著改善端到端性能。
使用前应先定位真实瓶颈。
十一、如何正确进行性能测试
不要使用下面这种方式得出结论:
long start = System.currentTimeMillis();
// 发送若干事件
long cost = System.currentTimeMillis() - start;
这种测试容易受到以下因素干扰:
- JVM 类加载;
- JIT 编译;
- 逃逸分析;
- 死代码消除;
- GC;
- 操作系统调度;
- CPU 动态频率;
- 测试线程与消费者线程竞争;
- 日志输出;
- 启动和关闭时间。
推荐使用 JMH,并区分以下指标:
吞吐量
每秒能够处理多少条事件
平均延迟
事件从发布到消费完成的平均时间
尾延迟
P99、P99.9、P99.99 延迟
稳态性能
系统长时间运行后的性能和 GC 情况
过载行为
生产速度高于消费速度时,
系统如何背压、延迟如何增长、是否丢失数据
同时,测试中不应在消费者热路径执行:
System.out.println(...)
控制台输出的开销往往远高于队列操作本身,会完全掩盖 Disruptor 的性能差异。
十二、Disruptor 适合哪些场景
Disruptor 比较适合:
1. 金融交易系统
订单接收 → 风控校验 → 日志记录 → 撮合处理
这类系统通常具有:
- 延迟敏感;
- 事件模型固定;
- 处理顺序明确;
- 线程部署可控;
- 对尾延迟要求高。
2. 实时行情处理
行情接收 → 数据解析 → 指标计算 → 策略执行
RingBuffer 可以作为不同处理阶段之间的高速事件通道。
3. 异步日志
业务线程仅负责写入日志事件,消费者线程负责格式化和落盘。
此时可以选择对生产者影响较小的等待策略,并通过 endOfBatch 批量写文件。
4. 实时风控
同一笔交易可以同时广播给:
- 账户风控;
- 额度检查;
- 反欺诈规则;
- 审计记录。
多个检查器可以并行执行,之后再汇聚到最终决策处理器。
5. 游戏服务器
可以用于:
- 玩家指令处理;
- 状态更新;
- 战斗结算;
- 行为日志;
- 消息广播。
6. 高性能指标采集
多个业务线程发布指标,后台消费者进行聚合、计算和批量写出。
十三、哪些场景不适合使用 Disruptor
1. 跨进程或跨服务通信
Disruptor 只负责同一进程内的线程通信。
跨服务场景应考虑:
- Kafka;
- RocketMQ;
- RabbitMQ;
- Pulsar;
- gRPC;
- HTTP。
2. 需要消息持久化和重放
Disruptor 本身不提供:
- 消息持久化;
- 消费确认;
- 消费位点管理;
- 宕机恢复;
- 消息重放;
- 分布式副本。
如果需要这些能力,应使用消息中间件,或者在 Disruptor 消费链中自行实现日志与恢复机制。
3. 消费者包含大量慢速 I/O
如果消费者主要执行远程调用,吞吐量由外部系统决定,Disruptor 的低延迟优势可能无法体现。
4. 简单后台任务
对于普通异步任务:
ExecutorService executor =
Executors.newFixedThreadPool(8);
通常已经足够。
不要为了追求技术复杂度,在没有性能需求的系统中引入 Disruptor。
5. 需要天然任务竞争消费
Disruptor 4.x 的主要模型是事件广播与依赖编排。
如果需求只是“多个工作线程竞争处理任务,每个任务只执行一次”,有界阻塞队列和线程池往往更加直接。
十四、Disruptor、线程池与消息队列如何选择
可以按照下面的方式判断:
是否跨进程?
├── 是
│ └── 使用 Kafka、RocketMQ、RabbitMQ 等消息中间件
│
└── 否
│
├── 是否需要极低延迟和高吞吐?
│ ├── 否
│ │ └── BlockingQueue 或线程池
│ │
│ └── 是
│ │
│ ├── 是否需要广播消费或依赖编排?
│ │ ├── 是:优先评估 Disruptor
│ │ └── 否:Disruptor 或专用无锁队列
│ │
│ └── 是否能控制线程与 CPU 资源?
│ ├── 是:可使用激进等待策略
│ └── 否:优先使用阻塞或休眠策略
简化对照如下:
| 需求 | 推荐方案 |
|---|---|
| 普通异步任务 | 线程池 |
| 任务竞争消费 | BlockingQueue + 线程池 |
| 进程内低延迟事件流水线 | Disruptor |
| 一条事件广播给多个消费者 | Disruptor |
| 跨服务可靠消息 | Kafka、RocketMQ、RabbitMQ |
| 消息持久化与重放 | 分布式消息队列 |
| 网络连接事件处理 | Netty |
| 定时任务调度 | Quartz、XXL-JOB 等 |
十五、总结
Disruptor 的高性能并不是由某一个“神奇算法”带来的,而是多个设计共同作用的结果:
- 使用固定大小的环形数组保存事件;
- 启动阶段预分配 Event,减少运行时对象创建;
- 使用持续增长的 Sequence 表示生产和消费进度;
- 通过最慢消费者序号防止槽位被提前覆盖;
- 通过 CAS 和内存屏障减少锁竞争;
- 针对伪共享进行缓存行隔离;
- 支持连续批量消费;
- 支持消费者广播和依赖图;
- 提供多种等待策略,在延迟和 CPU 使用率之间进行权衡。
但 Disruptor 并不是所有并发系统的最佳选择。
它最适合的是:
同一进程内
+ 高吞吐
+ 低延迟
+ 有界事件流
+ 可控线程模型
+ 明确消费依赖
如果系统主要瓶颈来自数据库、网络请求或外部服务,那么首先应该优化真正的耗时环节,而不是盲目替换队列。
最终,判断是否采用 Disruptor 的标准不应该是“它是否足够快”,而应该是:
它的事件模型、背压方式、消费者语义和资源消耗,是否真正符合当前系统的需求。
更多推荐




所有评论(0)