做高并发系统,一提到异步解耦、流量削峰,很多人第一反应就是上 Kafka、RocketMQ。
但现实是,大量场景根本不需要重量级的分布式消息队列:
下单成功后发通知、记日志、统计数据;接口请求埋点采集;批量数据异步处理……
这些都是进程内的异步操作,上 MQ 太重,维护成本高,延迟也上去了。
用 JDK 自带的 ArrayBlockingQueue?高并发下锁竞争严重,上下文切换开销大,吞吐量上不去,根本扛不住峰值流量。
今天给大家介绍一款工业级的高性能内存队列框架——LMAX Disruptor。
它是英国外汇交易公司 LMAX 开源的无锁并发框架,专为极致低延迟、高吞吐量设计,官方基准测试下单线程每秒可处理 600 万+ 事件,是 JDK 原生队列的几十倍性能。
SpringBoot 整合极其简单,几行代码就能接入,完美解决进程内高并发异步处理的性能瓶颈。
一、传统阻塞队列,为什么扛不住高并发?
在说 Disruptor 之前,先搞明白为什么我们常用的 BlockingQueue 在高并发下性能上不去。
JDK 里的 ArrayBlockingQueue、LinkedBlockingQueue 是大家最常用的内存队列,但它们天生有三个绕不开的性能瓶颈:
1. 重量级锁,上下文切换开销大
ArrayBlockingQueue 底层用 ReentrantLock 保证线程安全,不管是生产者放数据还是消费者取数据,都要先抢锁。
高并发下,大量线程阻塞等待锁,CPU 大量时间花在线程上下文切换和锁竞争上,真正处理业务的时间占比很低。
并发越高,锁的开销越大,吞吐量直接被锁卡死。
2. 伪共享,缓存频繁失效
CPU 缓存是以「缓存行(Cache Line)」为单位加载的,通常 64 字节。
普通队列的数组里,相邻元素会被加载到同一个缓存行。
当生产者修改队尾下标、消费者修改队头下标时,两个变量在同一个缓存行里,导致对方的缓存失效,必须重新从主内存读取。
看似不起眼的缓存失效,在高频率读写下,会带来数量级的性能损耗。
3. 对象频繁创建,GC 压力大
LinkedBlockingQueue 是链表结构,每入队一个元素就要新建一个 Node 对象,高并发下对象创建速度极快,导致 Young GC 频繁,严重影响吞吐量和延迟稳定性。
而 Disruptor 的设计,就是从根源上把这三个问题全部解决了。
二、Disruptor 凭什么这么快?
Disruptor 不是什么黑科技,而是把计算机底层的性能优化原理做到了极致,每一个设计都精准命中传统队列的性能痛点。
1. 环形数组 RingBuffer:预分配,零 GC
Disruptor 的核心数据结构是环形数组(RingBuffer),而非链表。
-
• 数组长度固定为 2 的 n 次方,通过位运算计算下标,比取模运算快得多; -
• 初始化时就预创建所有事件对象,全程复用,不会频繁创建销毁对象,几乎没有 GC 压力; -
• 数组天然对 CPU 缓存友好,顺序访问时缓存命中率极高。
简单理解:
环形数组就是一个循环转圈的数组,生产者和消费者各自维护自己的序列号,写满了就从头覆盖(配合消费者追赶策略,不会乱覆盖),全程不用扩容、不用新建对象。
2. 无锁设计:CAS + 序列号,告别锁竞争
Disruptor 完全不用重量级锁,靠 Sequence 序列号 + CAS 原子操作 实现并发安全。
-
• 每个生产者、消费者都维护自己的序列号,表示自己处理到了哪个位置; -
• 生产者写入前,通过 CAS 申请下一个可用位置,成功就写入,失败就自旋等待; -
• 消费者只需要盯着自己的序列号,有新事件就处理,没有就等待。
整个过程没有锁,只有轻量级的 CAS 和自旋,高并发下的开销比重量级锁小得多。
3. 缓存行填充:从根源解决伪共享
为了解决伪共享问题,Disruptor 对序列号变量做了缓存行填充(Padding)。
简单说就是在序列号变量的前后,各填充一堆无意义的 long 类型字段,让每个序列号独占一个 64 字节的缓存行。
这样生产者和消费者的序列号不会出现在同一个缓存行里,修改时不会导致对方缓存失效,从根源上解决了伪共享的性能损耗。
4. 单线程消费:避免线程安全开销
Disruptor 推荐每个消费者用单线程处理事件,不用考虑线程安全问题,省掉了锁、同步的开销。
同时支持灵活的消费者编排:可以并行消费、串行依赖、分组消费,轻松实现复杂的处理流程。
性能对比到底差多少?
官方基准测试(单生产者单消费者,纯内存转发):
|
|
|
|
|
|
|
|
|
|
|
|
注意:600万是纯事件转发的极限性能,实际业务因为有逻辑、IO,吞吐量会低很多,但相比原生队列依然有数倍提升。
三、SpringBoot 完整整合
原理讲完了,直接上可运行代码。基于 Spring Boot 3.x + Disruptor 3.4.4 实现一套完整的订单异步处理示例。
3.1 Maven 核心依赖
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- Disruptor 核心依赖 -->
<dependency>
<groupId>com.lmax</groupId>
<artifactId>disruptor</artifactId>
<version>3.4.4</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
3.2 第一步:定义事件实体
事件就是队列里传递的数据对象,预创建复用,用完需清空字段避免脏数据。
/**
* 订单事件
*/
@Data
public class OrderEvent {
/** 订单ID */
private Long orderId;
/** 用户ID */
private Long userId;
/** 订单金额 */
private BigDecimal amount;
/** 事件类型:创建、支付、取消 */
private String eventType;
/**
* 清理数据,对象复用时调用
*/
public void clear() {
this.orderId = null;
this.userId = null;
this.amount = null;
this.eventType = null;
}
}
3.3 第二步:事件工厂
Disruptor 初始化时,会用工厂预创建所有事件对象,填满整个环形数组。
/**
* 订单事件工厂
*/
public class OrderEventFactory implements EventFactory<OrderEvent> {
@Override
public OrderEvent newInstance() {
return new OrderEvent();
}
}
3.4 第三步:定义消费者(事件处理器)
消费者实现 EventHandler 接口,单线程消费事件。
示例中两个并行消费者:一个记录订单日志,一个发送短信通知,互不干扰。
订单日志处理器
@Component
@Slf4j
public class OrderLogEventHandler implements EventHandler<OrderEvent> {
@Override
public void onEvent(OrderEvent event, long sequence, boolean endOfBatch) {
log.info("[订单日志] 记录订单事件:orderId={}, type={}",
event.getOrderId(), event.getEventType());
// 业务逻辑:异步记录操作日志、落库等
}
}
短信通知处理器
@Component
@Slf4j
public class OrderSmsEventHandler implements EventHandler<OrderEvent> {
@Override
public void onEvent(OrderEvent event, long sequence, boolean endOfBatch) {
log.info("[短信通知] 给用户{}发送订单通知:orderId={}",
event.getUserId(), event.getOrderId());
// 业务逻辑:调用短信网关发送通知
}
}
3.5 第四步:Disruptor 配置类
核心配置,构建 RingBuffer,绑定消费者,管理生命周期。
@Configuration
@RequiredArgsConstructor
public class DisruptorConfig {
private final OrderLogEventHandler logEventHandler;
private final OrderSmsEventHandler smsEventHandler;
/**
* 环形缓冲区大小,必须是 2 的 n 次方
* 根据业务峰值调整,建议为峰值流量的 2~3 倍
*/
private static final int RING_BUFFER_SIZE = 1024;
@Bean
public RingBuffer<OrderEvent> orderEventRingBuffer() {
// 1. 构建 Disruptor 实例
Disruptor<OrderEvent> disruptor = new Disruptor<>(
new OrderEventFactory(),
RING_BUFFER_SIZE,
// 自定义线程工厂,命名线程方便排查问题
r -> {
Thread thread = new Thread(r);
thread.setName("disruptor-order-handler");
thread.setDaemon(true);
return thread;
},
// 生产者类型:单生产者用 SINGLE 性能更好,多生产者用 MULTI
ProducerType.SINGLE,
// 等待策略:阻塞策略,CPU 占用低,适合普通业务
new BlockingWaitStrategy()
);
// 2. 配置消费者:两个消费者并行消费所有事件
disruptor.handleEventsWith(logEventHandler, smsEventHandler);
// 3. 所有消费完成后清理事件对象,防止脏数据
disruptor.handleEventsWith((event, sequence, endOfBatch) -> event.clear());
// 4. 全局异常处理器:避免单个异常搞垮消费者线程
disruptor.setDefaultExceptionHandler(new ExceptionHandler<OrderEvent>() {
@Override
public void handleEventException(Throwable ex, long sequence, OrderEvent event) {
log.error("Disruptor消费异常,sequence={}, event={}", sequence, event, ex);
}
@Override
public void handleOnStartException(Throwable ex) {
log.error("Disruptor启动异常", ex);
}
@Override
public void handleOnShutdownException(Throwable ex) {
log.error("Disruptor关闭异常", ex);
}
});
// 5. 启动 Disruptor
disruptor.start();
// 6. 注册 JVM 关闭钩子,优雅停机
Runtime.getRuntime().addShutdownHook(new Thread(disruptor::shutdown));
return disruptor.getRingBuffer();
}
}
3.6 第五步:生产者组件
封装事件投递逻辑,业务代码只调用生产者,不感知 Disruptor 底层细节。
@Component
@RequiredArgsConstructor
public class OrderEventProducer {
private final RingBuffer<OrderEvent> ringBuffer;
/**
* 发布订单事件(传统写法)
*/
public void publish(Long orderId, Long userId, BigDecimal amount, String eventType) {
// 1. 申请下一个序列号
long sequence = ringBuffer.next();
try {
// 2. 获取对应位置的预创建对象,直接赋值
OrderEvent event = ringBuffer.get(sequence);
event.setOrderId(orderId);
event.setUserId(userId);
event.setAmount(amount);
event.setEventType(eventType);
} finally {
// 3. 发布事件,必须放 finally,防止异常导致序列号错乱
ringBuffer.publish(sequence);
}
}
/**
* Lambda 风格发布,更简洁
*/
public void publishEvent(Consumer<OrderEvent> consumer) {
ringBuffer.publishEvent((event, sequence) -> consumer.accept(event));
}
}
3.7 第六步:业务层调用示例
下单接口同步处理核心逻辑,异步处理旁路操作,解耦又提效。
@Service
@RequiredArgsConstructor
@Slf4j
public class OrderService {
private final OrderEventProducer eventProducer;
/**
* 创建订单
*/
public OrderVO createOrder(OrderCreateDTO dto) {
// 1. 同步处理核心逻辑:校验、扣库存、生成订单
Order order = doCreateOrder(dto);
log.info("订单创建成功,orderId={}", order.getId());
// 2. 异步投递事件:日志、通知等旁路操作,不影响主流程
eventProducer.publish(
order.getId(),
order.getUserId(),
order.getAmount(),
"CREATE"
);
// 3. 返回结果
return new OrderVO(order.getId());
}
private Order doCreateOrder(OrderCreateDTO dto) {
// 核心订单创建逻辑
return new Order(1L, dto.getUserId(), dto.getAmount());
}
}
3.8 测试接口
@RestController
@RequestMapping("/order")
@RequiredArgsConstructor
public class OrderController {
private final OrderService orderService;
@PostMapping("/create")
public Result<OrderVO> createOrder(@RequestBody OrderCreateDTO dto) {
return Result.success(orderService.createOrder(dto));
}
}
至此,一套完整的 SpringBoot + Disruptor 异步处理框架就落地了。
业务代码只需要注入生产者调用发布方法,所有异步处理逻辑都在消费者里,完全解耦,性能拉满。
四、消费者灵活编排,应对复杂业务
Disruptor 的强大不止于快,更在于灵活的消费者依赖编排,轻松实现复杂的处理流程。
4.1 串行依赖:严格按顺序执行
比如订单事件要先做数据校验,再做业务处理,最后发通知,严格按顺序执行:
// 配置串行依赖:校验 -> 业务处理 -> 通知
disruptor.handleEventsWith(validateHandler)
.then(businessHandler)
.then(notifyHandler);
上一个消费者处理完,下一个才会收到事件,天然保证执行顺序,不用自己写依赖逻辑。
4.2 分组消费:分摊压力,集群模式
如果单消费者处理不过来,可以用 WorkerPool 实现分组消费。
同一个组内的消费者分摊事件,每个事件只被组内一个消费者处理,类似消息队列的集群消费模式,线性提升消费能力。
WorkerPool<OrderEvent> workerPool = new WorkerPool<>(
ringBuffer,
ringBuffer.newBarrier(),
new FatalExceptionHandler(),
new OrderWorkHandler(),
new OrderWorkHandler()
);
workerPool.start(executor);
4.3 批量消费:进一步提升吞吐量
消费者可以通过 endOfBatch 标记判断批次结束,攒一批数据再批量落库、批量调用接口,大幅减少 IO 次数,吞吐量再上一个台阶。
@Override
public void onEvent(OrderEvent event, long sequence, boolean endOfBatch) {
batch.add(event);
if (endOfBatch || batch.size() >= 100) {
// 批量处理
batchProcess(batch);
batch.clear();
}
}
4.4 等待策略选型
不同等待策略平衡 CPU 占用和延迟,根据业务场景选择:
-
• BlockingWaitStrategy:默认,阻塞等待,CPU 占用低,延迟一般,适合普通业务场景; -
• YieldingWaitStrategy:自旋 + 让出 CPU,延迟更低,CPU 占用适中,推荐性能要求高的场景; -
• BusySpinWaitStrategy:持续自旋,延迟最低,CPU 占用最高,适合极致低延迟场景; -
• SleepingWaitStrategy:自旋 + 睡眠,CPU 占用最低,延迟波动大,适合对延迟不敏感的后台任务。
五、注意事项
Disruptor 性能强,但用不对反而会出大问题,这些都是生产环境踩出来的经验。
1:消费者里做阻塞 IO,性能直接报废
Disruptor 消费者是单线程串行执行的,如果在消费者里调用慢接口、查数据库、sleep,会把整个消费线程卡住,后面的事件全部堆积。
✅ 解决:慢操作单独丢线程池处理,或者拆分多个消费者,不要让重逻辑阻塞轻逻辑。
2:环形缓冲区大小乱设
缓冲区太小,生产速度快于消费速度时,生产者会阻塞等待,等于白用 Disruptor;太大浪费内存。
✅ 解决:大小必须是 2 的 n 次方,按峰值流量的 2~3 倍设置,比如峰值每秒 1 万请求,设成 4096、8192 足够,不是越大越好。
3:事件对象不复用,脏数据满天飞
RingBuffer 里的对象是预创建复用的,如果用完不清理字段,下次取到的对象可能带着上次的数据,导致业务逻辑错误。
✅ 解决:在所有消费者处理完之后,加一个清理处理器,调用 event.clear() 清空字段。
4:多生产者用了 SINGLE 模式
生产者类型设成 SINGLE,但实际有多个线程调用发布事件,会导致序列号错乱,数据丢失。
✅ 解决:单生产者用 SINGLE(性能更好),多生产者必须改成 MULTI。
5:把 Disruptor 当分布式 MQ 用
很多人觉得 Disruptor 快,就想用来替代 Kafka,这是完全的误用。
Disruptor 是进程内内存队列,宕机数据全丢,不支持持久化、不支持跨服务、不支持消费堆积。
✅ 定位:进程内异步解耦、低延迟流量削峰;跨服务、持久化场景老老实实用 MQ。
6:异常不处理,消费者线程直接挂
消费者里抛出未捕获的异常,默认会导致整个消费者线程终止,后面的事件再也没人处理。
✅ 解决:配置全局异常处理器,捕获异常、打印日志、告警,不要让单个异常搞垮整个消费者。
六、全文总结
Disruptor 不是用来替代消息队列的,它解决的是进程内高并发、低延迟的事件流转问题。
它的核心价值,是用极致的工程优化(环形数组、无锁、缓存行填充),把内存队列的性能推到了极致,填补了 JDK 原生队列和分布式 MQ 之间的空白。
对于高并发系统来说,它是非常趁手的工具:
-
• 开发简单,几行配置就能接入; -
• 性能强劲,比原生队列高数倍吞吐量; -
• 编排灵活,轻松实现复杂的消费者依赖关系。
但也要记住:技术是为业务服务的,不要为了炫技盲目引入。
如果你的系统有高并发的进程内异步需求,BlockingQueue 已经扛不住了,Disruptor 绝对是最优解。
高并发、性能优化是后端进阶的核心方向,很多系统的瓶颈,往往就藏在这些基础组件的选型里。
后续持续更新高并发实战专栏:无锁编程、内存队列、分布式缓存、流量削峰全套生产源码干货。
喜欢高并发、性能优化、SpringBoot 实战内容,欢迎点赞、收藏、关注,持续跟进后端进阶开发专栏!
微信赞赏
支付宝扫码领红包
