如何使用ForkJoinPool优化50万元素列表校验并返回两个结果数组
50万条数据校验场景的ForkJoinPool优化方案
先做单线程基础优化,这步做完通常就能砍掉70%以上的耗时,再上并行才能拿到最大收益。
前置单线程优化
- 正则匹配是校验逻辑里的性能大头,首先确认
SKIP_PATTERN是全局预编译的静态常量,不要在循环里每次调用Pattern.compile()编译正则。其次Matcher对象不是线程安全的,单线程下可以复用同一个Matcher,每次调用reset(目标字符串)重置即可,避免循环里频繁创建对象:// 全局预编译正则,不要写在循环里 private static final Pattern SKIP_PATTERN = Pattern.compile("你的邮箱匹配正则"); // 多线程场景下用ThreadLocal给每个线程分配独立Matcher private static final ThreadLocal<Matcher> EMAIL_MATCHER = ThreadLocal.withInitial(() -> SKIP_PATTERN.matcher("")); - ArrayList默认初始容量只有10,插入50万条数据会触发十几次扩容+数组拷贝,直接给初始化容量就能省掉这部分开销:
// 按预估数据量设置初始容量,坏数据量按实际坏件率预估即可 List<String[]> list = new ArrayList<>(500000); List<NotValidRow> badList = new ArrayList<>(10000); - 排查正则本身的性能问题:如果你的正则写了大量回溯逻辑(比如嵌套
.*、不必要的捕获组、范围过大的通配),会导致单条匹配耗时飙升,这种情况先优化正则写法,否则上并行也拿不到预期收益。 - 注意你现有逻辑里
tmp.length!=2的行直接被跳过,不会进入合法/坏数据列表,如果业务上需要记录长度不合法的行,要补对应逻辑。
ForkJoinPool并行实现
你的校验逻辑是纯CPU密集型任务,单条数据处理完全独立,没有IO阻塞,非常适合用ForkJoinPool的工作窃取机制做并行加速。注意不要直接用JDK全局共享的ForkJoinPool.commonPool(),公共池会被所有并行流、其他ForkJoin任务共用,容易被阻塞任务影响性能,建议单独创建自定义线程池,线程数和CPU核心数持平即可(CPU密集型任务不需要太多线程,过多线程反而会因为上下文切换降低性能)。
方式1:手写RecursiveAction(性能最优,无锁)
核心思路是把大的数据集拆成多个小分片,每个分片大小小于阈值就直接串行处理,每个子任务独立维护自己的合法/坏数据列表,最后合并所有子任务的结果,全程没有锁竞争,性能最高。
// 自定义独立的校验线程池,线程数等于CPU核心数 private static final ForkJoinPool VALIDATE_POOL = new ForkJoinPool( Runtime.getRuntime().availableProcessors(), ForkJoinPool.defaultForkJoinWorkerThreadFactory, null, false ); // 分片校验任务 private static class ValidateTask extends RecursiveAction { // 分片阈值,单分片2000-5000条比较合适,太小会增加任务调度开销,太大无法充分利用多核 private static final int THRESHOLD = 2000; private final List<String[]> source; private final int start; private final int end; // 每个子任务内部维护结果,不共享 private List<String[]> validList; private List<NotValidRow> invalidList; public ValidateTask(List<String[]> source, int start, int end) { this.source = source; this.start = start; this.end = end; } @Override protected void compute() { // 分片大小达标,直接处理当前分片 if (end - start <= THRESHOLD) { validList = new ArrayList<>(end - start); invalidList = new ArrayList<>(16); // 当前任务线程独立持有Matcher,无线程安全问题 Matcher matcher = SKIP_PATTERN.matcher(""); for (int i = start; i < end; i++) { String[] tmp = source.get(i); if (tmp.length != 2) continue; if (tmp[0] == null || !matcher.reset(tmp[0]).matches()) { invalidList.add(new NotValidRow(tmp[0], tmp[1], NotValidRowReason.NOT_VALID_EMAIL)); } if (tmp[1] == null || tmp[1].isBlank()) { invalidList.add(new NotValidRow(tmp[0], tmp[1], NotValidRowReason.EMPTY_NAME)); } validList.add(tmp); } return; } // 分片过大,拆成两个子任务 int mid = (start + end) >>> 1; ValidateTask left = new ValidateTask(source, start, mid); ValidateTask right = new ValidateTask(source, mid, end); invokeAll(left, right); // 合并子任务结果 validList = new ArrayList<>(left.validList.size() + right.validList.size()); validList.addAll(left.validList); validList.addAll(right.validList); invalidList = new ArrayList<>(left.invalidList.size() + right.invalidList.size()); invalidList.addAll(left.invalidList); invalidList.addAll(right.invalidList); } }
调用逻辑:
// 第一步:单线程读完输入流所有数据,IO操作不要放到并行任务里 List<String[]> sourceList = new ArrayList<>(500000); Iterator<String[]> iterator = parser.iterate(request.getInputStream()).iterator(); while (iterator.hasNext()) { sourceList.add(iterator.next()); } // 第二步:提交并行任务等待执行完成 ValidateTask rootTask = new ValidateTask(sourceList, 0, sourceList.size()); VALIDATE_POOL.invoke(rootTask); // 第三步:拿最终结果 List<String[]> validResult = rootTask.validList; List<NotValidRow> invalidResult = rootTask.invalidList;
方式2:并行流实现(代码简单,够用)
如果不想手写任务拆分逻辑,也可以用JDK并行流实现,底层同样基于ForkJoinPool,注意提交到自定义线程池即可:
List<String[]> sourceList = new ArrayList<>(500000); // 省略读流逻辑,和上面一致 List<String[]> validResult = Collections.synchronizedList(new ArrayList<>(sourceList.size())); List<NotValidRow> invalidResult = Collections.synchronizedList(new ArrayList<>(10000)); // 提交到自定义线程池执行 VALIDATE_POOL.submit(() -> sourceList.parallelStream().forEach(tmp -> { if (tmp.length != 2) return; Matcher matcher = EMAIL_MATCHER.get(); if (tmp[0] == null || !matcher.reset(tmp[0]).matches()) { invalidResult.add(new NotValidRow(tmp[0], tmp[1], NotValidRowReason.NOT_VALID_EMAIL)); } if (tmp[1] == null || tmp[1].isBlank()) { invalidResult.add(new NotValidRow(tmp[0], tmp[1], NotValidRowReason.EMPTY_NAME)); } validResult.add(tmp); }) ).join();
性能预期与注意事项
- 8核CPU机器上,做完基础优化+并行处理后,整体耗时可以稳定在3-5秒,比原串行逻辑快8-10倍。
- 不要在子任务里共享非线程安全的对象(比如ArrayList、HashMap、Matcher),否则会出现脏数据甚至并发异常。
- 分片阈值不要设太小,否则会生成过多子任务,调度开销会吃掉并行收益。
- 如果读流解析的耗时占比很高,优先换高性能CSV解析库,低性能解析器的IO+对象创建开销可能占总耗时的一半以上。
内容的提问来源于stack exchange,提问作者JD_UA
相关产品推荐
相关产品推荐

