Lesson 23 · 并发编程

Fork/Join 框架与并行 Stream

中级·#并发·#框架

第 1 站

面试官:Java 有什么并行计算框架?

"你有一个 1000 万元素的数组要做求和,单线程太慢,你怎么利用多核?"

大部分候选人会说 parallelStream(),但面试官真正想听的是它背后的东西——Fork/Join 框架。parallelStream 只是冰山一角,底下的引擎是 ForkJoinPool

JDK 7 引入 Fork/Join 框架,JDK 8 的 Stream API 在它之上构建了 parallelStream。理解这套机制,你才能回答:

  • parallelStream 什么时候快、什么时候反而更慢?
  • 并行流的线程安全陷阱在哪里?
  • 自定义并行任务怎么写?
面试考什么?

Fork/Join 的考法通常是"连环追问":先用 parallelStream 引出 ForkJoinPool,再问 work-stealing 算法,最后挖线程安全陷阱。能把这条链路讲清楚的候选人,并发功底不会差。

第 2 站

分治思想:把大问题拆成小问题

Fork/Join 的核心思想就是经典的分治法(Divide and Conquer)

  1. 分解(Fork):把大任务拆成若干小任务
  2. 解决:小任务足够小时直接计算
  3. 合并(Join):把小任务的结果汇总

典型例子:归并排序、大数组求和、MapReduce。我们用一个"对 100 万元素求和"的例子来看分解过程:

SumTask(0, 1000000) sum = left + right SumTask(0, 500000) sum = left + right SumTask(500000, 1000000) sum = left + right SumTask(0,250000) 继续拆分... SumTask(250000,500000) 继续拆分... SumTask(500000,750000) 继续拆分... SumTask(750000,1000000) 继续拆分... ... 直到 size < 阈值 → 直接遍历求和 ↓ join:自底向上合并结果 ↓ 叶子结果 → 子任务合并 → 根任务得到最终 sum 每一层 fork 拆任务,每一层 join 收结果,递归直到完成
图 1 分治思想——大数组求和的任务分解树

分治三步骤:

Task(size) = Task(size/2).fork() + Task(size/2).fork() → join()

当子任务规模低于阈值(threshold)时停止拆分,直接计算。

阈值怎么定?

太小,任务拆分和调度的开销超过并行收益;太大,无法充分利用多核。经验值通常在 1000~10000 之间。Doug Lea(Fork/Join 框架作者)建议让每个子任务的计算量至少是顺序执行计算量的 100~10000 倍。

第 3 站

ForkJoinPool 架构:Work-Stealing 算法

ForkJoinPool 和普通 ThreadPoolExecutor 最大的区别在于工作窃取(Work-Stealing)算法。普通线程池是"共享一个队列,所有线程抢任务",而 ForkJoinPool 是每个工作线程有自己的双端队列(Deque)

ForkJoinPool — Work-Stealing 可视化 Worker Thread 0 task A3 task A2 task A1 task A0 push → pop ← 双端队列 (Deque) 底部 pop 自己的任务 Worker Thread 1 task B4 task B3 task B2 task B1 task B0 队列很深(忙) Worker Thread 2 (空) 没有任务了... Worker Thread 3 task D1 task D0 STEAL: 偷 B4 Work-Stealing 规则: 1. 每个线程从自己 Deque 的【底部】pop 任务执行(LIFO — 最新 fork 的子任务优先) 2. 空闲线程从其他线程 Deque 的【顶部】steal 任务(FIFO — 偷最老的大任务) 3. 为什么从顶部偷?因为顶部的任务是"大任务",被偷走后还能继续 fork 拆分,不会和原主人竞争
图 2 ForkJoinPool Work-Stealing 算法——空闲线程从繁忙线程偷任务

为什么自己的任务从底部 pop(LIFO),偷来的从顶部取(FIFO)?

LIFO 执行自己的任务:最新 fork 出来的子任务最"小"、最"新鲜",优先处理它可以尽快完成并释放结果给 join。同时,先 fork 的大任务留在 Deque 顶部,正好方便被其他线程偷走。

FIFO 偷任务:顶部的任务是最早 fork 的,通常是"大块"任务。偷来之后它还会继续 fork 拆分,产生新的子任务填入偷窃者的 Deque,从而让偷窃者也有活干。

