先说我前阵子接到的一个需求吧要对一份百万行的账单文件做多维度聚合统计单线程跑一次要四十多秒批处理任务每天要跑好几轮明显是块硬骨头。我第一反应就是分治——把文件按行号均匀切块核心越多切得越多每块独立统计最后再合并结果。用 ExecutorService 写出来后代码又长又绕主线程要反复提交子任务通过 Future 一个个 get() 取结果任务层级稍微深一点线程们就互相干等并行效率惨不忍睹。后来换成 ForkJoin 重写结构一下就清楚了很多这也是我想写这篇文章的原因。ForkJoin 从 JDK 7 就有了它是 Java 里为“分而治之”这类递归并行任务专门设计的基础设施。很多人知道这个框架的名字也见过 RecursiveTask 的示例代码但真到项目里要用往往说不清拆分阈值怎么定、工作窃取到底怎么工作、为什么任务里不能做阻塞 IO以及它和 Java 8 的 parallelStream 到底什么关系。这篇我会把这些东西按我实际使用的顺序串一遍从原理到代码从调参到踩坑尽量说到看完就能用、能判断哪些场景该上车。如果你正在写并发相关的代码或者在准备 Java 面试想把这个话题讲得比别人有深度这篇文章应该对你有用。1. ForkJoin 解决的核心痛点递归任务为什么不能直接丢给普通线程池1.1 传统线程池擅长“一次性任务”却不擅长“递归任务”我们先拆一下问题。Executors.newFixedThreadPool(8)加上 Callable适合处理什么适合一批彼此独立、没有父子关系的任务。比如同时下载八个文件、同时调八个接口谁先完成都不影响别人主线程只管把结果都收集齐。但分治算法在结构上是一棵树。父任务必须等子任务的结果子任务可能要再拆出孙任务每个节点之间天然存在依赖。你把这种任务用普通线程池跑通常的做法是每个任务内部再往线程池里提交子任务然后用 Future.get() 等待。问题在于get() 是阻塞的线程会等在结果上而不是转头去执行其他任务。于是任务层级一深线程池几乎所有线程都阻塞在 get() 上队列里堆着一大堆没人执行的子任务CPU 其实很闲。我之前做过一个测试用普通线程池跑一个一千万数据的归并排序每个归并节点都提交到池里再 get 等结果最终 8 核机器上只跑出了单核 1.3 倍的速度排完一次花了快 5 秒。那不是并行是给自己加调度负担。1.2 Fork-Join 模型的核心动作只有三个fork、join、computeForkJoin 的模型非常干净。每个任务继承RecursiveTaskV或RecursiveAction主要逻辑写在compute()里。compute 里要做两件事如果任务足够小直接计算并返回结果如果任务还不够小拆成两个子任务调用fork()把其中一个放入当前工作线程的任务队列另一个自己继续用compute()处理最后调用join()等被 fork 出去的子任务完成合并两个结果。这里最关键的优化是拆出来的两个子任务一个 fork 出去另一个在当前线程“顺手”就做了。这样就不会像普通线程池那样白白阻塞等待当前线程始终有活儿干。后面讲工作窃取的时候你会看到这个设计直接决定了整个框架的效率上限。1.3 工作窃取让“有人闲着”变成“谁来干活都一样”想象一个厨房流水线四个厨师各自站在自己的操作台前切菜每个人手边都有一个备料筐切好的料放在自己这筐里。干得快的人发现自己的筐空了不会傻站着而是跑到别人操作台后面把人家的备料筐拿过来接着切。这就是 ForkJoin 的工作窃取Work Stealing思想。在 ForkJoinPool 内部每个工作线程都有一个自己的双端队列存放被 fork 出来的子任务。正常工作状态下线程从自己的队列头部取出任务执行后 fork 的先执行类似 LIFO这样能最大化 CPU 缓存的局部性——刚创建的子任务数据很可能还在缓存里。而当一个线程的队列空了它不会闲着而是去偷其他线程队列尾部的任务来执行尾部是别人最老的任务这样偷取时和原线程的竞争最少。这也是 ForkJoin 在处理递归任务时能大幅跑赢普通线程池的根本原因普通线程池是队列顶着一个共享锁所有线程抢同一个队列ForkJoin 是“人人有队、余力互帮”把竞争最小化把空转时间压到最低。2. 核心 API 拆解一步一步写出第一个 RecursiveTask 程序2.1 四个核心类其实各管一件事ForkJoin 框架涉及的类不多真正要搞明白的就四个类名职责ForkJoinPool真正干活的线程池。负责管理工作线程、任务队列和窃取逻辑ForkJoinTaskV任务抽象基类。fork、join、invoke 等方法都定义在这里RecursiveTaskVForkJoinTask 的子类用于有返回值的分治任务RecursiveActionForkJoinTask 的子类用于没有返回值的分治任务比如排序、遍历还有一个ForkJoinWorkerThread是池里的工作线程一般我们不需要直接碰它。框架的核心思想就体现在这套类设计上RecursiveTask 是“compute 方法要返回结果”RecursiveAction 是“compute 方法只做事、不返回”两种场景基本覆盖了你能想到的所有分治计算。2.2 完整的求和示例超过阈值就算没超过就拆直接看代码。下面是一个典型的大数组求和的 RecursiveTaskpublic class SumTask extends RecursiveTaskLong { private static final int THRESHOLD 10_000; private final long[] array; private final int start; private final int end; public SumTask(long[] array, int start, int end) { this.array array; this.start start; this.end end; } Override protected Long compute() { // 如果区间足够小直接循环累加 if (end - start THRESHOLD) { long sum 0; for (int i start; i end; i) { sum array[i]; } return sum; } // 否则拆成两半左边 fork 到队列右边自己算 int mid (start end) 1; SumTask left new SumTask(array, start, mid); SumTask right new SumTask(array, mid, end); left.fork(); long rightSum right.compute(); long leftSum left.join(); return leftSum rightSum; } }调用方式也非常简单long[] data new long[10_000_000]; // 初始化 data ... ForkJoinPool pool new ForkJoinPool(); long sum pool.invoke(new SumTask(data, 0, data.length)); pool.shutdown();说两个初学者容易忽略的点。第一为什么是left.fork()然后right.compute()而不是两个都 fork ?两个都 fork 也不是不行但左侧任务已经提交到队列里了当前线程如果把右侧任务也 fork 出去自己就闲下来了还要等别人来窃取凭空多了一次任务入队和出队性能没有好处。正确姿势是本线程“顺手”处理一个子任务再 join 另一个这样当前线程利用率最高。第二join 和 get 的等待语义不一样。join 在任务没有完成时会阻塞并且如果 compute 抛出了异常join 会把异常包装后再次抛出get()会把任务执行的异常包装成ExecutionException。两种方式触发的异常类型不同这也是第 4 章会展开讲的重点之一。2.3 execute、invoke、submit 三个入口什么时候用哪个ForkJoinPool 对外提供三个主要入口很多人记混先看一张表方法是否等待结果返回类型异常行为适合场景execute(task)不等待void任务异常存储在 task 中需自行检查触发后不管比如异步后台统计invoke(task)等待并返回结果任务本身类型异常直接抛出同步调用主流程需要结果submit(task)可等待ForkJoinTask可通过 Future.get() 获取异常想要 Future但又不完全同步我说下实际感受项目里 90% 的同步计算场景直接用pool.invoke(task)最省心。它内部等价于submit(task)加join()但代码短异常处理也直观。只有当你需要把任务丢到后台、过一段时间再取结果的时候才用submit。需要注意的是execute不等于“没有异常”。任务执行异常后不会主动打印任何日志你如果不去调task.isCompletedAbnormally()异常会被悄悄吞掉。这个行为在线上会坑人我自己就遇到过一次报表数据突然少了一段排查半天才发现是某个 RecursiveAction 在计算某个中间区块时抛了空指针但因为用了 execute异常完全没暴露。3. 用 ForkJoin 写归并排序阈值和并行度的调参实录3.1 阈值不能拍脑袋定也不能直接照抄网上的回到我自己做账单聚合时的另一个痛点阈值 THRESHOLD 到底填多少网上示例经常写 1000 或者 10000但这不是万能公式。阈值本质上是“任务小到不值得再拆”的一个判断线它跟两个因素相关单任务的最小操作开销如果单个任务只做 10 次加法拆到这份上光 fork/join 的调度成本就比计算成本还高反而更慢数组长度和并行度假设数组一千万并行度 8每个线程理想上要处理一百二十五万个元素。如果阈值设成 10 万那任务总数大约会有 100 个任务粒度太粗可能出现某线程处理完了其他线程还有很多活窃取机制来不及平衡。如果设成 100任务总数几千个窃取可以充分发生但 fork/join 次数也多栈帧压力大。我给一个经过测试的经验参考区间CPU 密集的分治任务每个子任务至少承担 1 万到 10 万次基本操作。对数组求和这种超轻量操作阈值设在 1000050000 比较稳对归并排序这种需要额外复制数组的操作可以放宽到 50000。一句话总结阈值要大到让单任务做“有意义的工作”而不是让框架忙着拆任务。3.2 归并排序完整实现与性能对比归并排序没有返回值所以继承RecursiveAction更合适。下面是我测试用的版本注意我特意复用了临时数组避免每次 merge 都重新分配否则内存分配开销会吃掉并行收益public class MergeSortTask extends RecursiveAction { private static final int THRESHOLD 50_000; private final int[] arr; private final int[] tmp; private final int left; private final int right; public MergeSortTask(int[] arr, int[] tmp, int left, int right) { this.arr arr; this.tmp tmp; this.left left; this.right right; } Override protected void compute() { if (right - left THRESHOLD) { Arrays.sort(arr, left, right 1); return; } int mid (left right) 1; MergeSortTask leftTask new MergeSortTask(arr, tmp, left, mid); MergeSortTask rightTask new MergeSortTask(arr, tmp, mid 1, right); invokeAll(leftTask, rightTask); merge(arr, tmp, left, mid, right); } private void merge(int[] arr, int[] tmp, int left, int mid, int right) { System.arraycopy(arr, left, tmp, left, right - left 1); int i left, j mid 1, k left; while (i mid j right) { if (tmp[i] tmp[j]) { arr[k] tmp[i]; } else { arr[k] tmp[j]; } } while (i mid) { arr[k] tmp[i]; } while (j right) { arr[k] tmp[j]; } } }invokeAll(leftTask, rightTask)是框架提供的好东西它相当于把两个任务都提交并等待全部完成内部会处理“谁先谁后”的优化比手动 fork/join 更省心。在我本地环境JDK 178 核 i7下对一千万随机整数排序跑了三次取中位数数据大概是这样方案耗时说明单线程手写归并排序约 3100 ms作为基准ForkJoin 归并阈值 50000并行度 8约 1050 ms约 3 倍加速ForkJoin 归并阈值 5000并行度 8约 1250 ms任务拆分过细调度开销明显ForkJoin 归并阈值 50000并行度 4约 1650 ms并行度低加速有限数据本身并不重要不同机器差异很大但趋势是稳定的阈值过小会引入额外调度开销阈值太大则退化成近似串行。建议拿到自己机器上测一遍找到曲线拐点。3.3 parallelism 和 commonPool别在构造函数上随意调参ForkJoinPool 默认的并行度等于 CPU 可用核心数看起来合理但有两个小坑。第一CPU 密集任务里并行度设为核数就够了设太大反而变慢。我们试过把并行度设成 16在 8 核机器上跑归并排序结果和并行度 8 差不多有时候还慢一点点。因为线程上下文切换、队列竞争的边际成本超过了额外线程带来的收益。如果有超线程结论仍然类似——线程数等于物理核数通常是甜点。第二commonPool 默认并行度是 CPU 核数减 1。它给主线程留了一个位置避免主线程参与任务提交时直接被饿死。这个池是 JVM 全局共享的所有没指定池的 ForkJoinTask、包括 parallelStream默认都会跑在它上面。你要是想给 commonPool 调并行度得靠系统属性-Djava.util.concurrent.ForkJoinPool.common.parallelism4但我不建议全局乱调因为会影响 JVM 里所有使用并行流的代码。更稳妥的做法是重要业务单独创建专用 ForkJoinPool把并行度控制在你自己能掌控的范围。4. 工作窃取的底层机制双端队列和窃取方向为什么这样设计4.1 双端队列里的任务流转ForkJoinPool 里每个工作线程都持有自己的双端队列这个队列的操作规则很特殊正常执行时线程从队首取出任务执行自己fork()出的任务会 push 到队首当一个线程完成了自己的任务它去其他线程的队列尾部窃取一个任务执行。队首出队、队尾偷取这个方向设计不是随手定的。自己新 fork 的任务放在队首、优先从队首取意味着“年龄最新”的任务会被最先执行而这类任务的数据通常刚刚创建最有可能还在 CPU 缓存里执行效率高。反过来偷取的老任务在队尾离对方正在执行的队首最远偷的时候不容易和对方抢锁。如果偷取也从队首偷那两队线程会频繁竞争同一个队列头部性能直接下滑。4.2 实打实的坑在 ForkJoinTask 里做阻塞 IO线程池直接被拖死这是我在日志处理服务里遇到的最深一次坑。当时为了批量压缩一批日志文件我写了个 RecursiveAction每个子任务里直接调用了上传组件的同步接口上传文件。上传是网络 IO慢的能到几百毫秒。8 个 worker 全部阻塞在上传接口上而且上传接口底层还用的是同一个 ForkJoinPool 去等响应结果整个池的线程都在互相等没有线程可以去窃取任务任务队列越堆越长最终调用方超时。你可能会问工作窃取不是能把任务分给别人吗关键就在这里——窃取只能窃取“排队中的任务”不能窃取“正在执行且阻塞的任务”。当所有线程都阻塞在 IO 上时窃取机制完全失效。ForkJoin 的设计目标是 CPU 密集计算它不适合代替 IO 线程池。任何会阻塞线程的远程调用、磁盘 IO、锁等待都不应该出现在 compute() 里。后来我改成用主线程拆好文件块每个块交给普通 IO 线程池上传上传结果用 CompletableFuture 汇总。代码结构没复杂多少但吞吐提升了不止一倍。4.3 异常处理与递归深度两个容易被忽略的细节第一个细节是异常吞没。如果你用 execute() 提交任务任务里的异常不会自动打印用 invoke 或 submit 也一样要小心join()抛出的异常类型和get()不同前者还原成 RuntimeException 同样式后者包成 ExecutionException。我强烈建议在compute()的最外层加一层兜底日志把抛错位置打出来否则线上排查会非常痛苦。第二个细节是栈溢出。ForkJoin 并不是把所有子任务都 fork 出去并行执行而是靠 compute() 方法递归进入下一层。也就是说每个“用来继续拆分的节点”都会消耗一层 Java 调用栈。阈值设得太小递归深度可能几千上万层StackOverflowError 说来就来。如果你在日志里看到栈溢出第一件事不是调大 -Xss而是调大阈值让递归深度降下来。5. parallelStream 和 ForkJoin 的关系什么时候别用并行流5.1 parallelStream 的底层就是 commonPoolJava 8 的parallelStream()出来之后很多代码写起来确实爽。但它的底层就是 ForkJoinPool 的 commonPool。也就是说它共享了 JVM 里默认的那个池池里的线程数默认是 CPU 核数减 1。这意味着两件事第一如果你有好几处代码都在用 parallelStream它们之间会互相抢线程第二如果你在 parallelStream 的 lambda 里做了阻塞操作会连累 JVM 里所有其他并行流、以及其他默认跑在 commonPool 上的 ForkJoinTask。5.2 误用现场一个报表功能把整个服务的并行流带崩了我的一个业务系统里有个报表接口要对上千个分片文件做解析和汇总代码写了 parallelStream。另一个接口的某个统计功能也用 parallelStream 做计算平时都正常。后来报表接口的分片文件数量翻倍单次解析里还加了正则和字段校验跑一次要十几秒。结果报表接口每次一跑commonPool 的线程全被它占住另一个统计接口的并行流虽然代码没改却慢了好几倍。排查的过程很简单打开 jstack发现 commonPool-2 线程全阻塞在报表解析栈上。解决方式也不复杂把报表处理改成跑在自定义的 ForkJoinPool 里限制并行度同时把统计接口换成单线程流sequential因为它的数据量根本不需要并行。5.3 手动 ForkJoin 比并行流强在哪什么时候必须手动parallelStream 胜在代码量少但它帮你做的决策也多很多事你控制不了并行度由 commonPool 决定你不能只给某一段代码单独设线程数任务拆分粒度不可控框架内部是逐个元素拆分无法表达复杂的树形依赖和递归结构异常处理和 F/J 任务不太一样排查成本更高。所以我的选择标准很直接简单 map 操作、数据量大、CPU 密集、无副作用 用 parallelStream有父子依赖、需要递归拆分、要单独控并行度、任务里包含复杂的数据分区逻辑 手动 ForkJoin。后者的代码虽然长一点但每个边界条件你都看得见、测得了。6. 真实项目里的应用场景分治并行不是万金油6.1 什么场景真正值得用 ForkJoin我筛选过自己项目里能用上 ForkJoin 的场景基本满足这几个特征问题可以拆成互不相关的子问题、子问题计算量足够大、任务是 CPU 密集的、没有共享可变状态。典型代表大批量数组求和、统计、极值查找归并排序、快速排序、二分查找变体等算法密集型任务树形结构遍历比如目录占用大小统计、组织架构层级汇总蒙特卡洛模拟、矩阵分块运算这类需要大量独立计算的数值问题轻量版 MapReduce把一批数据拆到多核上做统计再归并。6.2 什么场景我明确不建议用反过来有几种场景我几乎不会选它。比如 IO 密集型任务包括文件读取、远程接口调用、数据库操作再比如任务本身非常轻量单个任务只有几微秒拆分的调度成本远大于计算本身还有任务之间存在共享可变状态的情况比如多个子任务要同时往一个 Map 里写数据锁竞争会直接把并行收益吃掉。有个例子我在写树形目录统计时最初每个子任务都往一个 ConcurrentHashMap 里累加文件大小结果 8 核跑下来比单线程还慢。改成每个子任务返回局部结果最后在父任务里合并速度才正常。这说明 ForkJoin 本身再快也快不过有缺陷的算法结构。6.3 和 CompletableFuture 配合CPU 并行和 IO 异步各司其职在真实业务系统里任务往往既有 CPU 计算又有 IO 等待。我的惯用组合是用 ForkJoin 处理纯计算部分用 CompletableFuture 做 IO 异步编排。比如一个统计报表先把几 GB 日志里的关键字段用 ForkJoin 并行解析成结构化数据然后每个分片文件的落库用 CompletableFuture 提交给独立线程池。两个框架各跑各的池互不干扰想继续扩展的时候也方便加超时、重试和回调。但注意别在一个 ForkJoin 任务里直接调 CompletableFuture 的阻塞 join否则又回到“线程被阻塞”的老问题。跨框架编排时最好在 compute() 里只返回“需要异步执行的子任务清单”让外层来处理后续编排。7. 写给后来者的几条实战经验监控、测试和设计取舍7.1 用 ForkJoinPool 的自带指标观察窃取效率很多人不知道 ForkJoinPool 自带监控方法都是上线后等报警。我常用的几个pool.getStealCount()返回被偷走执行的任务数量这个值如果长期为 0说明任务拆分粒度太粗没产生足够的窃取机会pool.getQueuedTaskCount()返回当前队列里的任务数如果非常大而又不下降说明线程都在忙别的事典型就是阻塞了pool.getActiveThreadCount()返回正在工作的线程数结合并行度看能判断线程是否大量空闲。我在性能调优时会在关键流程里把这些数字记到日志里和耗时指标一起看。有一次发现 stealCount 很高但整体耗时没下降就知道任务是偷了但每个任务太小、偷取的调度开销吃掉了收益于是把阈值调大问题就解决了。7.2 测试并行代码的要点别用默认并行度测测试 ForkJoin 最容易犯的错是拿自己本机默认并行度跑一次就下结论。多核 CPU 的波动很大我建议固定池的并行度比如 4 和 8 各测一轮做性能对比时先预热让 JIT 编译完成再跑多次取中位数不要取最大值测试环境尽量和生产环境核数一致单测先验证结果正确性和可重复性再谈性能。ForkJoin 的并行计算如果逻辑有误结果很可能是“大多数时候对偶尔不对”这种状态最难查务必在 compute 前后打印可对比的边界值。7.3 设计取舍上的几点体会最后说点框架之外的感受。ForkJoin 的 API 本身不难难的是任务划分。我在项目里真正受益的场景几乎都是“问题本身有清晰的分治结构”才能写出漂亮的递归任务。如果一个问题拆起来勉强甚至要引入很多共享状态那就说明它不适合这个框架。另一个体会是尽量保持每个子任务的独立性任务内不要依赖外部可变变量。只要所有子任务的结果可以单纯通过返回值汇总这个并行代码几乎不会出问题。遇到想共享中间结果的冲动时先停下来想想能不能改成“父任务合并子结果”的计算结构。我记得刚把 ForkJoin 用进账单聚合模块的时候同事看着那段代码说这不就是以前学算法时的递归吗。我说对分治思想学了很多年真正让它在多核机器上跑起来的就是这个框架。理解它之后再看 parallelStream、看各种大数据框架的“分片汇总”模型背后都是同一个套路。这个问题的价值不在于记住几个 API而在于你能不能在写业务代码的时候一眼看出“这段逻辑其实可以分而治之”。
阅读完成 · 觉得有帮助?