You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

ReentrantReadWriteLock逻辑迁移至带屏障LMAX Disruptor的方法

方案可行性结论

完全可以迁移到LMAX Disruptor实现,你初步的思路大方向是对的,只要调整几个设计细节,在高并发场景下确实能获得比ReentrantReadWriteLock+ArrayList方案更好的性能:Disruptor本身就是为消除锁竞争、伪共享、不必要内存屏障设计的,核心路径全程无锁,多线程场景下的延迟稳定性、吞吐量都比显式锁方案高一个量级。

原有思路的正误校验
  • 正确的部分:用RingBuffer替代ArrayList做核心存储、拆分生产/消费角色、用序列屏障做流程控制的思路完全符合Disruptor的设计范式,避开了从零实现无锁结构的常见坑。
  • 需要调整的误区:
    1. 别把删除操作做成独立生产者:RingBuffer是固定容量的环形顺序存储结构,不支持像ArrayList那样随机删除中间条目后移动元素缩容,删除逻辑应该下沉到维护角色做逻辑标记,硬要在生产者侧做随机删除会彻底破坏RingBuffer的顺序写入特性,反而引入额外开销。
    2. 别把读取做成普通事件消费者:你的场景是按需读取,不是事件到达就被动消费,不需要给读线程分配独立消费序列,否则会和写入、删除流程争抢序列权限,反而拖慢性能。
    3. 没必要给写入、删除配置双生产者:删除不属于事件写入流的一部分,双生产者模式本身会引入额外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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 06:30:44