Lesson 23 · 并发编程
Fork/Join 框架与并行 Stream
面试官:Java 有什么并行计算框架?
"你有一个 1000 万元素的数组要做求和,单线程太慢,你怎么利用多核?"
大部分候选人会说 parallelStream(),但面试官真正想听的是它背后的东西——Fork/Join 框架。parallelStream 只是冰山一角,底下的引擎是 ForkJoinPool。
JDK 7 引入 Fork/Join 框架,JDK 8 的 Stream API 在它之上构建了 parallelStream。理解这套机制,你才能回答:
- parallelStream 什么时候快、什么时候反而更慢?
- 并行流的线程安全陷阱在哪里?
- 自定义并行任务怎么写?
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 之间。Doug Lea(Fork/Join 框架作者)建议让每个子任务的计算量至少是顺序执行计算量的 100~10000 倍。
ForkJoinPool 架构:Work-Stealing 算法
ForkJoinPool 和普通 ThreadPoolExecutor 最大的区别在于工作窃取(Work-Stealing)算法。普通线程池是"共享一个队列,所有线程抢任务",而 ForkJoinPool 是每个工作线程有自己的双端队列(Deque)。
为什么自己的任务从底部 pop(LIFO),偷来的从顶部取(FIFO)?
LIFO 执行自己的任务:最新 fork 出来的子任务最"小"、最"新鲜",优先处理它可以尽快完成并释放结果给 join。同时,先 fork 的大任务留在 Deque 顶部,正好方便被其他线程偷走。
FIFO 偷任务:顶部的任务是最早 fork 的,通常是"大块"任务。偷来之后它还会继续 fork 拆分,产生新的子任务填入偷窃者的 Deque,从而让偷窃者也有活干。
| 对比维度 | ThreadPoolExecutor | ForkJoinPool |
|---|---|---|
| 任务队列 | 共享一个 BlockingQueue | 每个工作线程一个 Deque |
| 任务分配 | 线程从共享队列竞争 | work-stealing 自动均衡 |
| 任务类型 | 独立任务 | 可递归拆分的子任务 |
| 线程数 | 可配置 core/max | 默认 = CPU 核心数 |
| 适用场景 | IO 密集、独立请求处理 | CPU 密集、可分治的计算 |
RecursiveTask 与 RecursiveAction
Fork/Join 框架提供了两个抽象基类:
| 类 | 返回值 | 类比 | 典型场景 |
|---|---|---|---|
RecursiveTask<V> | 有返回值 V | Callable | 求和、查找最大值、归并排序 |
RecursiveAction | 无返回值 (void) | Runnable | 数组填充、批量修改、排序原地操作 |
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(),当前线程就要 join 两次——它在等待的时候什么都做不了,浪费了一个线程。正确做法是:fork 左半部分(让其他线程偷走执行),当前线程自己 compute 右半部分,最后 join 左半部分。这样当前线程始终在干活,不浪费时间。
在 RecursiveTask 的 compute() 中,只 fork 一侧,另一侧直接 compute(),最后 join fork 出去的那一侧。这是 Doug Lea 在 JDK 源码中反复使用的模式。
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,每个子任务在 ForkJoinPool 的工作线程中并行执行,最后合并结果。
parallelStream 何时更快?经验公式:
N > 10,000 && 操作是 CPU 密集型
其中 N 是数据量。数据太少,线程拆分和合并的开销反而超过并行收益。
| 场景 | 数据量 | 操作类型 | parallelStream 表现 |
|---|---|---|---|
| 大数组数值计算 | 100 万+ | CPU 密集(加减乘除) | 快 2~8 倍 |
| 小集合过滤 | < 1000 | CPU 密集 | 更慢(开销 > 收益) |
| 大量 IO 操作 | 任意 | IO 密集(HTTP/DB) | 危险!阻塞公共池 |
| 简单 map 转换 | 10 万+ | CPU 密集 | 略快,收益有限 |
| 复杂 reduce/聚合 | 50 万+ | CPU 密集 | 明显提升 |
面试官追问:parallelStream 里能调用 HTTP 接口吗?
绝对不行!parallelStream 默认使用 ForkJoinPool.commonPool(),这个池是整个 JVM 共享的。如果你在 lambda 里发起 HTTP 调用(IO 阻塞),会阻塞公共池的工作线程,导致所有使用 commonPool 的任务(包括其他 parallelStream)全部被拖慢。IO 密集任务应该用 CompletableFuture + 自定义线程池。
parallelStream 的 5 大陷阱
这是面试最高频的考点——知道怎么用不难,知道什么不能用才是功力。
// ❌ 错误示范:多个线程同时写同一个 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 内部处理了并发安全
// ❌ sorted() 在并行流下需要先收集所有元素才能排序 // 并行优势被完全抵消,甚至更慢(多了一次并行→串行→并行的切换) list.parallelStream() .sorted() // 💥 有状态中间操作,破坏并行性 .distinct() // 💥 同样有状态,需要全局去重 .limit(100) // ⚠️ 在无序流上还行,有序流上会等待前序完成 .collect(Collectors.toList()); // ✅ 如果不需要原始顺序,用 unordered() 释放约束 list.parallelStream() .unordered() // 告诉框架"我不关心顺序" .distinct() .limit(100) .collect(Collectors.toList());
// ❌ 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());
使用前逐项检查:(1) Lambda 无副作用、不修改外部共享变量;(2) 避免 sorted/distinct 等有状态操作;(3) 用 IntStream/LongStream 代替装箱流;(4) 数据量 > 10,000;(5) 不嵌套 parallelStream;(6) 不在 Lambda 里做 IO。
总结:Fork/Join 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(); }
全文核心要点回顾
- Fork/Join 思想:分治 = fork 拆分 + 递归计算 + join 合并
- Work-Stealing:每个线程有自己的 Deque,空闲线程从繁忙线程顶部偷任务
- RecursiveTask:只 fork 一侧、compute 另一侧、最后 join
- parallelStream:Fork/Join 的语法糖,底层用 commonPool
- 五大陷阱:共享可变状态、有状态操作、装箱开销、小数据集、嵌套并行
- 选型:CPU 密集 + 集合 → parallelStream;递归算法 → ForkJoinPool;IO 编排 → CompletableFuture
"Fork/Join 框架基于分治思想,通过 work-stealing 算法让空闲线程从繁忙线程的 Deque 偷任务,实现负载均衡。parallelStream 是它的语法糖,默认用 commonPool。使用时要注意五个陷阱:不在 Lambda 里修改共享变量、避免 sorted/distinct 等有状态操作、用 IntStream 避免装箱、数据量要大于 1 万、不嵌套并行流。IO 密集场景应该用 CompletableFuture 而不是 parallelStream。"