Disruptor:高性能队列的颠覆者

为什么需要 Disruptor?

场景

  1. 高性能日志收集与处理(如日志聚合系统)
  2. 金融交易系统中高频交易数据的异步处理
  3. 游戏领域高并发事件的分发与处理
  4. 大数据厂家下的数据流式处理(如实时数据清洗、转换)
  5. 文档中实时协同编辑
  6. 需要高效传递消息的各种场景......

传统解决方案

使用线程池处理高并发的线程交互

核心组件:

  1. 阻塞队列:如ArrayBlockingQueue(数组实现,有界),LinkedBlockingQueue(链表实现,无界)等,负责存储待处理事件 、任务。
  2. 线程池:负责从队列中获取任务并异步处理(如写入磁盘、网络请求)。

处理流程:

  1. 业务线程(生产者)将任务(日志事件、消息)放入阻塞队列,立即返回,不等待处理完成。
  2. 线程池中的工作线程(消费者)从队列中取任务并执行,实现生成与消费解耦。

局限性

  1. 锁竞争开销:阻塞队列内部依赖 ReentrantLock 保证线程安全,高并发下生产者与消费者抢锁会导致大量上下文切换,内核态用户态相互转换,吞吐量下降。
  2. GC频繁:任务对象(如 Runable、事件对象)动态创建,频繁入队、出队,会产生大量临时对象,触发 JVM 频繁 GC,增加风险。
  3. 无界队列风险:若使用无界队列,当消费速度跟不上生产速度时,队列无限膨胀,可能导致 OOM。
  4. 线程调度成本:线程池的工作线程由操作系统调度,多线程切换会消耗 CPU 资源,尤其是任务粒度较小时,调度开销占比更高

它的优势是什么?

Disruptor(颠覆者) 是一个基于内存的高性能异步处理队列,本质上是线程之间无锁传递消息的队列。颠覆了传统高并发消息交互问题解决方案。

优势:

  1. 环形缓冲区:一个固定大小的数组,启动时创建所有事件对象,运行时循环复用,无序频繁GC;
  2. 序号机制:生产者与消费者通过 "序号" 跟踪进度
  3. 无锁设计,生产者和消费者通过 CAS 获取环形缓冲区的序号,获取到之后,直接在该位置写入或者读取数据;

实际上它还采用了缓冲行填充(在关键变量前后添加占位符,确保 其独占一个缓冲行),彻底避免伪共享,由于其涉及知识较多,不再赘述,如需深入了解,请看美团技术团队Disruptor原文

核心组件

  1. RingBuffer:环形缓冲区,存储事件对象的环形数组
  2. Event:传递的数据载体,也就是事件
  3. EventFactory:预分配 Event 对象的工厂
  4. EventHandler:消费者逻辑(实现 onEvent() 方法)
  5. Sequence:跟踪生产者、消费者的处理进度
  6. SequenceBarrier:协调生产者与消费者的序号同步
  7. WaitStrategy:缓冲区满时的等待策略(如阻塞、自旋、休眠)

工作流程

初始化

  1. 定义 Event 事件
  2. 创建 EventFactory 工厂
  3. 配置 Disruptor(大小、线程池、等待策略)初始化生产者 Sequeuce 序号

生产

  1. 业务线程先通过 CAS 获取下一个可用的写入序号,序号是递增的
  2. 根据序号对缓冲区大小取模获取数组索引,填充数据
  3. 发布事件,将该序号标记为 "已发布",此时 SequenceBarrier 感知到新的可处理序号

消费

  1. 消费者线程通过 SequenceBarrier.waitFor(nextSequence) 等待可处理的事件
  • nextSequence:是消费者期望处理的下一个序号(例如,消费者当前已经处理到10,期望处理11)
  • SequenceBarrier 会检查:该序号是否 ≤ 生产者已发布的最大序号,且是否满足消费者依赖(如存在依赖链,需等待前置消费者处理完成)。若满足返回可处理的序号范围
  1. 消费者循环处理该范围内的事件
  2. 处理完成后,更新自身的 Sequeuce 序号,告知其它依赖的消费者,该序号已经处理完成

