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

Java多线程下带容量限制的集合刷盘非阻塞方案优化咨询

非阻塞实现固定容量刷盘集合的解决方案

核心约束回顾

  • 集合容量严格不超过指定阈值(如25)
  • 刷盘期间禁止写入旧数据
  • 刷盘数据量准确,无超量
  • 避免全量独占锁的性能开销

问题根源分析

你之前的方案失败主要因为:

  • 未保证添加操作+容量检查的原子性,多线程并发添加时会突破容量限制
  • 全局共享缓冲区无线程安全保护,导致size统计失真
  • 同步逻辑存在竞态条件,比如flushOnGoing检查与add之间的间隙,以及锁的错误使用导致死锁

可行非阻塞方案

方案1:双缓冲区原子切换(自定义实现)

通过双缓冲区交替使用,结合原子引用和计数器控制写入,刷盘时切换缓冲区,避免写入线程等待刷盘完成。

import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;

public class NonBlockingBoundedWriter {
    private static final int MAX_CAPACITY = 2;
    // 两个交替使用的缓冲区
    private final List<Object> buffer1 = new ArrayList<>(MAX_CAPACITY);
    private final List<Object> buffer2 = new ArrayList<>(MAX_CAPACITY);
    // 当前可写入的缓冲区(原子引用保证线程安全切换)
    private final AtomicReference<List<Object>> activeBuffer = new AtomicReference<>(buffer1);
    // 当前缓冲区的元素数量(原子计数器,精确控制容量)
    private final AtomicInteger currentSize = new AtomicInteger(0);
    // 标记是否正在刷盘(避免多线程重复触发)
    private volatile boolean flushing = false;

    public void add(Object data) {
        while (true) {
            List<Object> current = activeBuffer.get();
            int size = currentSize.get();

            // 当前缓冲区已满,尝试触发刷盘
            if (size >= MAX_CAPACITY) {
                // 只有一个线程能成功触发刷盘
                if (!flushing) {
                    synchronized (this) {
                        if (!flushing) {
                            flushing = true;
                            // 原子切换到新缓冲区
                            List<Object> oldBuffer = activeBuffer.getAndSet(current == buffer1 ? buffer2 : buffer1);
                            // 重置计数器,让新缓冲区接收写入
                            currentSize.set(0);
                            // 异步刷盘(也可同步,这里不阻塞写入线程)
                            CompletableFuture.runAsync(() -> {
                                flush(oldBuffer);
                                flushing = false;
                            });
                        }
                    }
                }
                // 继续循环,等待新缓冲区可用
                continue;
            }

            // 原子增加计数器,成功则添加元素
            if (currentSize.compareAndSet(size, size + 1)) {
                // 缓冲区非线程安全,对当前缓冲区加锁(竞争极小,因为只有CAS成功的线程能进入)
                synchronized (current) {
                    current.add(data);
                }
                break;
            }
            // CAS失败,重试
        }
    }

    private void flush(List<Object> buffer) {
        System.out.printf("Thread %s flushing %d elements%n", Thread.currentThread().getName(), buffer.size());
        // 模拟刷盘操作(写入磁盘、数据库等)
        buffer.clear();
    }

    public static void main(String[] args) {
        NonBlockingBoundedWriter writer = new NonBlockingBoundedWriter();
        ExecutorService ex = Executors.newFixedThreadPool(2);
        List<CompletableFuture<String>> cfList = new ArrayList<>();

        for (int n = 0; n <= 5; n++) {
            cfList.add(CompletableFuture.supplyAsync(() -> {
                writer.add(UUID.randomUUID().toString());
                return "Done";
            }, ex));
        }

        cfList.forEach(CompletableFuture::join);
        ex.shutdown();
    }
}

方案优势:

  • 写入操作大部分时间无锁,仅CAS操作和极小概率的缓冲区同步锁
  • 刷盘时切换缓冲区,写入线程无需等待刷盘完成,并发性能高
  • 原子计数器严格保证容量不超限,刷盘数据量准确

方案2:基于ArrayBlockingQueue的非阻塞实现

利用JDK原生的ArrayBlockingQueue(固定容量、线程安全),结合非阻塞的offer方法实现写入,满队列时触发批量刷盘。

import java.util.UUID;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicBoolean;

public class QueueBasedBoundedWriter {
    private static final int MAX_CAPACITY = 2;
    // 固定容量的线程安全队列
    private final ArrayBlockingQueue<Object> queue = new ArrayBlockingQueue<>(MAX_CAPACITY);
    // 避免多线程重复触发刷盘
    private final AtomicBoolean flushing = new AtomicBoolean(false);

    public void add(Object data) {
        // 非阻塞写入,失败说明队列已满
        while (!queue.offer(data)) {
            // 尝试触发刷盘(仅一个线程成功)
            if (flushing.compareAndSet(false, true)) {
                // 批量取出所有元素,保证刷盘数据量准确
                Object[] elements = new Object[MAX_CAPACITY];
                int count = queue.drainTo(List.of(elements));
                flush(elements, count);
                // 标记刷盘完成,允许后续触发
                flushing.set(false);
            }
            // 重试写入
        }
    }

    private void flush(Object[] elements, int count) {
        System.out.printf("Thread %s flushing %d elements%n", Thread.currentThread().getName(), count);
        // 模拟刷盘操作
    }

    public static void main(String[] args) {
        QueueBasedBoundedWriter writer = new QueueBasedBoundedWriter();
        ExecutorService ex = Executors.newFixedThreadPool(2);
        List<CompletableFuture<String>> cfList = new ArrayList<>();

        for (int n = 0; n <= 5; n++) {
            cfList.add(CompletableFuture.supplyAsync(() -> {
                writer.add(UUID.randomUUID().toString());
                return "Done";
            }, ex));
        }

        cfList.forEach(CompletableFuture::join);
        ex.shutdown();
    }
}

方案优势:

  • 完全复用JDK线程安全组件,代码简洁,无需手动实现同步逻辑
  • offer方法非阻塞,刷盘时drainTo批量取数,保证数据准确性
  • 队列本身严格限制容量,无需额外计数器

方案选择建议

  • 若需要自定义缓冲区结构或更高的性能控制,选择双缓冲区方案
  • 若追求代码简洁、可靠性优先,选择ArrayBlockingQueue方案

内容的提问来源于stack exchange,提问作者user1906450

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 18:50:26