JAVA · Vol.II · DAY 23 · 并发编程
Fork/Join 框架与并行 Stream
面试官:Java 有什么并行计算框架?
"你有一个 1000 万元素的数组要做求和,单线程太慢,你怎么利用多核?"
大部分候选人会说 parallelStream(),但面试官真正想听的是它背后的东西——Fork/Join 框架。parallelStream 只是冰山一角,底下的引擎是 ForkJoinPool。
JDK 7 引入 Fork/Join 框架,JDK 8 的 Stream API 在它之上构建了 parallelStream。理解这套机制,你才能回答:
- parallelStream 什么时候快、什么时候反而更慢?
- 并行流的线程安全陷阱在哪里?
- 自定义并行任务怎么写?fork() 和 join() 到底在干什么?
Fork/Join 的考法通常是"连环追问":先用 parallelStream 引出 ForkJoinPool,再问 work-stealing 算法,最后挖线程安全陷阱。能把这条链路讲清楚的候选人,并发功底不会差。
分治思想:把大问题拆成小问题
Fork/Join 的核心思想就是经典的分治法(Divide and Conquer):
- 分解(Fork):把大任务拆成若干小任务
- 解决:小任务足够小时直接计算
- 合并(Join):把小任务的结果汇总
典型例子:归并排序、大数组求和、MapReduce。我们用一个"对 100 万元素求和"的例子来看分解过程:
分治三步骤:
Task(size) = Task(size/2).fork() + Task(size/2).fork() → join()
当子任务规模低于阈值(threshold)时停止拆分,直接计算。
阈值怎么定?
太小,任务拆分和调度的开销超过并行收益;太大,无法充分利用多核。经验值通常在 1000~10000 之间。判断标准不是"多小算小",而是让每个叶子任务的计算时间明显大于一次 fork/join 的调度成本(后者约几十微秒级)——叶子任务太碎,并行度再高也在交调度税。
ForkJoinPool 架构:Work-Stealing 算法
ForkJoinPool 和普通 ThreadPoolExecutor 最大的区别在于工作窃取(Work-Stealing)算法。普通线程池是"共享一个队列,所有线程抢任务",而 ForkJoinPool 是每个工作线程有自己的双端队列(Deque)。
为什么自己的任务从底部 pop(LIFO),偷来的从顶部取(FIFO)?
LIFO 执行自己的任务:最新 fork 出来的子任务最"小"、最"新鲜",优先处理它可以尽快完成并释放结果给 join。同时,先 fork 的大任务留在 Deque 顶部,正好方便被其他线程偷走。还有一个硬件层面的好处:LIFO 访问的是队列顶部——刚 push 进来的任务还在 L1 缓存里,命中率远高于偷队列尾部的老任务。
FIFO 偷任务:顶部的任务是最早 fork 的,通常是"大块"任务。偷来之后它还会继续 fork 拆分,产生新的子任务填入偷窃者的 Deque,从而让偷窃者也有活干。
| 对比维度 | ThreadPoolExecutor | ForkJoinPool |
|---|---|---|
| 任务队列 | 共享一个 BlockingQueue | 每个工作线程一个 Deque |
| 任务分配 | 线程从共享队列竞争 | work-stealing 自动均衡 |
| 任务类型 | 独立任务 | 可递归拆分的子任务 |
| 线程数 | 可配置 core/max | 默认 = CPU 核心数 |
| 适用场景 | IO 密集、独立请求处理 | CPU 密集、可分治的计算 |
共享队列的竞争点在取任务时(所有线程抢同一个队列头,锁/CAS 竞争 + 缓存行乒乓);work-stealing 的队列是私有的,取自己任务零竞争,只有"偷"时才有一次 CAS——把竞争频率从「每次取」降到「偶尔偷」。任务越多越均衡,这正是分治场景的理想模型。
RecursiveTask 与 RecursiveAction:框架的两个基类
Fork/Join 框架提供了两个抽象基类,区别只有"要不要返回值":
| 类 | 返回值 | 类比 | 典型场景 |
|---|---|---|---|
RecursiveTask<V> | 有返回值 V | Callable | 求和、查找最大值、归并排序 |
RecursiveAction | 无返回值 (void) | Runnable | 数组填充、批量修改、原地排序 |
两个类都是 ForkJoinTask 的子类,都实现了 compute() 抽象方法——你的全部算法逻辑都写在这里。类名里的 "Recursive" 不是限制,是提示:这个框架就是为「自己 fork 出同类子任务」设计的,但不递归直接计算也完全合法(比如一次性任务)。
追问:ForkJoinTask 还有第三个常用子类吗?
有——CountedCompleter<T>:fork 的子任务不需要等,而是靠"计数归零"触发完成回调。适合子任务彼此独立、父任务想在最后统一收口的场景(比所有 join 轻量得多,省掉每次 join 的挂起/唤醒)。JDK 里 parallelStream 的部分内部实现就用了它。面试能主动提这一嘴,算超出预期。
完整示例:用 RecursiveTask 对大数组求和
import java.util.concurrent.RecursiveTask; import java.util.concurrent.ForkJoinPool; public class ArraySumTask extends RecursiveTask<Long> { private static final int THRESHOLD = 5000; // 拆分阈值 private final int[] array; private final int start, 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) / 2; ArraySumTask left = new ArraySumTask(array, start, mid); ArraySumTask right = new ArraySumTask(array, mid, end); // fork() 将子任务提交到 ForkJoinPool 的队列中异步执行 left.fork(); // 异步执行左半部分 // 当前线程直接计算右半部分(比再 fork 一个线程更高效) long rightResult = right.compute(); // 注意:不是 right.fork()! // join() 等待左半部分完成并获取结果 long leftResult = left.join(); return leftResult + rightResult; } }
int[] data = new int[1_000_000]; // ... 填充数据 ... // 方式一:使用公共 ForkJoinPool(parallelStream 也用这个) ForkJoinPool commonPool = ForkJoinPool.commonPool(); long result = commonPool.invoke(new ArraySumTask(data, 0, data.length)); // 方式二:创建自定义 ForkJoinPool(生产推荐,隔离任务) ForkJoinPool customPool = new ForkJoinPool(4); // 4 个工作线程 try { result = customPool.invoke(new ArraySumTask(data, 0, data.length)); } finally { customPool.shutdown(); }
为什么只 fork 一个、compute 另一个?
如果左右两边都 fork(),当前线程就要 join 两次——它在等待的时候什么都做不了,浪费了一个线程。正确做法是:fork 左半部分(让其他线程偷走执行),当前线程自己 compute 右半部分,最后 join 左半部分。这样当前线程始终在干活,不浪费时间。
在 RecursiveTask 的 compute() 中,只 fork 一侧,另一侧直接 compute(),最后 join fork 出去的那一侧。这是 Doug Lea 在 JDK 源码中反复使用的模式。注意一个不对称点:fork 出去的那侧由"别的线程"执行,当前线程只做 join;自己 compute 的那侧不 join(结果在手里),所以全程只有一次 join。
需要 fork 多个子任务怎么办?
用 invokeAll(t1, t2, ..., tn)(JDK 8 起有 varargs 版本):一次性提交全部子任务并等待完成,内部用 CountedCompleter 计数归零的方式收口,比逐个 fork + 逐个 join 高效。适用场景是"拆成 N 份并行跑完再合并"(不是二叉递归)。
fork / join / invoke 三个动词的精确语义
这三个词在源码注释里经常混着出现,面试被追问"它们到底有什么区别"时,要能一句话一个:
| 方法 | 阻塞? | 精确语义 | 谁调用 |
|---|---|---|---|
task.fork() | 不阻塞 | 把子任务提交到当前线程所在池的 Deque,立即返回;任务由池里某个工作线程(不一定是别人)执行 | 必须在池的工作线程内部调用(commonPool 的工作线程里) |
task.join() | 阻塞(仅到任务完成) | 等待 task 完成并返回结果。特殊之处:等待期间当前线程会帮忙干活——偷别的队列的任务(compensating thread 补偿机制),而不是傻等 | 调 fork 的那个线程 |
pool.invoke(task) | 阻塞到完成 | 把 task 作为"根任务"提交给池,调用线程直接参与执行它(而不是扔进队列),并阻塞到出结果 | 任何线程(通常主线程) |
两个容易说错的点:
- fork() 不是"派给别的线程":它只是把任务压入当前线程的 Deque 顶部。如果当前线程还有别的子任务要处理,这个 fork 出去的任务可能暂时留在队列里,甚至被别的线程偷走——"谁执行"由 work-stealing 决定,不由 fork 的调用者指定。
- join() 的"帮忙"机制是它比 wait() 高级的地方:普通 wait 是干等被通知;ForkJoinTask.join 在等待目标完成的同时,会顺手偷其他队列的活干(compensating)——等待时间被用来创造吞吐量,这是 Fork/Join 框架吞吐稳定的核心技巧之一。注意补偿有次数上限,防着无限帮忙。
join() 能拿到子任务的异常吗?
compute() 里抛出的运行时异常/错误会被任务对象记住,join() 时重新抛给调用者(不会静默吞掉)。所以分治代码的异常处理惯例是:叶子节点尽量捕获并聚合,根任务的 join 处做最终兜底——和线程池 submit 后 future.get() 抛 ExecutionException 是一个思想。
parallelStream:一行代码的并行化
JDK 8 的 parallelStream() 本质上就是 Fork/Join 的语法糖。它使用公共 ForkJoinPool.commonPool()(线程数默认 = CPU 核心数 - 1),自动帮你做任务拆分和合并。
List<Integer> numbers = IntStream.rangeClosed(1, 10_000_000) .boxed() .collect(Collectors.toList()); // 串行:单线程 long serialSum = numbers.stream() .filter(n -> n % 2 == 0) .mapToLong(Integer::longValue) .sum(); // 并行:自动利用多核 long parallelSum = numbers.parallelStream() .filter(n -> n % 2 == 0) .mapToLong(Integer::longValue) .sum(); // reduce 写法 long reduced = numbers.parallelStream() .mapToLong(Integer::longValue) .reduce(0L, Long::sum);
底层流程:parallelStream 将数据源通过 Spliterator(可分裂迭代器)拆分成多个分段,每个分段封装成一个 ForkJoinTask 投入 commonPool 并行执行,最后按分段合并(combine)规则汇总结果。拆分不是死板的对半分——Spliterator 带着数据的结构信息(数组可以精确切半,链表只能逐步分),这也是"数组并行流快、链表并行流慢"的底层原因之一。
追问:为什么 commonPool 是「核心数 - 1」而不是「核心数」?
给系统留一个核给不经过池的线程用(比如 main 线程在等池里的结果、JIT、GC)。如果池占满所有核,调用者线程自己都没核可跑,"等池出结果"的线程可能永远拿不到时间片——留一个核是防全局死锁的工程保险。
parallelStream 何时更快:经验公式与边界
parallelStream 何时更快?经验公式:
N > 10,000 && 操作是 CPU 密集型
其中 N 是数据量。数据太少,线程拆分和合并的开销反而超过并行收益。更精确的判据是单次元素的平均处理时间 > 并行化开销(约 1~10 微秒级)——每条记录处理要 1 微秒,1 万条就是 10 毫秒的总计算量,才值得拆给 8 个核。
| 场景 | 数据量 | 操作类型 | parallelStream 表现 |
|---|---|---|---|
| 大数组数值计算 | 100 万+ | CPU 密集(加减乘除) | 快 2~8 倍 |
| 小集合过滤 | < 1000 | CPU 密集 | 更慢(开销 > 收益) |
| 大量 IO 操作 | 任意 | IO 密集(HTTP/DB) | 危险!阻塞公共池 |
| 简单 map 转换 | 10 万+ | CPU 密集 | 略快,收益有限 |
| 复杂 reduce/聚合 | 50 万+ | CPU 密集 | 明显提升 |
为什么"快 2~8 倍"而不是"快 8 倍"?
阿姆达尔定律的朴素版:并行加速比受"串行部分"限制——拆分、合并、结果收集都是串行的;加上分段不均衡(有的段先干完空等)、cache 冷启动,实际加速比 = 核数 × 一个折扣系数。8 核机器上 2~4 倍是常态,宣传"线性加速"的都是没实测过。
commonPool 深水区:整个 JVM 共享的一颗雷
commonPool 的危险性不在"慢",在于它是全 JVM 单例——所有并行消费者共担它的死活:
- 所有
parallelStream()(没指定池的话) - 没传 executor 的
CompletableFuture异步方法(thenApply等) - JDK 内部部分并行实现(如
Arrays.sort()的大数组并行快排)
// ❌ 订单服务里某位同学的"优化": orders.parallelStream() .map(order -> httpCall("http://inventory/stock/" + order.getSku())) // 💥 lambda 里发 HTTP .collect(Collectors.toList()); // 库存服务那天响应从 20ms 涨到 3s: // → 8 个 commonPool 工作线程全部卡在 HTTP 读上 // → 同一 JVM 里另一个报表模块的 parallelStream 求和:排队,无核可用 // → 报表导出超时告警——和订单毫无业务关系,但同一颗池子炸了 // 排查半天才发现:jstack 里 commonPool worker 全部 TIMED_WAITING (parking) // 栈顶停在 HttpURLConnection 的 read 上
三条生产纪律:
- commonPool 只跑纯 CPU 短任务:任何阻塞(IO、锁、sleep、park)都不允许出现在它的 lambda 里。IO 密集走 CompletableFuture + 独立 IO 线程池(见「CompletableFuture 异步编程」篇)。
- 重要业务流不要搭便车:核心链路的并行计算用自建 ForkJoinPool,把业务从"全 JVM 共担"里摘出来。第 15 站给出写法。
- 把它当共享资源监控:
ForkJoinPool.commonPool().getActiveThreadCount()/getQueuedTaskCount()打点——活跃线程数持续贴近池大小、队列堆积,就是有人在占用它。
陷阱 1~2:共享可变状态 & 有状态操作
这是面试最高频的考点——知道怎么用不难,知道什么不能用才是功力。五个陷阱里,前两个最致命(直接数据错乱),后三个最阴险(性能陷阱,不报错只变慢)。
// ❌ 错误示范:多个线程同时写同一个 ArrayList List<String> results = new ArrayList<>(); list.parallelStream() .map(User::getName) .forEach(results::add); // 💥 ArrayList 不是线程安全的! // 结果:数据丢失、ArrayIndexOutOfBoundsException、甚至死循环 // ✅ 正确做法:用 collect 代替外部集合 List<String> results = list.parallelStream() .map(User::getName) .collect(Collectors.toList()); // Collectors 内部处理了并发安全
为什么 collect(Collectors.toList()) 就安全了?因为 Collector 的契约把并发安全内置了:每个线程各自维护一个局部容器(supplier),处理完自己那一段,最后由框架按分段归并(combine)——和 Fork/Join 的 join 合并是同一套思想。而 forEach(list::add) 是"所有线程挤同一个容器",Collector 机制完全没参与。判断标准:结果收集必须走 Collector/终端操作,不许走外部共享集合。
// ❌ sorted() 在并行流下需要先收集所有元素才能排序 // 并行优势被完全抵消,甚至更慢(多了一次并行→串行→并行的切换) list.parallelStream() .sorted() // 💥 有状态中间操作,破坏并行性 .distinct() // 💥 同样有状态,需要全局去重 .limit(100) // ⚠️ 在无序流上还行,有序流上会等待前序完成 .collect(Collectors.toList()); // ✅ 如果不需要原始顺序,用 unordered() 释放约束 list.parallelStream() .unordered() // 告诉框架"我不关心顺序" .distinct() .limit(100) .collect(Collectors.toList());
怎么快速识别"有状态"操作?
问一句:"处理第 100 个元素时,需不需要知道前面 99 个元素的完整信息?"需要 → 有状态(sorted 要知道全部才能排、distinct 要知道全部才能去重、limit 在有序流上要按序截取)→ 并行流在这里必须"收拢"处理。无状态操作(filter/map/peek)每个元素独立决策,天然并行。这个判断法比背清单更耐用。
陷阱 3~5:装箱开销、小数据集、嵌套并行
// ❌ Stream<Integer> 每个元素都是 Integer 对象,内存浪费 + GC 压力 int sum = list.parallelStream() .map(e -> e * 2) // 返回 Stream<Integer>,装箱! .reduce(0, Integer::sum); // 拆箱再求和 // ✅ 使用原始类型流(IntStream / LongStream / DoubleStream) long sum = list.parallelStream() .mapToInt(Integer::intValue) // 转 IntStream,无装箱 .map(e -> e * 2) .asLongStream() .sum();
// ❌ 只有 100 个元素,拆分/调度/合并的开销远大于并行收益 List<Integer> small = Arrays.asList(1, 2, 3, ... , 100); int sum = small.parallelStream().mapToInt(Integer::intValue).sum(); // 比 stream() 慢 5~10 倍! // ✅ 小数据用普通 stream 即可 int sum = small.stream().mapToInt(Integer::intValue).sum();
// ❌ 外层和内层都用 parallelStream —— 争抢同一个 commonPool outerList.parallelStream() .flatMap(item -> innerList.parallelStream() // 💥 嵌套并行,线程互相等待 .map(x -> process(item, x)) ) .collect(Collectors.toList()); // ✅ 只在一个层级使用 parallelStream outerList.parallelStream() .flatMap(item -> innerList.stream() // 内层用串行 stream .map(x -> process(item, x)) ) .collect(Collectors.toList());
为什么嵌套并行不是"更并行"而是"互相饿死"?
内外层都从同一个 commonPool 取工作线程。外层任务占着线程执行到 flatMap 时,内层并行流又向同一个池要线程——池的线程全被外层任务占着,内层任务排队等线程,而"等线程"这件事本身就在占线程。8 线程池跑两层嵌套,很容易出现全员互相等待、吞吐跌到比纯串行还低。规则铁板一块:一条流水线里只允许一个并行级。
顺序的真相:encounter order 与 forEachOrdered
并行执行 = 各段乱序完成,那结果顺序怎么保证?这里有个高频误区要拆干净:
- 有序流(ordered stream)的终端结果保留 encounter order:即使各分段是乱序完成的,
collect(Collectors.toList())出来的 List 元素顺序与串行完全一致——框架在合并分段结果时按原始位置归位(类似并行快排的 merge 阶段)。所以「parallelStream 结果顺序随机」是错的。 - unordered() 是性能开关:声明"我不在乎顺序"后,框架可以跳过按位归位的合并步骤——对 distinct/limit 这类有状态操作,这个差别是数量级的。
- forEach 不保证顺序,forEachOrdered 保证:并行流的
forEach中,各元素的执行时机是乱的(副作用操作可能任意先后发生);forEachOrdered强制按 encounter order 逐个执行副作用——代价是前面一段没做完,后面即使算完了也得等,并行优势被顺序依赖吃掉大半。能不用就不用,要用的场景是"顺序本身有业务含义"(比如按时间序写日志、按优先级通知)。
// ① collect 结果顺序:有序流 → 与串行一致(框架归位合并) list.parallelStream().map(f).collect(Collectors.toList()); // 顺序稳定 ✅ // ② forEach 副作用顺序:不保证 list.parallelStream().forEach(e -> log.info(e)); // 日志顺序是乱的 ⚠️ // ③ forEachOrdered:强制按序执行副作用(慢) list.parallelStream().forEachOrdered(e -> log.info(e)); // 按序但收益打折
追问:并行流的中间结果能"偷看"吗?
不能。流是惰性的、一次性的,中间分段在合并完成前不存在"可观察的中间态"——你唯一能拿到的就是终端操作的结果。这也解释了为什么并行流不适合"边算边输出":要么全部算完拿结果(collect),要么放弃顺序保证(forEach)。
红线场景:这四类情况 parallelStream 绝不用
前几站的陷阱是"性能账",这一站的四类是"安全账"——踩了不是慢,是错:
| 红线场景 | 后果 | 正确替代 |
|---|---|---|
| lambda 里做 IO / 阻塞调用 (HTTP、DB、锁、sleep) |
阻塞 commonPool 工作线程,拖垮全 JVM 的并行任务 | CompletableFuture + 独立 IO 线程池;或批量 IO 后串行处理 |
| 操作有副作用 / 改共享状态 (累加外部变量、add 进外部集合、改对象字段) |
数据错乱、丢失、偶发异常——且并发下极难复现 | 无状态化:结果全部通过 Collector 归并;必须写则改串行或显式同步 |
| 操作不可重入 / 依赖全局状态 (调 SimpleDateFormat、读全局单例的可变字段) |
偶发时间解析错、脏读——线上"每周出一次"的玄学 bug | ThreadLocal 化或不可变对象(DateTimeFormatter);依赖参数化传入 |
| 在请求处理路径里处理小数据 (几百条记录也 parallelStream) |
比串行慢 5~10 倍 + 占用共享池的线程做无用调度 | 串行 stream;确需并行且数据大,用自建池隔离 |
parallelStream 只干一件事:对大份数据做无副作用的纯计算。凡是需要"等外面"(IO)、"改外面"(副作用)、"记外面"(全局状态)的,都该走别的方案。
选型:ForkJoinPool vs parallelStream vs CompletableFuture
面试中经常被问到这三者的选型。核心区别在于任务的性质:
| 维度 | ForkJoinPool + RecursiveTask | parallelStream | CompletableFuture |
|---|---|---|---|
| 任务模型 | 递归分治,自定义拆分逻辑 | 集合数据的并行处理 | 多个独立异步任务的编排 |
| 控制粒度 | 最高:自定义阈值、拆分策略 | 中:框架自动拆分 | 高:手动组合异步阶段 |
| 线程池 | 自定义 ForkJoinPool | 默认 commonPool(可自定义) | 自定义 ExecutorService |
| 适合场景 | 大规模 CPU 密集计算 (矩阵运算、图遍历) |
集合的并行 map/filter/reduce (数据量 > 1万、纯计算) |
IO 密集 + 异步编排 (调多个 RPC 再聚合) |
| IO 操作 | 不适合 | 不适合(阻塞 commonPool) | 最适合 |
| 代码复杂度 | 高(写 RecursiveTask) | 低(一行 parallelStream) | 中(链式 API) |
// 你的任务是什么类型? if (任务是 "集合数据的并行计算" && 数据量 > 10000 && 纯 CPU 操作) { → 用 parallelStream(简洁,够用) } else if (任务是 "递归分治" && 需要精细控制拆分阈值) { → 用 ForkJoinPool + RecursiveTask(自定义并行算法) } else if (任务是 "多个异步 IO 操作的编排") { → 用 CompletableFuture(异步回调 + thenCompose/thenCombine) } else { → 用普通 ThreadPoolExecutor(最简单,覆盖 80% 场景) }
// 生产环境:不要依赖 commonPool,用自定义池隔离 ForkJoinPool customPool = new ForkJoinPool(8); List<String> result; try { // submit 一个 Callable,在里面使用 parallelStream // 此时 parallelStream 会使用 customPool 而不是 commonPool result = customPool.submit(() -> dataList.parallelStream() .filter(e -> e.length() > 5) .map(String::toUpperCase) .collect(Collectors.toList()) ).get(); // get() 阻塞等待完成 } finally { customPool.shutdown(); }
为什么 submit 进去之后 parallelStream 就换了池?
因为 parallelStream 取池的逻辑是"当前线程如果正在某个 ForkJoinPool 的工作线程上执行,就用那个池;否则 fallback 到 commonPool"。submit 进 customPool 后,Callable 跑在 customPool 的工作线程里,里面的 parallelStream 就"看见"了 customPool。这是官方认可的隔离手法,也是它最容易被误解的地方——换池靠的是"在哪执行",不是"怎么声明"。
这一篇你掌握了什么
核心知识点回顾
- Fork/Join 思想:分治 = fork 拆分 + 递归计算 + join 合并;阈值定在"叶子计算时间 > 调度成本"
- Work-Stealing:每线程私有 Deque;自己从底部 pop(LIFO,热缓存),闲者从别人顶部 steal(FIFO,偷大任务还能再拆);竞争点从"每次取"降到"偶尔偷"
- RecursiveTask / RecursiveAction / CountedCompleter:要返回值、不要返回值、计数收口三种基类
- fork / join / invoke 语义:fork 不阻塞(压入当前线程 Deque);join 阻塞但等待时顺手帮忙(compensating);invoke 由调用线程直接执行根任务;fork 一侧 compute 一侧 + invokeAll 收多子任务
- commonPool:全 JVM 单例、核数-1、parallelStream/CF/Arrays.sort 共用;只跑纯 CPU 短任务,阻塞 IO 会炸全场;重要业务自建池(submit 换池靠"在哪执行")
- 顺序真相:有序流 collect 结果与串行一致(框架归位);unordered() 是性能开关;forEach 乱序、forEachOrdered 按序但慢
- 五大陷阱:共享可变状态(走 Collector 归并)、有状态操作(识别法:"需不需要知道前面元素")、装箱(IntStream)、小数据集(串行)、嵌套并行(单流水线单并行级)
- 选型:CPU 密集 + 大集合 → parallelStream;递归分治算法 → ForkJoinPool;IO 编排 → CompletableFuture;其余 → 普通线程池
一句话总结
Fork/Join 是"分治 + 偷活"的并行引擎:大任务拆小、各线程私有队列、闲者偷忙者的顶部任务;parallelStream 是它的语法糖,好用但只许喂"大份 + 纯计算 + 无副作用"的数据,IO 请绕道 CompletableFuture。
🎤 面试 30 秒总结
"Fork/Join 框架基于分治思想:任务递归拆分(fork),叶子直接算,结果向上合并(join)。线程模型是 work-stealing——每个工作线程有自己的双端队列,自己从底部 LIFO 取最新任务,空闲时从别的线程队列顶部 FIFO 偷最老的大任务,把锁竞争降到最低。写 RecursiveTask 的惯例是只 fork 一侧、当前线程 compute 另一侧、最后 join 一次;fork 不阻塞,join 阻塞但等待时会帮忙偷任务。parallelStream 是它的语法糖,底层用 commonPool(核数减一、全 JVM 共享),所以三条纪律:只跑纯 CPU、lambda 无副作用无 IO、数据量上万再并行;顺序上,有序流的 collect 结果和串行一致,forEach 乱序、forEachOrdered 按序但慢。IO 密集编排应该用 CompletableFuture 加独立线程池,不该碰 commonPool。"
Comments · 评论