迭代器批量处理改并行?20万条数据库数据处理优化求助
关于ArrayList的线程安全问题
是的,你必须替换当前的ArrayList<OutX> built——ArrayList是线程不安全的,多个并行任务同时调用built::add会导致数据丢失、数组越界或元素重复等并发问题。
但不推荐用Vector:Vector是通过给每个方法加synchronized实现线程安全的,锁粒度大,高并发场景下性能损耗严重。更优的方案有两种:
- 方案一(优先选择):让每个并行任务独立生成子结果列表,最后合并到主列表,完全避免并发写入竞争:
然后调整并行任务逻辑:// 修改processBatch方法,返回子结果列表 private List<OutX> processBatch(List<X> list) { return list.stream().map(m -> { X xEntry = repository.findLatestUpdateById(m.getEntryId()); if (xEntry == null) { return null; } return buildXOut(xEntry); }).filter(Objects::nonNull).collect(Collectors.toList()); }List<List<OutX>> subResults = new ArrayList<>(); for (int i = 0; i < lists.size(); i += batchSize) { int endIndex = Math.min(i + batchSize, lists.size()); List<X> batch = lists.subList(i, endIndex); ForkJoinTask<List<OutX>> task = forkJoinPool.submit(() -> processBatch(batch)); subResults.add(task.join()); } // 最后合并所有子结果 List<OutX> built = new ArrayList<>(); subResults.forEach(built::addAll); - 方案二:如果必须共享列表,使用
Collections.synchronizedList(new ArrayList<>())或CopyOnWriteArrayList。前者是全局锁实现,性能比Vector好;后者适合读多写少场景,写操作会复制整个数组,写密集场景开销大。
进一步优化建议
解决N+1数据库查询问题
当前processBatch中每个元素单独调用findLatestUpdateById,20万条数据会产生20万次数据库查询,这是核心性能瓶颈之一。改成批量查询:private List<OutX> processBatch(List<X> list) { // 收集当前批次所有entryId List<Long> entryIds = list.stream().map(X::getEntryId).collect(Collectors.toList()); // 新增批量查询方法,一次获取所有最新记录并转成Map Map<Long, X> latestEntryMap = repository.findLatestUpdatesByIds(entryIds) .stream().collect(Collectors.toMap(X::getEntryId, Function.identity())); return list.stream().map(m -> { X xEntry = latestEntryMap.get(m.getEntryId()); if (xEntry == null) { return null; } return buildXOut(xEntry); }).filter(Objects::nonNull).collect(Collectors.toList()); }这会把每批的查询次数从5000次降到1次,大幅减少数据库交互耗时。
避免一次性加载全量数据到内存
当前用findBySourceFileCreationDate一次性加载20万条数据,会占用大量堆内存,甚至触发OOM。建议回到流式查询streamBySourceFileCreationDate,边读边处理,降低内存占用。调整ForkJoinPool线程数
ForkJoinPool默认线程数等于CPU核心数,但你的任务是IO密集型(大量数据库等待),核心数线程会导致资源利用率低。可根据数据库连接池大小设置线程数,比如new ForkJoinPool(16)(若连接池最大连接数为20),避免线程过多耗尽数据库连接。优化批次大小
测试不同批次大小(如2000、5000、10000),找到最优值:批次太大单任务耗时过长,太小会增加任务调度开销。用并行流简化代码
若不需要自定义ForkJoinPool,可直接用Java并行流简化逻辑(注意先解决N+1问题):List<OutX> built = repository.streamBySourceFileCreationDate(getDateFromFilename(sourceFileName)) .parallel() .map(m -> { X xEntry = repository.findLatestUpdateById(m.getEntryId()); return xEntry != null ? buildXOut(xEntry) : null; }) .filter(Objects::nonNull) .collect(Collectors.toList());
内容的提问来源于stack exchange,提问作者cUser

