ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Java并行编程利器:ForkJoin框架原理与实践

Java并行编程利器:ForkJoin框架原理与实践 多核CPU早已不是新鲜事从四核到十六核再到服务器上动辄几十上百个逻辑核心Java开发者再靠“单线程跑到底”的思路显然吃不下这块红利。想要把一个大任务真正跑起来让多个核心同时干活自然要请并行编程出场。而Java里最贴合“分而治之”思想的并行编程工具就是ForkJoin框架。它不像手动new线程那样原始也不像ExecutorService那样只知道盲目分发任务而是专门为“把一个大事递归拆成小块、分别计算、再合并结果”这类场景设计的。这篇内容不光讲清楚ForkJoin的运行原理还会把我实际踩过的坑一起说给你听适合正在学Java并发的同学也适合面试前突击并发体系的人。1. ForkJoin框架是什么一个专门为“分而治之”设计的并行框架1.1 并行编程老问题线程池为什么解决不了一切传统并行编程里最常用的方式就是ThreadPoolExecutor Future。把一堆任务丢进线程池让池里的线程从阻塞队列里取任务执行。这个模型解决的是“一堆相互独立的任务并发执行”的问题比如同时处理100个HTTP请求、同时写10个文件每个任务都不依赖别人。可一旦遇到“一个大任务需要递归拆解成无数个子任务”的场景线程池就有点力不从心了。拿数组排序举例。假设要对1亿个整数排序串行排序可能太慢你想并行。如果用普通线程池你得手动把数组切成8段丢给8个线程等它们算完再手动合并结果。切得均匀倒还好切不均匀就会出现某些线程早早就干完了蹲在队列旁边看别人忙活CPU利用率上不去。更麻烦的是如果排序过程中又衍生出新任务比如排序里某个子段又太大需要继续拆这个“任务套任务”的结构用普通线程池写起来非常别扭代码会变得又长又绕。线程序切换和调度本身也有成本。每切一个任务都要创建对象、进入队列、出队列线程内上下文切换更是走钢丝。如果任务粒度太小比如一个任务只算10个数的和那并行调度开销可能比计算本身还大性能反而不如串行。这就是为什么需要一套专门针对“可递归拆分任务”的并行框架。1.2 分而治之如何和并行编程结合“分而治之”的核心套路是把复杂问题拆成若干个独立子问题子问题继续拆直到拆到足够简单可以直接算最后把子结果合并。这个思想在算法课里非常常见快速排序、归并排序、二分查找、树遍历全是它的影子。ForkJoin框架就是把“分而治之”直接翻译成了并行编程模型。Fork对应“拆”也就是把父任务切成左半部分和右半部分分别提交给池里的工作线程Join对应“合并”也就是等待各个分出去的任务执行完取回它们的返回值汇总成结果。整个递归过程会一直持续直到任务小到某个阈值threshold你就可以用普通for循环或者直接计算搞定不再继续拆。这套模型好就好在它跟人类的思维方式是一致的。你不用关心调度细节只需要在compute方法里描述“怎么拆、怎么算、怎么合”。框架负责把任务分发给空闲线程如果你这边的线程正在等待子任务结果它还会顺手帮别人执行队列里的任务尽量不让CPU闲着。这种“自己等的时候帮别人干活”的机制就是后面要详细讲的工作窃取。1.3 典型适用场景什么时候真正需要它ForkJoin不是万能银弹它的主场是“计算密集型且任务可递归拆解”的场景。我实际用得最多的是这几类。第一类是数组聚合计算比如求和、求最大值、求均值、统计满足某个条件元素的数量。数据量一旦上百万甚至上千万串行循环就会明显慢而ForkJoin平均能发挥出多核优势。第二类是排序和搜索类算法例如自定义对象的批量排序、二分搜索的前置结果合并。Java里的Arrays.parallelSort底层就用了ForkJoin相关的排序逻辑可见这类场景有多成熟。第三类是树和图结构遍历比如统计一个大目录下所有文件的行数、计算某个组织架构树的整体绩效这类数据天然有层级结构用递归拆分特别自然。不适合用ForkJoin的场景也有一个共性任务本身太小或者任务不是计算型而是IO型。比如每次任务只读一个文件、发一个HTTP请求这时候瓶颈在网络和磁盘拆再多线程也没用反而可能把池子占死。ForkJoin池的设计是假设任务执行是短促的、纯粹的计算如果你在任务内部sleep或等锁池里的工作线程会被长期占用别人也没法干活。这一点一定要记牢。2. ForkJoin核心组件深度拆解2.1 ForkJoinPool任务池的架构与关键参数ForkJoin框架的门面是ForkJoinPool它继承自AbstractExecutorService但内部结构跟普通线程池完全不同。普通线程池是一个共享阻塞队列加多个线程ForkJoinPool则给每个工作线程分配了一个自己的双端任务队列同时还有一个全局的提交队列用来接收外部发来的初始任务。构造函数有四个关键参数我挨个说下。parallelism表示并行度也就是池里工作线程的数量默认值通常是CPU核数减1。为什么减1因为调用方线程比如main线程在执行invoke和join时也会参与任务计算占了一个“天然线程”。如果你不指定可以手动控制new ForkJoinPool(4)表示最多4个工作线程适合共享机器上的应用。第二个参数是线程工厂可以自定义线程名称、优先级排查问题时非常有用。第三个参数是未捕获异常处理器子任务抛出异常而你没接住时这个处理器能帮你记录错误。第四个参数是asyncMode如果设为true工作线程的本地队列会从LIFO改成FIFO。对于常规的fork/join递归任务保持默认false就好asyncMode更适合事件流式任务。还有一个重头戏是commonPool。Java 8以来parallelStream和很多标准库默认都使用这个公共池而不是你new出来的池。commonPool默认并行度也是CPU核数减1但可以通过JVM参数-Djava.util.concurrent.ForkJoinPool.common.parallelism16修改。公共池归JVM管理属于守护线程进程退出时会自动关闭。你自己创建的ForkJoinPool则记得要shutdown否则非守护线程会卡住进程退出。2.2 ForkJoinTask家族RecursiveTask、RecursiveAction、CountedCompleter使用ForkJoin时你的任务类必须继承ForkJoinTask的子类。最常见的两个是RecursiveTask 和RecursiveAction。RecursiveTask 是有返回值的任务compute方法返回类型T。比如数组求和任务compute返回Long。RecursiveAction则没有返回值compute方法返回void适合“拆分后不需要合并结果”的任务比如给每个文件写入某些数据、遍历整棵树并打印节点。这两个类的compute方法里你需要写清楚“基线条件”和“递归拆分条件”。基线条件达成时直接在本线程内计算不要再fork。递归条件触发时可以创建子任务并调用fork()或invokeAll()。ForkJoinTask本身可以看作轻量级Future。它比Runnable/Callable多了几个关键操作fork()把任务异步提交到当前工作线程的队列join()等待任务结果如果任务还没执行完当前线程不是傻等而是会去偷其他任务来做直到结果就绪invoke()则等价于forkjoin适合一次提交等待完成。这里有个高频面试点ForkJoinTask和普通Future的区别是什么我习惯这样理解Future.get会阻塞当前线程等待任务完成后唤醒ForkJoinTask.join则是一种“协作式等待”它会参与偷取和计算其他任务相当于把空等的时间利用起来。所以写递归任务时推荐使用join而不是get否则会退化成为串行加阻塞。CountedCompleter是另一个分支适合任务之间有完成回调的复杂场景。比如你希望所有子任务都完成后执行某个聚合动作用RecursiveTask得在join处手工聚合用CountedCompleter可以设置onCompletion回调由框架保证在所有pending子任务搞定后触发。不过日常业务里RecursiveTask和RecursiveAction已经覆盖九成需求大家先把这两个玩熟再深入也不迟。2.3 工作窃取Work Stealing算法原理这是ForkJoin最值得讲也是面试官最爱问的部分。先想一个问题如果递归拆出来的任务量不均匀有的分支很快算完有的分支还在深度递归其他线程傻等怎么办工作窃取的思路是这样每个工作线程都有一个自己的任务队列。你在某线程里fork子任务子任务会进入这个线程的队列当前线程执行其他任务时是从队列的“自己这一端”取任务这种取法遵循LIFO后进先出能最大程度利用CPU缓存局部性。因为新拆出来的子任务往往和刚执行的数据关联越晚放进去的任务越可能在缓存里这种策略能让热点数据留在同一线程上减少缓存失效。当一个线程把自己的队列任务全部执行完它不会闲着而是随机挑一个“倒霉”线程偷偷从对方队列的另一端取出任务执行。这样一来繁忙线程仍在处理自己队列尾部的热任务空闲线程则从队列头拿走最老、最冷门的任务两头不会互相挤竞争被降到了最低。我经常用食堂打菜来类比A窗口排队人少B窗口排队人多A窗口的工作人员直接走到B窗口队尾拿一份菜来帮忙炒。不会出现B窗口忙死、A窗口只能干瞪眼的情况。工作窃取的本质就是让忙碌的线程继续干手头的活让空闲的线程主动去分担别人的积压任务。这比“全局共享一个队列大家抢锁拿任务”高效得多因为锁竞争被拆散成了局部队列之间的少量窃取。2.4 提交方式execute、submit、invoke怎么选ForkJoinPool提交任务有六七个方法但归纳起来就三种语义。execute是异步提交没有返回值调用后立即返回适合不需要等结果的火花型任务。submit会返回一个ForkJoinTask对象你可以用它来join或调用get获取结果算异步带结果的提交。invoke最直接提交并等待任务完成返回结果适合在主线程中一口气把一个大任务跑完再继续后续逻辑。还有一个容易忽略的点池内部还能提交“非ForkJoinTask”的Callable和Runnable。提交Runnable后ForkJoinPool会把任务包装成AdaptedRunnable这种任务不走递归拆分直接由某个工作线程执行因此不会有工作窃取带来的额外收益。如果只是为了跑一个Runnable完全没必要用ForkJoinPool普通线程池更合适。实际编写递归任务的时候子任务之间通常走fork/join或invokeAll不会各自往全局队列丢。只有最外层的大任务才通过pool.invoke进入。千万别在compute内部调用pool.execute(subTask)那会把任务扔进全局队列重新分配破坏局部性和窃取效率。3. 实操实现一个“百万级数组求和”的全过程3.1 案例设计从串行到ForkJoin的思路我用一个非常经典的案例来演示给定一个长度为1000万的int数组求所有元素之和。串行版本就是一个for循环累加性能瓶颈很明显。ForkJoin的思路是把数组对半拆开左半部分一个子任务右半部分一个子任务两个子任务都算完后再把结果相加。这里要定两个问题阈值设多大合适任务切到什么时候停止。我的经验法则是不要切到只剩一个元素那样生成的子任务数量爆炸反而拖累性能。一般让单个任务负责的元素个数在1万到10万之间。1000万的数组阈值设成5万会产生200个左右的叶子任务这个规模对ForkJoinPool来说很舒服。当然这个数字不是死的需要按你机器的CPU核数和数组复杂度实测调整。还有个细节值得注意数组可以不拷贝。子任务只需要持有数组的引用以及一个左闭右开的索引范围[start, end)。这样拆分时不会产生拷贝开销只是生成几个轻量的task对象。3.2 完整代码实现下面是一个可直接运行的实现我把关键代码都贴出来。import java.util.concurrent.ForkJoinPool; import java.util.concurrent.RecursiveTask; import java.util.concurrent.ThreadLocalRandom; public class ArraySumTask extends RecursiveTaskLong { private static final int THRESHOLD 50_000; private final int[] array; private final int start; private final int end; public ArraySumTask(int[] 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; } int mid (start end) 1; ArraySumTask leftTask new ArraySumTask(array, start, mid); ArraySumTask rightTask new ArraySumTask(array, mid, end); leftTask.fork(); long rightResult rightTask.compute(); long leftResult leftTask.join(); return leftResult rightResult; } public static void main(String[] args) { int[] array new int[10_000_000]; for (int i 0; i array.length; i) { array[i] ThreadLocalRandom.current().nextInt(1000); } ForkJoinPool pool new ForkJoinPool(); long start System.nanoTime(); long result pool.invoke(new ArraySumTask(array, 0, array.length)); long end System.nanoTime(); pool.shutdown(); System.out.println(结果: result); System.out.println(耗时: (end - start) / 1_000_000.0 ms); } }这段代码看起来不长其实藏着好几个值得展开的门道。我把它们逐个拆开讲。3.3 代码里的关键细节与性能陷阱第一处(start end) 1。这个写法比(start end) / 2更不容易出问题。因为start和end都是正整数且end start直接用加法则可能出现int溢出导致mid变成负数。用无符号右移一位代替除以2可以避免这个隐患。第二处递归分支里我写了leftTask.fork()然后让当前线程直接算rightTask.compute()最后再leftTask.join()。这是ForkJoin官方文档和源码注释推荐的标准写法。为什么如果你两个子任务都fork父线程就会空等着两个子任务回来少干一次活而一个fork一个compute可以利用父线程本身参与计算把任务真正摊到可用线程上。按照经验多fork一个任务不仅不会变得更快反而多了子任务入队、调度、窃取的开销。我见过有人写成两个都fork结果性能比串行还差一半就是因为这个细微区别。第三处基线任务用的是普通for循环。当end - start小于等于阈值时直接在本线程顺序累加。这里不建议再嵌套并行流因为一个叶子任务可能只负责5万个元素并行流拆分开销太大了纯属画蛇添足。第四处pool.invoke会阻塞等待整个任务树完成。invoke返回时最外层任务的结果已经合并完毕。如果你想做异步请求应该改用submit返回的Future再join。另外main函数里我加了System.nanoTime计时这种粗粒度测试只能看个大概严谨的基准最好交给JMH并做预热。3.4 运行结果与调优记录我在一台8核16线程的机器上跑过这个例子阈值设5万1000万数组求和大概是串行版本的5到6倍耗时缩短。具体数字和你机的CPU架构、内存带宽、JVM版本都有关系重点是观察趋势数组越大、核心越多收益越明显当你把阈值调得过高或过低收益都会缩水。我专门做了个小实验验证阈值影响。阈值设为1000时任务量暴增耗时比10万阈值还多出将近一倍原因就是创建了几千个RecursiveTask对象队列压入弹出和窃取竞争全成了瓶颈。阈值设为100万时任务只拆成10个子任务8核机器上很多核心全程空闲效率也不理想。最终5万在这个场景里比较平衡响应时间最短。这个实验说明ForkJoin调优不是玄学核心就一句话别让任务粒度太小也别让任务粒度太大让每个叶子任务执行尽量均匀并且单个任务耗时至少要超过1微秒才好。你把阈值设成让叶子任务数量大约是parallelism的10到50倍通常表现就不错。4. 踩坑记录我遇到过的常见问题与排查方案4.1 任务粒度太细性能不升反降这是我带新人的时候最常见的翻车点。有人写求和任务阈值直接设成1或者只检查end - start是否等于1结果一百万数组就生成一百万个任务对象。ForkJoin虽然能处理大量任务但每个任务入队、fork、join都有固定开销。当任务数量百万级别时光是对象分配就能把堆内存撑起来GC跟着紧张调度时间也远超计算时间。正确的做法是让基线任务执行得久一点。有人喜欢用元素个数当阈值我更推荐用“任务预计执行时长”来判断。如果一个任务执行时间少于1微秒说明拆太细了多于1毫秒说明可能拆得不够。实际项目中我通常先写一个粗略的阈值然后用JFR或JMH测一下叶子任务的平均耗时再做微调。还有一个观察技巧运行完程序后用jstack抓线程栈如果发现大量线程都停在fork/join调度方法上而很少真正在跑你的compute逻辑说明任务拆得太过火先把阈值调大一个数量级试试。4.2 join写错造成假并行、死锁与栈溢出第3节我特意强调了“一个fork、另一个compute”的写法这里再展开说下为什么错。如果左右两个子任务都调用fork再分别join代码仍然正确因为join会等待子结果只是父线程在join等待期间无法像compute那样“亲自干活”整体上少利用了一个执行线程并行度会下降。某些极端情况下如果子任务又继续创建更多子任务两个都fork会导致任务树深度迅速加深递归里每个分支都多一层task最终栈溢出不是开玩笑的。还有人在子任务的compute里写了个left.join()完全没调用fork这种情况就是自己等自己产出的任务任务还没提交就等待结果必然死锁。还有人在任务内部使用get()注意ForkJoinTask.get会抛出InterruptedException和ExecutionException而且它不像join那样会帮助执行其他任务真遇到阻塞场景可能把整个池拖住。如果遇到StackOverflowError第一步看阈值第二步看是不是拆分的维度出了问题。比如改成每次只拆一个元素进去递归深度瞬间变成十万栈肯定扛不住。基线阈值要保证递归深度在几百层以内。4.3 任务里塞了阻塞操作把整个池拖垮这个坑是我在真实业务里踩得最深的一次。当时用parallelStream批量处理一批数据其中一条分支里调了一个同步的远程接口每个接口要花几十毫秒。结果8个并行线程全部卡在等待远程返回上其余几十万条数据排着队饿死最后整个接口超时。ForkJoinPool跟普通线程池最大的不同是它不会为等待中的任务创建额外线程。默认并行度就是CPU核数减1线程数是固定的。你让一个计算型线程池去执行IO等待任务等于把珍贵的执行单位白白占住其他CPU密集型任务根本抢不到执行权。这是设计问题不是调优能解决的。如果你在compute里遇到了网络请求、数据库查询、读文件、持有锁这些场景要么把IO部分分离出来用普通线程池要么换一种完全不同于ForkJoin的思路。把阻塞任务混进ForkJoin池纯属给生产环境埋雷。另外在并行流里写递归任务时也要注意stream递归调用自身并最终parallel执行很容易在commonPool里形成深度依赖和饿死问题。4.4 共享状态并发写导致数据错乱有人图方便在多个子任务里同时往一个共享List里塞结果或者用静态变量累加计数。这必然带来并发写问题。ArrayList在多线程并发add会丢数据HashMap并发put可能死循环哪怕你用AtomicInteger多个核心同时CAS累加也会产生严重的缓存行争抢。我的建议是子任务尽量无状态各自返回结果再由父任务合并。如果确实需要收集结果可以使用ConcurrentLinkedQueue或者最后汇总后统一合并。例如我给一批用户计算积分每个子任务返回一个局部金额父任务累加后把结果统一写入目标表整个过程没有任何共享可变状态既安全又容易排查。一个我自己用过的实用模式是在compute里用ArrayList作为局部变量收集子结果最后返回合并外部只读一个输入数组不写任何共享结构。这样任务天然线程安全不需要加锁性能也不会被锁拖后腿。4.5 排查技巧与工具排查ForkJoin并发问题时我习惯按下面这套顺序来。先用jstack -l 进程号看线程栈。ForkJoinPool的线程名通常是ForkJoinPool-1-worker-1这种格式。如果发现大量worker线程处于WAITING状态并且停靠在某个任务的join或CAS等待上说明任务没有正常完成很可能是死锁或阻塞I/O。如果一堆worker都在执行compute但系统CPU占用率不高可能是任务粒度太大、拆分不够或者机器核心数少。再用VisualVM看堆内存。如果任务对象数量巨大且频繁GC那基本就是任务粒度过细。也可以开JFR录制看ForkJoinPool相关的监控事件能直观看到窃取次数和任务排队时长。最后是性能基准。我习惯把串行版本和ForkJoin版本写在同一测试类里预热几千次后再测量。没有预热的数据不可信JVM还没有完成JIT编译并行代码统计出来的结果会失真。我建议用JMH写带Fork的基准测试省心又准确。5. ForkJoin在Java生态中的身影与面试考点5.1 parallelStream底层与公共池Java 8推出的Stream并行流大多数人早就用上了但未必知道底层是ForkJoin。只要你调用.parallel()它就会通过ForkJoinPool.commonPool()来执行任务。一个List的parallelStream().filter(x).count()背后就是ForkJoin在帮你拆集合。这里有个非常实际的问题commonPool是JVM全局共享的如果一个应用里多个业务团队都使用parallelStream它们会抢同一批线程。某个接口的ForkJoin任务一旦出现阻塞整个应用的并行流都会遭殃。更讽刺的是如果用parallelStream去跑递归型任务其实还是同一个池只是换了个壳子并没有成长生不老。如果想让某个并行流的执行隔离到独立池里可以自己创建ForkJoinPool然后再提交这个流任务ForkJoinPool customPool new ForkJoinPool(4); long result customPool.submit(() - list.parallelStream().mapToLong(Long::longValue).sum() ).get();这种写法能让你控制并发度也让commonPool免受污染强烈推荐在生产项目里使用。5.2 Arrays.parallelSort、ConcurrentHashMap、CompletableFuture中的ForkJoinJava标准库其实早就把ForkJoin用到了底层。Arrays.parallelSort是一个典型它在排序阈值以下用普通快速排序阈值以上用并行归并类的排序实现这套并行逻辑正是基于ForkJoinCommonPool。ConcurrentHashMap的批量操作比如forEachEntry、reduce、search也默认使用commonPool来并行执行。这说明ForkJoin是你日常开发的“隐形基础设施”不一定直接写task但可能因为你调了一个API就间接跑起来。CompletableFuture里有几个静态方法设计到了ForkJoinPool.commonPool作为默认执行器用它来处理包含异步依赖的链式任务。当然更专业的做法是通过工厂方法传入自定义Executor这也是为什么很多高并发项目里会单独定义一个线程池给CompletableFuture用。了解这些生态关联对面试和实际排查都很有帮助当你发现一个并发问题很诡异时先想到底层可能是公共池你就能快速排查出是谁在抢占线程。5.3 经常被面试官问的ForkJoin问题结合最近的Java面试热词我把ForkJoin高频考点整理一下差不多是这几条。ForkJoinPool和ThreadPoolExecutor的区别核心是从工作窃取、本地队列、固定并行度、协作式join四个角度来答。工作窃取算法是怎么实现的要讲双端队列、LIFO自取、FIFO窃取。阈值怎么选从任务粒度、执行时间、经验范围、性能测试来答。为什么不建议两个子任务都fork要从父线程参与计算、减少调度开销来答。commonPool的默认并行度是多少CPU核数减1怎么修改先说系统属性和自定义池。在ForkJoin任务里做阻塞操作有什么问题固定线程数、饥饿、不应在计算任务里塞IO。ForkJoinTask和Future的区别join不会空等而是帮助执行其他任务。并行流默认用的是哪个线程池ForkJoinPool.commonPool。不要死记硬背理解了前面几节的内容这些问题都能自然答出来。答的时候注意别只讲概念最好带一个小例子说明自己的实践面试官更看重你有没有真正写过、踩过坑。6. 我对ForkJoin框架的一些实操经验与建议如果说有什么是我反复强调也不嫌多的那就一句话先想清楚任务到底是计算密集型还是IO密集型再决定要不要用ForkJoin。我见过太多人拿着一个数据结构就开始parallelStream最后被IO卡得比串行还慢。正确解法是默认串行只有明确测量出“瓶颈在多核计算”时再考虑并行方案。项目里真正值得用ForkJoin的场景通常是数据量较大且拆分后子任务相对均匀的聚合计算比如排行榜积分统计、库存汇总、日志压缩、向量运算。这类任务我可以自信地切成几百个块跑收益立竿见影。而那些只拆成两三个子任务就能完成的小功能老老实实用for循环就好别为了“性能”而性能。工作中的另一个心得是ForkJoin代码的可读性很重要。不要把所有逻辑塞进一个巨大的compute方法里尽量把“拆分逻辑”、“基线计算逻辑”、“结果合并逻辑”抽成独立方法哪怕多写几行后续维护的人会感谢你。你可以在关键位置加上注释说明阈值是怎么试出来的不然下个同事接手可能直接改成神秘数字。最后分享一个调优小技巧改阈值前先记录串行基线和当前并发基线每次只改一个变量。你可以把阈值从1000、5000、1万、5万、10万依次测一遍画一条趋势线最优值通常出现在CPU核心数的20到60倍任务数量折算出的范围内。这样得到的数字有据可依而不是拍脑袋。ForkJoin是个好工具但也需要你用工程化的方式去驾驭它。
返回列表