如何在EventClass中向LMAX Disruptor推送列表?MongoDB大数据场景Disruptor数据接入咨询
关于LMAX Disruptor的两个实践问题解答
1. 如何在EventClass中向LMAX Disruptor推送列表?
Disruptor的核心逻辑是围绕**事件(Event)**运转的,要推送列表数据,核心有两种思路,你可以根据业务场景选择:
方案1:将整个列表作为单个Event的字段
如果你的列表是一个完整的业务单元(比如一批待入库的订单),直接把列表塞进自定义Event里最省心。
先定义包含列表的Event类(注意要写reset方法,Disruptor会复用Event实例,避免频繁GC):
public class BatchDataEvent { private List<YourBizModel> dataList; // 重置方法,用于复用Event实例 public void reset(List<YourBizModel> dataList) { this.dataList = dataList; } public List<YourBizModel> getDataList() { return dataList; } }
生产者推送时直接把列表传入Event:
// 生产者发布事件 disruptor.publishEvent((event, sequence) -> { List<YourBizModel> yourDataList = // 从外部获取的列表数据 event.reset(yourDataList); });
消费者处理时直接读取列表:
public class BatchDataHandler implements EventHandler<BatchDataEvent> { @Override public void onEvent(BatchDataEvent event, long sequence, boolean endOfBatch) throws Exception { List<YourBizModel> dataList = event.getDataList(); // 这里做批量处理,比如批量写入MongoDB } }
方案2:拆分列表为单个Event推送
如果列表里的每个元素都是独立的业务事件(比如单条用户行为日志),拆分后逐个推送更贴合Disruptor的高性能设计:
先定义单条数据的Event类:
public class SingleDataEvent { private YourBizModel data; public void reset(YourBizModel data) { this.data = data; } public YourBizModel getData() { return data; } }
生产者遍历列表逐个发布:
List<YourBizModel> yourDataList = // 你的数据源 for (YourBizModel item : yourDataList) { disruptor.publishEvent((event, sequence) -> { event.reset(item); }); }
这种方式适合需要对每条数据做独立处理的场景,Disruptor的高吞吐量能轻松支撑批量推送。
2. 如何将自定义数据融入LMAX Disruptor机制并结合MongoDB存储海量数据?
结合MongoDB存储海量数据的场景,核心是围绕你的业务数据定义Event,通过Disruptor的生产者-消费者模型完成数据接收与持久化。下面给你一套完整的实践流程:
步骤1:定义自定义业务Event
根据MongoDB的文档结构,定义对应的Event类,务必实现reset方法:
public class MongoBizEvent { // 对应MongoDB文档的字段 private String docId; private String content; private long createTime; // 重置方法,复用Event实例 public void reset(String docId, String content, long createTime) { this.docId = docId; this.content = content; this.createTime = createTime; } // Getter方法 public String getDocId() { return docId; } public String getContent() { return content; } public long getCreateTime() { return createTime; } }
步骤2:创建EventFactory
Disruptor需要通过Factory来初始化环形缓冲区的Event实例池:
public class MongoBizEventFactory implements EventFactory<MongoBizEvent> { @Override public MongoBizEvent newInstance() { return new MongoBizEvent(); } }
步骤3:编写生产者(接收外部数据并发布)
生产者负责接收外部数据,然后发布到Disruptor的环形缓冲区:
public class MongoDataProducer { private final RingBuffer<MongoBizEvent> ringBuffer; // 通过构造函数注入RingBuffer public MongoDataProducer(RingBuffer<MongoBizEvent> ringBuffer) { this.ringBuffer = ringBuffer; } // 发布单条数据的方法 public void publishData(String docId, String content, long createTime) { // 获取下一个可用的序列 long sequence = ringBuffer.next(); try { // 获取对应的Event实例 MongoBizEvent event = ringBuffer.get(sequence); // 重置Event内容 event.reset(docId, content, createTime); } finally { // 必须在finally块发布,避免序列泄漏 ringBuffer.publish(sequence); } } }
步骤4:编写消费者(读取Event并写入MongoDB)
海量数据场景下,一定要用MongoDB的批量插入来提升效率,消费者可以积累到一定数量再批量写入:
public class MongoDataHandler implements EventHandler<MongoBizEvent> { private final MongoCollection<Document> mongoCollection; private final List<Document> batchBuffer = new ArrayList<>(1000); // 批量大小按需调整 private static final int BATCH_THRESHOLD = 1000; public MongoDataHandler(MongoCollection<Document> mongoCollection) { this.mongoCollection = mongoCollection; } @Override public void onEvent(MongoBizEvent event, long sequence, boolean endOfBatch) throws Exception { // 将Event转换为MongoDB的Document Document doc = new Document("_id", event.getDocId()) .append("content", event.getContent()) .append("create_time", event.getCreateTime()); batchBuffer.add(doc); // 达到批量阈值或当前批次结束时,执行插入 if (batchBuffer.size() >= BATCH_THRESHOLD || endOfBatch) { mongoCollection.insertMany(batchBuffer); batchBuffer.clear(); } } }
步骤5:组装并启动Disruptor
把所有组件整合起来,启动Disruptor开始工作:
public class DisruptorMongoDemo { public static void main(String[] args) { // 1. 初始化MongoDB连接 MongoClient mongoClient = MongoClients.create("mongodb://localhost:27017"); MongoDatabase db = mongoClient.getDatabase("your_biz_db"); MongoCollection<Document> collection = db.getCollection("your_data_collection"); // 2. 配置Disruptor参数 int bufferSize = 1024 * 1024; // 环形缓冲区大小,必须是2的幂 ThreadFactory threadFactory = Executors.defaultThreadFactory(); // 3. 创建Disruptor实例 Disruptor<MongoBizEvent> disruptor = new Disruptor<>( new MongoBizEventFactory(), bufferSize, threadFactory, ProducerType.MULTI, // 支持多生产者,单生产者可以用SINGLE new BlockingWaitStrategy() // 等待策略,按需选择,低延迟场景可用YieldingWaitStrategy ); // 4. 注册消费者 disruptor.handleEventsWith(new MongoDataHandler(collection)); // 5. 启动Disruptor disruptor.start(); // 6. 获取生产者实例,模拟推送海量数据 RingBuffer<MongoBizEvent> ringBuffer = disruptor.getRingBuffer(); MongoDataProducer producer = new MongoDataProducer(ringBuffer); for (int i = 0; i < 1000000; i++) { producer.publishData( UUID.randomUUID().toString(), "biz_content_" + i, System.currentTimeMillis() ); } // 生产环境记得优雅关闭资源 // disruptor.shutdown(); // mongoClient.close(); } }
关键注意点
- 环形缓冲区大小:必须是2的幂,根据你的数据吞吐量调整,太小会导致生产者阻塞,太大浪费内存。
- 等待策略:不同策略对应不同性能特性,
BlockingWaitStrategy适合CPU紧张场景,YieldingWaitStrategy适合低延迟场景。 - 批量写入:海量数据场景下,批量插入能大幅降低MongoDB的IO开销,一定要用。
- Event复用:必须实现
reset方法,Disruptor复用Event实例能避免频繁GC影响性能。
内容的提问来源于stack exchange,提问作者patoCapongo93
相关产品推荐
相关产品推荐

