如何在以下场景中避免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
相关产品推荐
相关产品推荐

