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

如何在以下场景中避免ConcurrentModificationException且不影响性能

高效解决并发修改异常的方案

问题根源

原代码中使用的ArrayList不是线程安全集合,多个线程同时执行addAll()、size()判断、clear()操作时,会触发ConcurrentModificationException;同时count的批量提交逻辑存在竞态条件,可能导致重复提交或数据丢失。

方案一:ThreadLocal+局部批量处理(低锁竞争)

让每个线程维护自己的本地列表,仅在需要全局批量提交时加锁,大幅减少线程间的竞争,性能损耗极低。

修改后的代码示例:

public void processFiles(List<FileIdBothDirectoryInformation> files) {
    AtomicInteger count = new AtomicInteger(0);
    ExecutorService executorService = Executors.newFixedThreadPool(30);
    // 全局批量列表,仅在提交/持久化时加锁
    List<Xxx> globalBatchList = Collections.synchronizedList(new ArrayList<>());
    int totalFiles = files.size();

    List<CompletableFuture<Void>> completableFutures = files.stream().map(file -> CompletableFuture.runAsync(() -> {
        // 每个线程绑定独立的本地列表,避免跨线程竞争
        ThreadLocal<List<Xxx>> localDslList = ThreadLocal.withInitial(ArrayList::new);
        try {
            File localFile = nasDataLoader.downloadFile(file.getFilePath());
            if(localFile != null) {
                List<?> lines = parse(localFile);
                List<Xxx> parsedLines = lines.stream().map(line -> {
                    apply(line);
                    return new Xxx(line);
                }).collect(Collectors.toList());
                localDslList.get().addAll(parsedLines);

                // 本地列表达到批量阈值时,拆分提交到全局列表
                while (localDslList.get().size() >= BATCH_SIZE) {
                    List<Xxx> batch = new ArrayList<>(localDslList.get().subList(0, BATCH_SIZE));
                    localDslList.get().subList(0, BATCH_SIZE).clear();
                    // 仅在全局批量提交时加锁,保证持久化和清空的原子性
                    synchronized (globalBatchList) {
                        globalBatchList.addAll(batch);
                        if (globalBatchList.size() >= BATCH_SIZE) {
                            persistenceManager.persist(globalBatchList);
                            globalBatchList.clear();
                        }
                    }
                }
            }
        } catch (Exception e) {
            // 按需添加异常处理逻辑
        } finally {
            // 最后一个任务执行时,提交本地剩余数据并清理ThreadLocal
            if (count.incrementAndGet() == totalFiles) {
                synchronized (globalBatchList) {
                    globalBatchList.addAll(localDslList.get());
                    if (!globalBatchList.isEmpty()) {
                        persistenceManager.persist(globalBatchList);
                        globalBatchList.clear();
                    }
                }
                localDslList.remove();
            }
        }
    }, executorService)).collect(Collectors.toList());

    completableFutures.forEach(CompletableFuture::join);
    executorService.shutdown();
}

优势:

  • 线程本地列表操作无锁,仅全局批量提交时短暂加锁,性能影响极小
  • 避免了跨线程对同一集合的并发修改,彻底解决异常问题

方案二:无锁并发队列(高并发场景最优)

使用ConcurrentLinkedQueue(无锁线程安全队列)代替ArrayList,利用其内置的线程安全操作实现批量处理,全程无显式锁,性能最优。

修改后的代码示例:

public void processFiles(List<FileIdBothDirectoryInformation> files) {
    AtomicInteger count = new AtomicInteger(0);
    ExecutorService executorService = Executors.newFixedThreadPool(30);
    ConcurrentLinkedQueue<Xxx> dslQueue = new ConcurrentLinkedQueue<>();
    int totalFiles = files.size();

    List<CompletableFuture<Void>> completableFutures = files.stream().map(file -> CompletableFuture.runAsync(() -> {
        try {
            File localFile = nasDataLoader.downloadFile(file.getFilePath());
            if(localFile != null) {
                List<?> lines = parse(localFile);
                lines.stream().map(line -> {
                    apply(line);
                    return new Xxx(line);
                }).forEach(dslQueue::offer);

                // 尝试批量处理队列中的数据
                processBatch(dslQueue);
            }
        } catch (Exception e) {
            // 按需添加异常处理逻辑
        } finally {
            // 最后一个任务执行时,处理队列中剩余的所有数据
            if (count.incrementAndGet() == totalFiles) {
                processRemainingBatch(dslQueue);
            }
        }
    }, executorService)).collect(Collectors.toList());

    completableFutures.forEach(CompletableFuture::join);
    executorService.shutdown();
}

// 批量处理队列数据,无锁操作
private void processBatch(ConcurrentLinkedQueue<Xxx> dslQueue) {
    List<Xxx> batch = new ArrayList<>(BATCH_SIZE);
    // 队列元素达到批量阈值时,循环取出处理
    while (dslQueue.size() >= BATCH_SIZE) {
        for (int i = 0; i < BATCH_SIZE; i++) {
            Xxx elem = dslQueue.poll();
            if (elem == null) break;
            batch.add(elem);
        }
        if (!batch.isEmpty()) {
            persistenceManager.persist(batch);
            batch.clear();
        }
    }
}

// 处理队列中剩余的所有数据
private void processRemainingBatch(ConcurrentLinkedQueue<Xxx> dslQueue) {
    List<Xxx> batch = new ArrayList<>();
    Xxx elem;
    while ((elem = dslQueue.poll()) != null) {
        batch.add(elem);
    }
    if (!batch.isEmpty()) {
        persistenceManager.persist(batch);
    }
}

优势:

  • ConcurrentLinkedQueue基于CAS实现无锁线程安全,高并发下性能远超同步集合
  • 无需手动加锁,代码逻辑更简洁,避免死锁风险

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 11:45:25