对比维度ThreadPoolExecutorForkJoinPool
任务队列共享一个 BlockingQueue每个工作线程一个 Deque
任务分配线程从共享队列竞争work-stealing 自动均衡
任务类型独立任务可递归拆分的子任务
线程数可配置 core/max默认 = CPU 核心数
适用场景IO 密集、独立请求处理CPU 密集、可分治的计算
第 4 站

RecursiveTask 与 RecursiveAction

Fork/Join 框架提供了两个抽象基类:

返回值类比典型场景
RecursiveTask<V>有返回值 VCallable求和、查找最大值、归并排序
RecursiveAction无返回值 (void)Runnable数组填充、批量修改、排序原地操作
完整示例:用 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;
    }
}
调用方式:fork() vs invoke()
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 源码中反复使用的模式。

第 5 站

parallelStream:一行代码的并行化

JDK 8 的 parallelStream() 本质上就是 Fork/Join 的语法糖。它使用公共 ForkJoinPool.commonPool()(线程数默认 = CPU 核心数 - 1),自动帮你做任务拆分和合并。

parallelStream 典型用法
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 倍
小集合过滤< 1000CPU 密集更慢(开销 > 收益)
大量 IO 操作任意IO 密集(HTTP/DB)危险!阻塞公共池
简单 map 转换10 万+CPU 密集略快,收益有限
复杂 reduce/聚合50 万+CPU 密集明显提升

面试官追问:parallelStream 里能调用 HTTP 接口吗?

绝对不行!parallelStream 默认使用 ForkJoinPool.commonPool(),这个池是整个 JVM 共享的。如果你在 lambda 里发起 HTTP 调用(IO 阻塞),会阻塞公共池的工作线程,导致所有使用 commonPool 的任务(包括其他 parallelStream)全部被拖慢。IO 密集任务应该用 CompletableFuture + 自定义线程池。

第 6 站

parallelStream 的 5 大陷阱

这是面试最高频的考点——知道怎么用不难,知道什么不能用才是功力。

陷阱 1:共享可变状态(副作用 Lambda)
// ❌ 错误示范:多个线程同时写同一个 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 内部处理了并发安全
陷阱 2:有状态操作在并行下性能暴降
// ❌ sorted() 在并行流下需要先收集所有元素才能排序
// 并行优势被完全抵消,甚至更慢(多了一次并行→串行→并行的切换)
list.parallelStream()
    .sorted()       // 💥 有状态中间操作,破坏并行性
    .distinct()     // 💥 同样有状态,需要全局去重
    .limit(100)     // ⚠️ 在无序流上还行,有序流上会等待前序完成
    .collect(Collectors.toList());

// ✅ 如果不需要原始顺序,用 unordered() 释放约束
list.parallelStream()
    .unordered()    // 告诉框架"我不关心顺序"
    .distinct()
    .limit(100)
    .collect(Collectors.toList());
陷阱 3:装箱/拆箱开销
// ❌ 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();
陷阱 4:小数据集用并行流 = 杀鸡用牛刀
// ❌ 只有 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();
陷阱 5:嵌套 parallelStream 导致线程饥饿
// ❌ 外层和内层都用 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());
parallelStream 安全清单

使用前逐项检查:(1) Lambda 无副作用、不修改外部共享变量;(2) 避免 sorted/distinct 等有状态操作;(3) 用 IntStream/LongStream 代替装箱流;(4) 数据量 > 10,000;(5) 不嵌套 parallelStream;(6) 不在 Lambda 里做 IO。

第 7 站

总结:Fork/Join vs parallelStream vs CompletableFuture

面试中经常被问到这三者的选型。核心区别在于任务的性质

维度ForkJoinPool + RecursiveTaskparallelStreamCompletableFuture
任务模型 递归分治,自定义拆分逻辑 集合数据的并行处理 多个独立异步任务的编排
控制粒度 最高:自定义阈值、拆分策略 中:框架自动拆分 高:手动组合异步阶段
线程池 自定义 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% 场景)
}
自定义 ForkJoinPool 给 parallelStream 用
// 生产环境:不要依赖 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
面试 30 秒总结

"Fork/Join 框架基于分治思想,通过 work-stealing 算法让空闲线程从繁忙线程的 Deque 偷任务,实现负载均衡。parallelStream 是它的语法糖,默认用 commonPool。使用时要注意五个陷阱:不在 Lambda 里修改共享变量、避免 sorted/distinct 等有状态操作、用 IntStream 避免装箱、数据量要大于 1 万、不嵌套并行流。IO 密集场景应该用 CompletableFuture 而不是 parallelStream。"