消费者线程从 RingBuffer 读取事件,执行写入文件、网络等操作

总结

  1. 环形缓冲区初始化:创建一个固定大小(如 8)的RingBuffer(索引范围为8),初始化序号为 0
  2. 生产者写入数据:生产者申请序号 0,将数据写入事件对象,提交后序号递增为 1,用于下次申请
  3. 消费者读取数据:消费者通过 SequenceBarrier 获取序号为 0 的事件,处理完成后提交,序号递增为1
  4. 环形缓冲区会循环使用,序号每次均是递增区分先后顺序,通过对 8 取模来获取环形缓冲区的索引
  5. 如果生产者追上消费者,消费者没有处理数据,则会根据等待策略进行相应的等待

实战

模拟高并发日志处理场景,Github 代码

引入依赖

xml
复制代码
<dependency> <groupId>com.lmax</groupId> <artifactId>disruptor</artifactId> <version>3.4.4</version> <!-- 稳定版本 --> </dependency>

定义事件

java
复制代码
package top.xiaoyijun.loggerdisruptor; import lombok.Data; /** * 日志事件 */ @Data public class LogEvent { private long id; // 日志ID private String message; // 日志内容 }

创建事件工厂

java
复制代码
package top.xiaoyijun.loggerdisruptor; import com.lmax.disruptor.EventFactory; /** * 日志事件工厂:负责创建日志事件实例 */ public class LoggerEventFactory implements EventFactory<LogEvent> { @Override public LogEvent newInstance() { // 预创建事件(Disruptor 会预先分配缓冲区大小的事件) return new LogEvent(); } }

消费者实现

java
复制代码
package top.xiaoyijun.loggerdisruptor.handler; import com.lmax.disruptor.EventHandler; import org.springframework.stereotype.Component; import top.xiaoyijun.loggerdisruptor.event.LogEvent; /** * 消费者 1:打印日志消费者实现,需要实现 EventHandler 接口 */ @Component public class PrintLogHandler implements EventHandler<LogEvent> { @Override public void onEvent(LogEvent logEvent, long sequence, boolean b) { // 处理事件,打印日志 System.out.println("[打印日志] ID = " + logEvent.getId() + ", 内容:" + logEvent.getMessage()); // 模拟处理耗时 try { Thread.sleep(10); }catch (Exception e){ Thread.currentThread().interrupt(); } } }
java
复制代码
package top.xiaoyijun.loggerdisruptor.handler; import com.lmax.disruptor.EventHandler; import org.springframework.stereotype.Component; import top.xiaoyijun.loggerdisruptor.event.LogEvent; import java.util.concurrent.atomic.AtomicLong; /** * 消费者2:统计日志总数 */ @Component public class CountLogHandler implements EventHandler<LogEvent> { private final AtomicLong count = new AtomicLong(0); // 原子类保证线程安全 @Override public void onEvent(LogEvent event, long sequence, boolean endOfBatch) { // 累加计数 count.incrementAndGet(); // 若到达最后一批事件,输出统计结果 if (endOfBatch) { System.out.println("\n[统计结果] 总日志数:" + count.get()); } } }

环形缓冲区配置

此次案例使用单生产者,多消费者实现,如需配置多生产者修改配置即可,测试时使用多个线程进行测试

java
复制代码
package top.xiaoyijun.loggerdisruptor.config; import com.lmax.disruptor.YieldingWaitStrategy; import com.lmax.disruptor.dsl.Disruptor; import com.lmax.disruptor.dsl.ProducerType; import jakarta.annotation.Resource; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import top.xiaoyijun.loggerdisruptor.factory.LoggerEventFactory; import top.xiaoyijun.loggerdisruptor.event.LogEvent; import top.xiaoyijun.loggerdisruptor.handler.CountLogHandler; import top.xiaoyijun.loggerdisruptor.handler.PrintLogHandler; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @Configuration public class EventDisruptorConfig { @Resource private PrintLogHandler printLogHandler; @Resource private CountLogHandler countLogHandler; @Bean("eventDisruptor") public Disruptor<LogEvent> messageModelRingBuffer() { // 1. 配置参数 int bufferSize = 1024; // 环形缓冲区大小(必须是2的幂,如1024、2048) ExecutorService executor = Executors.newCachedThreadPool(); // 消费者线程池 //2. 创建Disruptor实例 Disruptor<LogEvent> disruptor = new Disruptor<>( new LoggerEventFactory(), // 事件工厂 bufferSize, executor, ProducerType.SINGLE, // 单生产者(多生产者改为MULTI) new YieldingWaitStrategy() //等待策略 ); // 设置消费者 disruptor.handleEventsWith(printLogHandler, countLogHandler); // 开启 disruptor disruptor.start(); return disruptor; } }

生产者实现

实现生产者时,需要先获取序号,才可以填充事件,但是如果此时出现异常,将会导致该序列号对应事件会被永久占用,环形缓冲区会被耗尽,生产者阻塞,因此要使用 Translator 模式发布事件

java
复制代码
package top.xiaoyijun.loggerdisruptor.producer; import com.lmax.disruptor.RingBuffer; import com.lmax.disruptor.EventTranslatorOneArg; import com.lmax.disruptor.dsl.Disruptor; import jakarta.annotation.PreDestroy; import jakarta.annotation.Resource; import org.springframework.stereotype.Component; import top.xiaoyijun.loggerdisruptor.event.LogEvent; /** * 日志事件生产者:发布日志到Disruptor */ @Component public class LogEventProducer { @Resource(name = "eventDisruptor") private Disruptor<LogEvent> logEventDisruptor; /** * 发布事件 * @param id 消息id * @param message 消息内容 */ public void publishEvent(long id, String message){ // 获取 Disruptor 的环形缓冲区 RingBuffer<LogEvent> ringBuffer = logEventDisruptor.getRingBuffer(); // 获取可以生成的位置(手动获取可以生成的位置,如果出现异常,该序列号对应事件会被永久占用,环形缓冲区会被耗尽,生产者阻塞) long next = ringBuffer.next(); LogEvent logEvent = ringBuffer.get(next); logEvent.setId(id); logEvent.setMessage(message); // 发布事件 ringBuffer.publish(next); } /** * 优雅停机 */ @PreDestroy public void close(){ logEventDisruptor.shutdown(); } }
java
复制代码
package top.xiaoyijun.loggerdisruptor.producer; import com.lmax.disruptor.EventTranslatorOneArg; import com.lmax.disruptor.RingBuffer; import com.lmax.disruptor.dsl.Disruptor; import jakarta.annotation.PreDestroy; import jakarta.annotation.Resource; import org.springframework.stereotype.Component; import top.xiaoyijun.loggerdisruptor.event.LogEvent; /** * 日志事件生产者:发布日志到 Disruptor, 使用Translator模式发布事件 */ @Component public class LogEventProducer2 { @Resource(name = "eventDisruptor") private Disruptor<LogEvent> logEventDisruptor; /** * 事件转换器:封装事件设置逻辑,避免手动操作序列号 */ public static final EventTranslatorOneArg<LogEvent, LogData> TRANSLATOR = ((logEvent, l, logData) -> { logEvent.setId(logData.id); logEvent.setMessage(logData.message); }); /** * 发布事件 * * @param id 消息id * @param message 消息内容 */ public void publishEvent(long id, String message) { // 获取 Disruptor 的环形缓冲区 RingBuffer<LogEvent> ringBuffer = logEventDisruptor.getRingBuffer(); // 发布事件 ringBuffer.publishEvent(TRANSLATOR, new LogData(id, message)); } // 内部数据载体:封装待发布的参数 private static class LogData { private final long id; private final String message; public LogData(long id, String message) { this.id = id; this.message = message; } } /** * 优雅停机 */ @PreDestroy public void close(){ logEventDisruptor.shutdown(); } }

测试

java
复制代码
package top.xiaoyijun.loggerdisruptor; import jakarta.annotation.Resource; import org.springframework.boot.CommandLineRunner; import org.springframework.stereotype.Component; import top.xiaoyijun.loggerdisruptor.producer.LogEventProducer; import top.xiaoyijun.loggerdisruptor.producer.LogEventProducer2; @Component public class EventPublishTest implements CommandLineRunner { // 注入生产者 @Resource private LogEventProducer logEventProducer; @Resource private LogEventProducer2 logEventProducer2; @Override public void run(String... args) throws Exception { // 模拟发布10条日志 for (int i = 1; i <= 10; i++) { //logEventProducer.publishEvent(i, "这是第" + i + "条日志"); logEventProducer2.publishEvent(i, "这是第" + i + "条日志"); // 间隔50ms Thread.sleep(50); } } }

测试结果

java
复制代码
[统计结果] 总日志数:1 [打印日志] ID = 1, 内容:这是第1条日志 [打印日志] ID = 2, 内容:这是第2条日志 [统计结果] 总日志数:2 [统计结果] 总日志数:3 [打印日志] ID = 3, 内容:这是第3条日志 [统计结果] 总日志数:4 [打印日志] ID = 4, 内容:这是第4条日志 [统计结果] 总日志数:5 [打印日志] ID = 5, 内容:这是第5条日志 [统计结果] 总日志数:6 [打印日志] ID = 6, 内容:这是第6条日志 [统计结果] 总日志数:7 [打印日志] ID = 7, 内容:这是第7条日志 [统计结果] 总日志数:8 [打印日志] ID = 8, 内容:这是第8条日志 [统计结果] 总日志数:9 [打印日志] ID = 9, 内容:这是第9条日志 [打印日志] ID = 10, 内容:这是第10条日志 [统计结果] 总日志数:10

参考文献

高性能队列——Disruptor

鱼皮编程导航云图库项目

0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
Lucky
作者分享
Day 10 ✅ 今天做了: 学习安卓开发 activity,组件页面交互 六级词汇 200 📚 今日感悟:加油
1
Day 9 ✅ 今天做了:学习安卓开发基础内容 ⏰ 明天计划:继续学习kotlin、安卓开发内容,有时间的话要总结mq 📚 今日感悟:加油
1
Day 9 ✅ 今天做了: 学习kotlin基础 六级词汇 200 📚 今日感悟:加油
1
Day 8 前几天没有打卡,主要是没有学习,去参加了软考架构师阅卷,期间发生了一个有意思的事情: 试卷会通过发题机发到我们手里,但是当我们把试卷批改完成之后,发现还有20份试卷没有批改,所有人的电脑上都没有题目了,这时有意思的来了,一堆人在那里想这是哪里出现了bug,然后项目组联系技术人员,查了半小时,发现有一位老师,在开始的时候登错了一个账号,发现之后又登录了自己的账号,导致登错的账号被发题机发了20题,最后也是技术人员重新把题目发到了其它账号上。 这里需要吐槽几点 1.我们发现这个发题机有bug,总共5200份试卷,我们最后看到后台却批改了5400份试卷,说明有人重复批改了,个人感觉像是缓存的问题,因为中间网络波动了几次。 2.这后台系统,都不能看到题目发到了哪个账号上,导致我们所有人在那等了半个多小时,难受 3.这发题系统,最后将那20道题发到了三个账号上,其它人只能干等着,这里算法需要进行优化 ✅ 今天做了:学习完成kafka ⏰ 本月计划:完成RabbitMQ、RocketMQ、Kafka对比总结 📚 今日感悟:加油
1
Day 7 ✅ 今天做了: 明天有面试,发现八股忘光了,紧急复习:java基础、集合、JUC、JVM、Spring(部分) 六级词汇:100 ⏰ 明天计划:明天争取全部复习完成 📚 今日感悟:加油
1
下载 APP