LMAX Disruptor生产者未阻塞回绕,覆盖未读取数据问题求助
问题分析与解决方案
看起来你遇到的核心问题是生产者没有感知到消费者的处理进度,导致环形缓冲区被填满后直接回绕覆盖未处理的数据。让我一步步拆解原因并给出修复方案:
问题根源
Disruptor中生产者的ringbuffer.next()方法是否会阻塞,取决于它能否获取到足够的可用槽位——而可用槽位的计算依赖于最慢消费者的Sequence。但在你的实现里:
- 消费者维护了自己的
seq变量,但没有把这个Sequence注册到RingBuffer的"门控序列(Gating Sequences)"中 - RingBuffer无法感知到消费者的处理进度,会错误地认为所有槽位都是可用的,因此生产者会无限制地生产,直到覆盖旧数据
另外,你的消费者实现也不够规范:手动调用seqbar.waitFor(seq)并维护seq变量容易出错,而且没有把处理完成的序列反馈给Disruptor框架。
修复方案
我推荐两种修复方式,你可以根据场景选择:
方案1:使用Disruptor官方推荐的EventHandler(更简洁规范)
这种方式让Disruptor自动管理消费者的Sequence和屏障,不需要手动维护,也能自动让生产者感知消费者进度:
重构消费者为EventHandler
public class EventConsumer implements EventHandler<Event> { @Override public void onEvent(Event event, long sequence, boolean endOfBatch) throws Exception { // 处理event的逻辑,比如你之前在getData()里的操作 Data data = event.get(); ... Do stuff ... } }
重构启动逻辑
public class DisruptorTest { public static void main(String[] args) { // 创建Disruptor Disruptor<Event> disruptor = new Disruptor<>( Event.EVENT_FACTORY, 1024, Executors.defaultThreadFactory() ); // 注册消费者 disruptor.handleEventsWith(new EventConsumer()); // 启动Disruptor并获取RingBuffer RingBuffer<Event> ringbuffer = disruptor.start(); // 启动生产者线程 ExecutorService exec = Executors.newCachedThreadPool(); exec.submit(new Producer(ringbuffer)); // 主线程保持运行,避免进程退出 try { Thread.currentThread().join(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }
生产者代码无需修改
因为Disruptor已经把消费者的Sequence注册为门控序列,生产者调用ringbuffer.next()时,会自动检查最慢消费者的进度,当缓冲区满了就会阻塞,直到消费者处理完旧数据。
方案2:手动注册消费者的Sequence(适合自定义消费逻辑场景)
如果你必须手动维护消费者的Sequence(比如你的getData()是同步调用的场景),需要把消费者的Sequence注册到RingBuffer中:
修改Consumer类
public class Consumer { private final ExecutorService exec; private final RingBuffer<Event> ringbuffer; private final SequenceBarrier seqbar; // 用线程安全的Sequence对象代替long变量,方便注册到RingBuffer private final Sequence consumerSeq = new Sequence(0L); public Consumer() { exec = Executors.newCachedThreadPool(); Disruptor<Event> disruptor = new Disruptor<>( Event.EVENT_FACTORY, 1024, Executors.defaultThreadFactory() ); ringbuffer = disruptor.start(); // 关键:把消费者的Sequence加入到生产者的门控序列中 ringbuffer.addGatingSequences(consumerSeq); seqbar = ringbuffer.newBarrier(); Producer producer = new Producer(ringbuffer); exec.submit(producer); } public Data getData() throws InterruptedException { try { long nextSeq = seqbar.waitFor(consumerSeq.get()); Event e = ringbuffer.get(nextSeq); Data data = e.get(); // 更新消费者进度,告诉生产者该槽位已处理完成 consumerSeq.set(nextSeq); return data; } catch (TimeoutException | AlertException e) { // 根据业务场景处理异常 throw new RuntimeException(e); } } }
关键改动说明
- 用
Sequence对象代替long seq:Disruptor的Sequence是线程安全的,且能被RingBuffer实时感知 - 调用
ringbuffer.addGatingSequences(consumerSeq):让生产者在计算可用槽位时,必须考虑该消费者的处理进度 - 处理完事件后更新
consumerSeq:及时反馈进度,确保生产者不会提前覆盖未处理的数据
额外注意点
- 环形缓冲区的大小建议设置为2的幂(你已经用了1024,这是正确的),Disruptor的性能优化依赖于此
- 生产者代码中可以考虑添加异常捕获(比如
InsufficientCapacityException),不过注册门控序列后,默认的next()会自旋等待,不会抛出该异常
内容的提问来源于stack exchange,提问作者Alice Smith
相关产品推荐
相关产品推荐

