如果你写过千万级流量的后端服务,或者只是被线上 GC 问题折磨过,你大概率听说过 Disruptor 这个名字。很多人把它当成“最快的队列”,但真正上手之后才发现,想让它跑到单机 200 万 TPS,其实没那么玄乎,但也没那么简单。这篇内容不打算复述官方文档,我会从 RingBuffer 的设计原理讲到实战代码,再把调优过程中踩过的坑和关于“tps 虚高”的冷思考一并聊透,适合正在选型消息中间件、冲刺高吞吐场景的开发者参考。
1. 内容整体设计与思路拆解
1.1 为什么主流队列在高并发下不够用
先想想我们平时用的队列问题出在哪。以 JDK 自带的 LinkedBlockingQueue 为例,它的底层是链表加两把锁,take 和 put 各自锁住一段操作。链表的每个节点都是运行时 new 出来的,节点之间的内存地址并不连续,CPU 缓存线很难命中,而每次入队出队只要遇到锁竞争,线程就会从用户态陷入内核态,这个开销在高吞吐场景下非常致命。
ArrayBlockingQueue 稍好一点,底层是数组,内存连续,但它的读写共用同一把 ReentrantLock,读写互相阻塞。理论上它是有界队列,可以做缓冲,但实际上只要生产者速度上来,消费者就得在锁上等待,吞吐量很难突破每秒几十万这个量级。
如果去搜“高性能队列”,你还会看到 ConcurrentLinkedQueue,它确实是无锁的,用 CAS 实现入队出队,但它的无锁只是用了轻量的原子操作,链表节点仍要频繁创建和回收,而且它没有“槽位”的概念,消费者不知道下一个数据什么时候到,只能靠忙等或者懒加载触发,延迟和吞吐都不可控。
上面这些队列都服务于“方便好用”这个目标,而 Disruptor 的设计目标非常纯粹:把单线程下的顺序访问做到极致,再用无锁手段处理跨线程的发布与消费。这个思路来自 LMAX 公司对金融交易系统的极致追求,后来开源后大量用于日志处理、交易撮合、分布式存储等需要极低延迟和高吞吐的领域。
1.2 Disruptor 的整体架构:事件、序号器与依赖图
理解 Disruptor 一定不能只盯住 RingBuffer 一个类,它实际上是一整套基于事件的并发框架。核心组件包括:Event(真正要传递的数据对象)、EventFactory(预创建事件对象的工厂)、RingBuffer(存储事件的环形数组和序号分配器)、SequenceBarrier(消费者进度对齐的屏障)、WaitStrategy(消费者如何等待新事件)以及 EventHandler 等处理器。
值得一提的是,Disruptor 把生产者和消费者之间的协作变成了“序号”之间的对比。生产者申请序号、写入事件、发布序号,消费者等待序号、处理事件、更新自己的已消费序号。整个交互过程中没有锁,只有对 Sequence 对象中一个 Volatile 字段的原子操作,以及各种内存屏障保证顺序与可见性。
从设计上看,Disruptor 有三大特点:一是环形数组避免内存频繁分配,二是缓存行填充避免伪共享,三是无锁设计消除上下文切换。这三点互为支撑,缺一个都会让性能大幅跳水。接下来的部分会逐个拆开讲。
2. 核心细节解析与实操要点
2.1 环形队列(RingBuffer)的底层结构
RingBuffer 本质上就是一个固定大小的 Object 数组,数组长度通常被设置为 2 的整数次幂,比如 4096、65536。为什么要用 2 的幂次?因为下标计算可以用位运算sequence & (size - 1)代替取模,这个操作比%快一个数量级,在每秒百万级发布频率下能省下明显的 CPU 开销。
数组在初始化阶段就通过 EventFactory 把每个槽位的对象创建好,后续生产者写入事件时不是new一个新对象丢进去,而是直接从槽位里取出已有的对象,填充它的字段,再“发布”。这样做的好处是避免了高并发下频繁 GC 带来的停顿。你可以把 RingBuffer 想象成一个预先摆好一万个餐盘的旋转寿司台,厨师只把菜放到盘子上,食客拿走后再把空盘子放回去,而不是每次做菜都要新买一个盘子。
RingBuffer 内部用两个关键序号来维持读写平衡:一个是指示下一个可写槽位的序列号,由生产者持有;另一个是消费者组已经消费到的最小序列号,由各个消费者通过 Sequence 对象上报。生产者在写之前必须看消费者进度,如果消费者太慢,生产者就得停下来等待,或者在多生产者场景下通过 CAS 竞争申请可用序号。
2.2 缓存行填充与伪共享的真相
现代 CPU 读取数据不是按字节来的,而是以 64 字节的缓存行为单位加载到 L1/L2/L3 缓存中。如果两个核心分别频繁修改两个变量,而这两个变量刚好落在同一条缓存行里,那么每次修改都会导致整条缓存行在其他核心上失效,必须回写内存再重新加载。这在多线程程序里是非常隐蔽的性能杀手。
JDK 里的 LongAdder 就遇到过这种问题,后来靠@Contended注解加上 JVM 参数-XX:-RestrictContended才解决。Disruptor 的 Sequence 类更直接,它在核心的 volatile 值前后各填充了 7 个 long 类型的占位字段,保证了这个值单独占用一个缓存行,不会和相邻对象的数据互相干扰。
我第一次接触 Disruptor 源码时,看到 Padding 那一串无用的 p1、p2、p3,觉得这完全是多此一举。后来在压测中验证了一次,同样的代码,将 Sequence 和下游消费者的 index 字段放在一起,吞吐量下降了约 35%。伪共享不是理论问题,而是实实在在会拖垮性能的东西。
在实战中要注意,即使 Disruptor 自身做了缓存行填充,你的事件对象如果包含多个被不同线程读写的字段,依然可能产生伪共享。解决办法通常是给事件对象的头尾各加 7 个 long 占位,或者把发布频率极高的字段单独隔离到独立的类里。
2.3 生产、发布与消费的完整周期
正常使用 Disruptor 时,生产者的发布流程分为“申请”和“提交”两个阶段。你首先调用ringBuffer.next()获取一个可用的序号,这个操作在多生产者模式下通过 CAS 的compareAndSet实现,在单生产者模式下则是简单的递增。拿到序号后,通过ringBuffer.get(sequence)获得槽位对象,往里面填数据,最后调用ringBuffer.publish(sequence)让该序号对消费者可见。
为什么要分两步走?如果不发布,消费者是感知不到这个数据的。publish操作内部会调用内存屏障,保证先前对事件对象的写入对其他线程可见。这一步做不好,消费者看到的可能是不完整的脏数据。
消费者端则通过 SequenceBarrier 等待目标序号的出现,一旦发现生产者已经发布了比自己当前进度更靠后的序号,就取出事件进行处理。处理完成后更新自己的 Sequence,表示“我已经消费到这里”。多个消费者通过共同的 SequenceBarrier 协调进度,可以做到每个事件只被处理一次,也可以做到每个消费者都拿到完整事件流,取决于你配置的消费模式。
3. 实操过程与核心环节实现
3.1 一个能跑的最小 Disruptor 示例
直接从工程代码讲起。假设我们要实现一个每秒处理 100 万个订单事件的系统,订单事件包含订单号、金额和状态。第一步是定义事件类:
public class OrderEvent { private long orderId; private double amount; private String status; public void clear() { orderId = 0L; amount = 0D; status = null; } // getter / setter 省略 }接着创建 EventFactory,这一层只负责在 RingBuffer 初始化时生成对象实例,不参与业务逻辑:
public class OrderEventFactory implements EventFactory<OrderEvent> { @Override public OrderEvent newInstance() { return new OrderEvent(); } }然后是处理订单的消费者,实现 EventHandler 接口:
public class OrderEventHandler implements EventHandler<OrderEvent> { @Override public void onEvent(OrderEvent event, long sequence, boolean endOfBatch) throws Exception { // 模拟业务处理:校验、入库、扣减库存等等 processOrder(event); } }核心的组装代码如下所示:
int bufferSize = 1024 * 8; Disruptor<OrderEvent> disruptor = new Disruptor<>( new OrderEventFactory(), bufferSize, Executors.defaultThreadFactory(), ProducerType.SINGLE, new BusySpinWaitStrategy() ); EventHandler<OrderEvent> handler = new OrderEventHandler(); disruptor.handleEventsWith(handler); disruptor.start(); RingBuffer<OrderEvent> ringBuffer = disruptor.getRingBuffer(); // 发布事件 for (long i = 0; i < 1000000; i++) { long sequence = ringBuffer.next(); try { OrderEvent event = ringBuffer.get(sequence); event.setOrderId(i); event.setAmount(100.0); event.setStatus("NEW"); } finally { ringBuffer.publish(sequence); } } disruptor.shutdown();这里有几个关键点需要特别说明:
首先,ProducerType.SINGLE这个参数极其重要。如果你的生产者只有一个线程,必须选 SINGLE,远程它对应的 Sequencer 实现是单生产者序号器,没有 CAS 竞争,发布速度会有质的提升。如果在单生产者模式下错误选择了 MULTI,性能会白白损耗在 CAS 上。
其次,BusySpinWaitStrategy是吞吐量最高但 CPU 占用也最疯狂的等待策略,它会持续自旋直到事件到达。要在生产环境使用它,必须先确认系统有足够的 CPU 核数可用,否则会拖垮整个机器。
最后,事件对象是复用的,所以你必须在消费者处理完事件后清理字段。上面代码里我写的clear()方法就是用来做这个的,否则下一个生产者写入旧字段时,消费者读到的可能是上一个事件残留的数据,这种 bug 极难排查。
3.2 多生产者场景下的正确姿势
现实中很少只有一个线程在产生订单,典型的场景是多业务线程同时往 Disruptor 里灌事件。这时你需要把 ProducerType 改成 MULTI,并用ringBuffer.tryNext()代替next(),因为在高竞争下直接阻塞等待并不可取。
Disruptor<OrderEvent> disruptor = new Disruptor<>( new OrderEventFactory(), bufferSize, new ThreadFactory() { private final AtomicInteger counter = new AtomicInteger(); @Override public Thread newThread(Runnable r) { return new Thread(r, "disruptor-worker-" + counter.getAndIncrement()); } }, ProducerType.MULTI, new YieldingWaitStrategy() );多生产者模式下,next()返回的序号不保证连续,因为多个线程可能同时申请到不同的槽位,且中间有空隙。你写事件的时候必须严格使用自己拿到的sequence去定位槽位,不能假设它是上一个序号加一。
另外要重视tryNext()的失败处理。如果队列满了,它会抛异常或返回负值,你需要决定是丢弃事件、重试还是采用背压机制。我在实际项目中倾向于让生产者先做一些本地聚合再批量发布,这样可以有效降低序号竞争次数。
3.3 事件清理与批量处理的小技巧
很多人在消费者里处理单个事件,忽略了endOfBatch这个参数。Disruptor 已经帮你把“批量”的可能性暴露出来了,如果你的消费者能处理一批事件而不是一个个来,吞吐量还能再上一个台阶。
endOfBatch为 false 时表示后面还有连续的事件,你可以先把它们缓存到本地 List,直到 endOfBatch 为 true 时统一批量刷到数据库或发送到下游。比如批量攒 100 个订单再批量入库,数据库 IO 次数直接减少一个数量级,整体 TPS 自然就上去了。
但在批量处理时务必留意:事件对象是复用的,一旦你批量缓存了它们的引用,而生产者在同一槽位上写入了新数据,你缓存的对象就被污染了。正确的做法是在批量处理前把字段拷贝到自定义对象中,或者直接把批量处理需要的字段拷贝成基本类型数组。这是个非常隐蔽的坑,我踩过不止一次。
4. 性能调优与 200 万 TPS 背后的细节
4.1 等待策略选型的思维模型
Disruptor 提供了好几种 WaitStrategy,它们是在延迟和 CPU 占用之间的权衡。BlockingWaitStrategy 内部用了锁,消费者会进入阻塞状态,CPU 占用低但延迟最大,适合对吞吐要求不高、机器核数紧张的场景。
SleepingWaitStrategy 则采用先自旋再睡眠的折中方案,它在吞吐量和 CPU 利用率之间做了平衡,是大多数后台任务的稳妥选择。YieldingWaitStrategy 会主动让出 CPU 时间片,适合低延迟且多核环境,ArbitraryWaitStrategy 同理。
BusySpinWaitStrategy 性能最高,但它完全不让出 CPU,如果只有一个核心在跑消费者线程,它会吃掉整个核的所有时间片。选型核心原则是:先看你的业务能承受多大延迟,再考虑 CPU 资源是否充裕。无脑上 BusySpin 往往会让整台机器的其他服务遭殃。
4.2 200 万 TPS 到底怎么测出来的
很多人看到“单机 200 万 TPS”这个数字会觉得不可思议,但你要理解这个 TPS 是在什么条件下测出来的:空转的消费者、纯内存事件、固定 batch size、极端优化的等待策略、无 GC 干扰的预热环境。这其实是 Disruptor 本身在裸机上的能力上限,不代表你的业务系统也能跑这么高。
我做过一组对比实验。在同一台 8 核机器上,Disruptor 空转时吞吐可以到 220 万 TPS,消费者加了一个简单的字段累加后降到 130 万 TPS,再给消费者加一个网络发送逻辑后只剩下 40 万 TPS,等到消费者开始写数据库,TPS 稳定在 2 万左右。
这说明什么?Disruptor 解决的是“事件搬运”的问题,而不是“业务处理”的问题。你在设计系统时,应该把耗时操作从消费者的关键路径中剥离,让 Disruptor 只做轻量级的转发和路由,真正重的写库、调外部服务放到独立的线程池里异步执行。
4.3 关于“tps 虚高”的冷思考
这个热词说出了很多人的切身体会。不少人一看到 Disruptor 的宣传数字就热血上头,把核心交易链路改成 Disruptor,结果上线后发现 TPS 根本没有翻倍,甚至因为背压问题变得更糟。如果你也遇到过这种情况,多半是踩了三个典型陷阱。
第一个陷阱是拿 Disruptor 的吞吐量当业务吞吐量。裸 RingBuffer 确实是百万级 TPS,但你的业务 handler 只要执行一次磁盘 IO,这个数字就会断崖式下跌。所以衡量系统瓶颈时必须看整条链路的吞吐,而不是中间件的单点能力。
第二个陷阱是忽略了 GC。虽然 RingBuffer 预分配了事件对象,但如果你的业务代码里仍然大量创建临时对象,比如把订单号转成 String、把金额包装成 BigDecimal,这些对象依然会在堆里堆积,最终触发 GC,导致全链路暂停。Disruptor 再快也架不住频繁的 Full GC。
第三个陷阱是盲目选用 BusySpinWaitStrategy。它吃掉一个完整核心的 CPU,在容器环境或超卖环境下,可能直接导致其他线程饥饿。正确做法是先用 YieldingWaitStrategy,然后逐步尝试 SleepWaitStrategy,找到业务能接受的可调参数。
4.4 启动预热与 CPU 亲和性设置
高性能程序的通病是“跑久了才快”。Disruptor 第一次启动时,事件数组和消费者线程都有缓存冷启动问题,如果直接用生产流量压测,前几秒的数据会很难看。建议在系统启动时先发布一批预热事件,让 CPU 缓存把热点数据加载起来,再做正式发布。
如果你的系统允许,还可以给消费者线程设置 CPU 亲和性,绑定在固定的物理核心上。Linux 下可以用taskset命令完成,或者用 JNI 调用sched_setaffinity。这样避免了线程在核心间来回切换导致的缓存失效。
我实测启用 CPU 亲和性后,同样配置下的吞吐量提升了 8% 到 12%。但要注意,这个提升幅度取决于你的 CPU 架构和系统负载,不是每个环境都能拿到同样的收益。建议在不同机器上都跑一轮基准测试再决定是否启用。
5. 常见问题与排查技巧实录
5.1 消费者处理太慢导致事件覆盖
这是 Disruptor 新手的头号问题。有人觉得 RingBuffer 是无穷无尽的,拼命往里写,结果消费者跟不上,生产者申请到的序号已经超过了消费者的消费进度,事件在写入时就把还没处理的旧数据覆盖了。
排查方法很简单:在消费者里记录自己的 Sequence,在生产者里对比一下两者差值,如果差值长期接近 bufferSize,说明队列持续处于满状态,背压已经开始生效。此时你需要增加消费者数量,或者提升单个消费者的批量处理能力。
如果是单消费者模式,还可以考虑把事件拆分成多个独立 Disruptor,每个负责一种类型的事件,降低单一队列的负载。但要注意拆分后会增加事件流转的复杂度,不是所有场景都划算。
5.2 事件对象复用带来的脏数据
前面提到的clear()方法如果漏掉某个字段,就会出现“上一个事件的数据串到了下一个事件”的诡异现象。这种 bug 最坑的地方在于它不是稳定复现的,只有特定顺序下才出现,测试时极难发现。
我的经验是:给 OrderEvent 写一个fill()方法和一个clear()方法,在生产者里只调fill(),消费者处理完必须调clear(),并且通过单元测试验证 clear 后所有字段都恢复默认值。不要依赖编译器帮你检查,这种事情只有运行时才能暴露。
如果你使用的是 Kotlin 或 Scala,还要注意编译器可能会对字段做自动装箱和拆箱,这会给性能带来额外负担。建议事件对象里只放基本类型字段,不要放对象引用,除非你非常清楚自己在做什么。
5.3 等待策略导致 CPU 飙到 100%
这是压测时最常见的现象。你用 BusySpinWaitStrategy 跑出了一个亮眼的 TPS,但观察top发现消费者线程 CPU 占用一直顶在 100%,如果机器上还有其他服务,整体响应就会变得不可控。
解决办法有两个方向。一是限流:在生产者端做流量控制,让事件到达速率不超过消费者的处理速率。二是换更温和的等待策略,比如 SleepingWaitStrategy 或 BlockingWaitStrategy。
根据我个人经验,绝大多数业务场景根本不需要 BusySpin。如果真的追求极限吞吐,与其靠等待策略让消费者干转,不如把事件批量攒满再通知消费者。用RingBuffer.PROGRESS_BARRIER或自定义的批处理通知机制,能让消费者在事件积压到一定数量后才被唤醒,CPU 利用率也能降下来。
5.4 压测数据与生产数据差距巨大的归因
不少人在博客里晒出单机百万 TPS 的成绩,但跑到自己的代码上只有几万,就开始怀疑硬件或配置。其实这个问题多半出在业务处理与测试条件的差异上。
我做压测时会分成四层来测量:第一层只测 RingBuffer 的发布和消费,第二层加字段读取,第三层加简单计算,第四层加网络或磁盘 IO。每一层分别记录 TPS 和延迟,这样才能定位瓶颈在哪里。如果你直接从第四层反推 Disruptor 的性能,那得出的数字自然就是“tps 虚高”了。
有一点要特别提示:Disruptor 的吞吐量受 CPU 主频、核心数、缓存行大小影响很大,不同硬件之间的差异可以超过 50%,所以任何实验室数据都比不上你在一台真实机器上跑出来的数据。建议把基准测试脚本纳入 CI 流程,每次调优后跑一遍,用数据说话。
结尾
我在实际项目里大规模使用 Disruptor 已经有两年多,最大的感受是:它不是一个“装了就能快”的插件,而是一整套需要你理解并配合的高性能协作模型。RingBuffer 只是舞台,真正决定吞吐量的是你的业务逻辑如何在这个舞台上编排。如果你理解了缓存行填充和内存屏障的原理,理解了事件复用的必要性,理解了等待策略和背压控制的取舍,你就能在任意场景下把它调到最优。最开始接触 Disruptor 时我也会对着测出来的高数字兴奋,但踩过几次 tps 虚高的坑之后,我更倾向于在压测和上线前先想清楚业务链路的真实瓶颈在哪里,然后再决定值不值得用 Disruptor。如果你正在高速传输场景里寻找合适的中间件,建议先拿真实业务压一轮,再下结论。