Disruptor:高性能队列的颠覆者
为什么需要 Disruptor?
场景
- 高性能日志收集与处理(如日志聚合系统)
- 金融交易系统中高频交易数据的异步处理
- 游戏领域高并发事件的分发与处理
- 大数据厂家下的数据流式处理(如实时数据清洗、转换)
- 文档中实时协同编辑
- 需要高效传递消息的各种场景......
传统解决方案
使用线程池处理高并发的线程交互
核心组件:
- 阻塞队列:如ArrayBlockingQueue(数组实现,有界),LinkedBlockingQueue(链表实现,无界)等,负责存储待处理事件 、任务。
- 线程池:负责从队列中获取任务并异步处理(如写入磁盘、网络请求)。
处理流程:
- 业务线程(生产者)将任务(日志事件、消息)放入阻塞队列,立即返回,不等待处理完成。
- 线程池中的工作线程(消费者)从队列中取任务并执行,实现生成与消费解耦。
局限性
- 锁竞争开销:阻塞队列内部依赖 ReentrantLock 保证线程安全,高并发下生产者与消费者抢锁会导致大量上下文切换,内核态用户态相互转换,吞吐量下降。
- GC频繁:任务对象(如 Runable、事件对象)动态创建,频繁入队、出队,会产生大量临时对象,触发 JVM 频繁 GC,增加风险。
- 无界队列风险:若使用无界队列,当消费速度跟不上生产速度时,队列无限膨胀,可能导致 OOM。
- 线程调度成本:线程池的工作线程由操作系统调度,多线程切换会消耗 CPU 资源,尤其是任务粒度较小时,调度开销占比更高
它的优势是什么?
Disruptor(颠覆者) 是一个基于内存的高性能异步处理队列,本质上是线程之间无锁传递消息的队列。颠覆了传统高并发消息交互问题解决方案。
优势:
- 环形缓冲区:一个固定大小的数组,启动时创建所有事件对象,运行时循环复用,无序频繁GC;
- 序号机制:生产者与消费者通过 "序号" 跟踪进度
- 无锁设计,生产者和消费者通过 CAS 获取环形缓冲区的序号,获取到之后,直接在该位置写入或者读取数据;
实际上它还采用了缓冲行填充(在关键变量前后添加占位符,确保 其独占一个缓冲行),彻底避免伪共享,由于其涉及知识较多,不再赘述,如需深入了解,请看美团技术团队Disruptor原文。
核心组件
- RingBuffer:环形缓冲区,存储事件对象的环形数组
- Event:传递的数据载体,也就是事件
- EventFactory:预分配 Event 对象的工厂
- EventHandler:消费者逻辑(实现 onEvent() 方法)
- Sequence:跟踪生产者、消费者的处理进度
- SequenceBarrier:协调生产者与消费者的序号同步
- WaitStrategy:缓冲区满时的等待策略(如阻塞、自旋、休眠)
工作流程
初始化
- 定义 Event 事件
- 创建 EventFactory 工厂
- 配置 Disruptor(大小、线程池、等待策略)初始化生产者 Sequeuce 序号
生产
- 业务线程先通过 CAS 获取下一个可用的写入序号,序号是递增的
- 根据序号对缓冲区大小取模获取数组索引,填充数据
- 发布事件,将该序号标记为 "已发布",此时 SequenceBarrier 感知到新的可处理序号
消费
- 消费者线程通过 SequenceBarrier.waitFor(nextSequence) 等待可处理的事件
- nextSequence:是消费者期望处理的下一个序号(例如,消费者当前已经处理到10,期望处理11)
- SequenceBarrier 会检查:该序号是否 ≤ 生产者已发布的最大序号,且是否满足消费者依赖(如存在依赖链,需等待前置消费者处理完成)。若满足返回可处理的序号范围
- 消费者循环处理该范围内的事件
- 处理完成后,更新自身的 Sequeuce 序号,告知其它依赖的消费者,该序号已经处理完成
消费者线程从 RingBuffer 读取事件,执行写入文件、网络等操作
总结
- 环形缓冲区初始化:创建一个固定大小(如 8)的RingBuffer(索引范围为8),初始化序号为 0
- 生产者写入数据:生产者申请序号 0,将数据写入事件对象,提交后序号递增为 1,用于下次申请
- 消费者读取数据:消费者通过 SequenceBarrier 获取序号为 0 的事件,处理完成后提交,序号递增为1
- 环形缓冲区会循环使用,序号每次均是递增区分先后顺序,通过对 8 取模来获取环形缓冲区的索引
- 如果生产者追上消费者,消费者没有处理数据,则会根据等待策略进行相应的等待
实战
模拟高并发日志处理场景,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
参考文献
评论
问答助学
相关内容
0个评论
全部评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
内容推荐
Day 104✅ 今天做了:学习了Java反射及快速入门⏰ 明天计划:继续学习Java反射
0
Day 1✅ 今天做了:完成AI应用速通教程第一章、第二章、第三章前半节学习⏰ 明天计划:完成第三章后半节、第四章学习📚 今日感悟:理解机器学习、深度学习基础、transformer架构等
0
Day 12✅ 今天做了:AI 知识库面试通关120问(60-80);RAG知识库项目(知识库问答);MySQL 八股背诵;回溯算法题。⏰ 明天计划:AI 知识库面试通关120问(40-60);RAG知识库项目(Query优化);算法题。📚 今日感悟:验证RAG 知识库问答功能时,由于后端日志过于简略,人工排查难以定位问题所在,在借助 Claude Code后,仅提供项目目录与问题描述,AI 扫
1
day2(8.1)今日学习了rag的进阶知识,练习口诉了rag的流程。
0
独立开发一个企业级Web系统:斗篷系统(ABcloakPro)的开发实践
2
作者分享
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
