ReentrantReadWriteLock逻辑迁移至带屏障LMAX Disruptor的方法
方案可行性结论
完全可以迁移到LMAX Disruptor实现,你初步的思路大方向是对的,只要调整几个设计细节,在高并发场景下确实能获得比ReentrantReadWriteLock+ArrayList方案更好的性能:Disruptor本身就是为消除锁竞争、伪共享、不必要内存屏障设计的,核心路径全程无锁,多线程场景下的延迟稳定性、吞吐量都比显式锁方案高一个量级。
原有思路的正误校验
- 正确的部分:用RingBuffer替代ArrayList做核心存储、拆分生产/消费角色、用序列屏障做流程控制的思路完全符合Disruptor的设计范式,避开了从零实现无锁结构的常见坑。
- 需要调整的误区:
- 别把删除操作做成独立生产者:RingBuffer是固定容量的环形顺序存储结构,不支持像ArrayList那样随机删除中间条目后移动元素缩容,删除逻辑应该下沉到维护角色做逻辑标记,硬要在生产者侧做随机删除会彻底破坏RingBuffer的顺序写入特性,反而引入额外开销。
- 别把读取做成普通事件消费者:你的场景是按需读取,不是事件到达就被动消费,不需要给读线程分配独立消费序列,否则会和写入、删除流程争抢序列权限,反而拖慢性能。
- 没必要给写入、删除配置双生产者:删除不属于事件写入流的一部分,双生产者模式本身会引入额外CAS开销,写入端如果是单线程就直接用单生产者模式(性能最高),多线程写入再开多生产者模式即可。
最优实现思路
- 第一步:定义RingBuffer存储单元,除了业务字段,额外加两个标记位:
isValid(逻辑删除标记,默认true)、writeVersion(写入版本号,避免读到未提交的脏数据)。 - 第二步:角色拆分,各角色只做单一职责:
- 写入端:作为生产者,只负责申请RingBuffer的下一个可写槽位、写入业务数据、发布事件更新可读序列,不掺杂任何读、删逻辑。
- 定时删除端:作为独立维护角色,和读端共享序列定位权限,定时器触发时直接定位到目标槽位,把
isValid设为false即可,不需要移动RingBuffer指针,操作开销极低;遇到连续的已删除条目时,推进已清理序列,标记这些槽位可以被覆写复用。 - 读取端:不持有独立消费序列,按需查询时直接遍历当前可读区间内的所有槽位,过滤掉已删除、版本不匹配的脏数据,直接返回符合条件的结果,全程无锁。
- 第三步:参数配置:RingBuffer容量固定设为2的整数次幂(比如65536),开启Disruptor默认的伪共享填充;等待策略按业务选:CPU资源充足、追求低延迟用
YieldingWaitStrategy,CPU资源紧张、能接受毫秒级阻塞用BlockingWaitStrategy。RingBuffer写满时会自动等待清理序列推进,不会出现数据覆盖问题。
性能收益说明
不是所有场景迁移都能拿到收益:如果你的业务读写QPS低于1万/秒,迁移后几乎感知不到性能差异,甚至因为Disruptor的固定容量初始化开销,冷启动速度还不如ArrayList+锁的方案。
只有满足以下条件时,能拿到明确的性能提升:
- 多线程并发写入QPS超过10万,锁竞争已经成为性能瓶颈
- 读请求占比极高,对延迟稳定性要求高:锁方案下写线程持有写锁时,所有读线程都会阻塞,容易出现毫秒级毛刺;Disruptor读路径全程无锁,延迟波动稳定在纳秒级
- 删除逻辑能接受最终一致:逻辑删除不会立刻释放空间,等连续无效条目累积后才会统一标记可覆写,不需要强一致实时删除的场景完全适配
示例代码
import com.lmax.disruptor.*; import java.util.ArrayList; import java.util.List; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.function.Predicate; // RingBuffer存储单元定义 class DataEvent { private String bizData; private boolean isValid; private long writeVersion; public String getBizData() {return bizData;} public void setBizData(String bizData) {this.bizData = bizData;} public boolean isValid() {return isValid;} public void setValid(boolean valid) {isValid = valid;} public long getWriteVersion() {return writeVersion;} public void setWriteVersion(long writeVersion) {this.writeVersion = writeVersion;} // Disruptor预初始化对象用的工厂 public static final EventFactory<DataEvent> FACTORY = () -> { DataEvent event = new DataEvent(); event.setValid(false); event.setWriteVersion(-1); return event; }; } // 基于Disruptor实现的无锁数据存储 class DisruptorDataStore { private final RingBuffer<DataEvent> ringBuffer; // 可读最大序列:写入完成后更新,读线程最多读到这个位置 private final Sequence readableSeq = new Sequence(Sequencer.INITIAL_CURSOR_VALUE); // 已清理最小序列:删除逻辑推进后更新,标记之前的槽位可以被覆写 private final Sequence cleanedSeq = new Sequence(Sequencer.INITIAL_CURSOR_VALUE); private final int bufferSize; public DisruptorDataStore(int bufferSize) { if (Integer.bitCount(bufferSize) != 1) { throw new IllegalArgumentException("RingBuffer容量必须是2的整数次幂"); } this.bufferSize = bufferSize; // 单生产者模式性能最优,多线程写入替换为createMultiProducer即可 this.ringBuffer = RingBuffer.createSingleProducer( DataEvent.FACTORY, bufferSize, new YieldingWaitStrategy() ); // 配置门控序列,防止写入覆盖还没处理的有效数据 ringBuffer.addGatingSequences(cleanedSeq, readableSeq); } // 写入方法,对应原写入线程逻辑 public void write(String data) { long seq = ringBuffer.next(); try { DataEvent event = ringBuffer.get(seq); event.setBizData(data); event.setValid(true); event.setWriteVersion(seq); } finally { ringBuffer.publish(seq); readableSeq.set(seq); } } // 按需读取方法,对应原读取线程逻辑 public List<String> read(Predicate<String> filter) { List<String> result = new ArrayList<>(); long maxReadable = readableSeq.get(); long minReadable = Math.max(cleanedSeq.get() + 1, maxReadable - bufferSize + 1); for (long seq = minReadable; seq <= maxReadable; seq++) { DataEvent event = ringBuffer.get(seq); // 过滤已删除、未提交的脏数据 if (event.isValid() && event.getWriteVersion() == seq && filter.test(event.getBizData())) { result.add(event.getBizData()); } } return result; } // 定时删除方法,对应原定时器删除逻辑 public void deleteIf(Predicate<String> deleteCondition) { long maxReadable = readableSeq.get(); long minReadable = Math.max(cleanedSeq.get() + 1, maxReadable - bufferSize + 1); long lastCleaned = minReadable - 1; boolean meetValid = false; for (long seq = minReadable; seq <= maxReadable; seq++) { DataEvent event = ringBuffer.get(seq); if (event.isValid() && deleteCondition.test(event.getBizData())) { event.setValid(false); } // 只有从起始位置开始连续的无效条目,才能推进清理序列允许覆写 if (!meetValid && !event.isValid()) { lastCleaned = seq; } else { meetValid = true; } } cleanedSeq.set(lastCleaned); } } // 使用示例 public class DisruptorDemo { public static void main(String[] args) { DisruptorDataStore store = new DisruptorDataStore(65536); // 模拟写入线程 new Thread(() -> { for (int i = 0; i < Integer.MAX_VALUE; i++) { store.write("order-" + i); } }).start(); // 模拟读取线程 new Thread(() -> { while (true) { List<String> matched = store.read(s -> s.endsWith("666")); // 业务处理读取结果 } }).start(); // 模拟定时删除任务 ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); scheduler.scheduleAtFixedRate(() -> { // 示例删除条件:可替换为存入时间超过阈值等业务规则 store.deleteIf(s -> true); }, 10, 10, TimeUnit.SECONDS); } }
额外注意事项
如果你的业务需要无界存储所有历史数据、不允许旧数据被覆写,不要用Disruptor替换ArrayList,RingBuffer的环形设计天生适合有界缓存、流处理场景,不适合无界持久化存储。
内容的提问来源于stack exchange,提问作者Mixalis Navridis
相关产品推荐
相关产品推荐

