前言

Disruptor是英国外汇交易公司LMAX开发的一个高性能队列,研发的初衷是解决内存队列的延迟问题(在性能测试中发现竟然与I/O操作处于同样的数量级)。基于Disruptor开发的系统单线程能支撑每秒600万订单,2010年在QCon演讲后,获得了业界关注。2011年,企业应用软件专家Martin Fowler专门撰写长文介绍。同年它还获得了Oracle官方的Duke大奖。

在高并发系统中,线程之间经常需要传递数据。例如:

  • 业务线程将日志事件交给日志线程;
  • 网络线程将请求交给业务线程;
  • 订单接收线程将订单交给风控、撮合和持久化线程;
  • 数据采集线程将指标交给计算与存储线程。

Java 提供了 BlockingQueueConcurrentLinkedQueue 等并发队列,它们能够很好地满足大多数业务场景。但在金融交易、实时风控、行情处理、异步日志等对延迟和吞吐量极其敏感的系统中,传统并发队列可能面临以下问题:

  1. 锁竞争和线程上下文切换;
  2. 频繁创建队列节点,引发垃圾回收;
  3. 多线程竞争同一个队头或队尾;
  4. 缓存局部性较差;
  5. 多消费者协作关系难以表达。

为了解决这些问题,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 实战

下面实现一个简单的订单处理系统:

  1. 生产者发布订单;
  2. 校验处理器检查订单是否合法;
  3. 持久化处理器模拟记录订单;
  4. 业务处理器在前两个处理器完成后执行;
  5. 清理处理器释放 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 被写满。

对于不可避免的慢操作,可以考虑:

  1. 批量处理;
  2. 设置超时;
  3. 将慢操作放入独立处理阶段;
  4. 使用另一个队列进行速度隔离;
  5. 将非关键支路与核心低延迟链路拆分。

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 已移除了旧版的 WorkerPoolWorkProcessorhandleEventsWithWorkerPool 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 的高性能并不是由某一个“神奇算法”带来的,而是多个设计共同作用的结果:

  1. 使用固定大小的环形数组保存事件;
  2. 启动阶段预分配 Event,减少运行时对象创建;
  3. 使用持续增长的 Sequence 表示生产和消费进度;
  4. 通过最慢消费者序号防止槽位被提前覆盖;
  5. 通过 CAS 和内存屏障减少锁竞争;
  6. 针对伪共享进行缓存行隔离;
  7. 支持连续批量消费;
  8. 支持消费者广播和依赖图;
  9. 提供多种等待策略,在延迟和 CPU 使用率之间进行权衡。

但 Disruptor 并不是所有并发系统的最佳选择。

它最适合的是:

同一进程内
+ 高吞吐
+ 低延迟
+ 有界事件流
+ 可控线程模型
+ 明确消费依赖

如果系统主要瓶颈来自数据库、网络请求或外部服务,那么首先应该优化真正的耗时环节,而不是盲目替换队列。

最终,判断是否采用 Disruptor 的标准不应该是“它是否足够快”,而应该是:

它的事件模型、背压方式、消费者语义和资源消耗,是否真正符合当前系统的需求。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