首页 / Java 学习笔记 / 23

JAVA · Vol.II · DAY 23 · 并发编程

Fork/Join 框架与并行 Stream

中级中高#并发#框架
第 1 站

面试官: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 算法,最后挖线程安全陷阱。能把这条链路讲清楚的候选人,并发功底不会差。

第 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 之间。判断标准不是"多小算小",而是让每个叶子任务的计算时间明显大于一次 fork/join 的调度成本(后者约几十微秒级)——叶子任务太碎,并行度再高也在交调度税。

第 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 顶部,正好方便被其他线程偷走。还有一个硬件层面的好处:LIFO 访问的是队列顶部——刚 push 进来的任务还在 L1 缓存里,命中率远高于偷队列尾部的老任务。

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

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

共享队列的竞争点在取任务时(所有线程抢同一个队列头,锁/CAS 竞争 + 缓存行乒乓);work-stealing 的队列是私有的,取自己任务零竞争,只有"偷"时才有一次 CAS——把竞争频率从「每次取」降到「偶尔偷」。任务越多越均衡,这正是分治场景的理想模型。

第 4 站

RecursiveTask 与 RecursiveAction:框架的两个基类

Fork/Join 框架提供了两个抽象基类,区别只有"要不要返回值":

返回值类比典型场景
RecursiveTask<V>有返回值 VCallable求和、查找最大值、归并排序
RecursiveAction无返回值 (void)Runnable数组填充、批量修改、原地排序

两个类都是 ForkJoinTask 的子类,都实现了 compute() 抽象方法——你的全部算法逻辑都写在这里。类名里的 "Recursive" 不是限制,是提示:这个框架就是为「自己 fork 出同类子任务」设计的,但不递归直接计算也完全合法(比如一次性任务)。

追问:ForkJoinTask 还有第三个常用子类吗?

有——CountedCompleter<T>:fork 的子任务不需要等,而是靠"计数归零"触发完成回调。适合子任务彼此独立、父任务想在最后统一收口的场景(比所有 join 轻量得多,省掉每次 join 的挂起/唤醒)。JDK 里 parallelStream 的部分内部实现就用了它。面试能主动提这一嘴,算超出预期。

第 5 站

完整示例:用 RecursiveTask 对大数组求和

ArraySumTask.java —— 一个可以跑的分治求和任务
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();
}
第 6 站

为什么只 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 份并行跑完再合并"(不是二叉递归)。

第 7 站

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 是一个思想。

第 8 站

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(可分裂迭代器)拆分成多个分段,每个分段封装成一个 ForkJoinTask 投入 commonPool 并行执行,最后按分段合并(combine)规则汇总结果。拆分不是死板的对半分——Spliterator 带着数据的结构信息(数组可以精确切半,链表只能逐步分),这也是"数组并行流快、链表并行流慢"的底层原因之一。

追问:为什么 commonPool 是「核心数 - 1」而不是「核心数」?

给系统留一个核给不经过池的线程用(比如 main 线程在等池里的结果、JIT、GC)。如果池占满所有核,调用者线程自己都没核可跑,"等池出结果"的线程可能永远拿不到时间片——留一个核是防全局死锁的工程保险。

第 9 站

parallelStream 何时更快:经验公式与边界

parallelStream 何时更快?经验公式:

N > 10,000 && 操作是 CPU 密集型

其中 N 是数据量。数据太少,线程拆分和合并的开销反而超过并行收益。更精确的判据是单次元素的平均处理时间 > 并行化开销(约 1~10 微秒级)——每条记录处理要 1 微秒,1 万条就是 10 毫秒的总计算量,才值得拆给 8 个核。

场景数据量操作类型parallelStream 表现
大数组数值计算100 万+CPU 密集(加减乘除)快 2~8 倍
小集合过滤< 1000CPU 密集更慢(开销 > 收益)
大量 IO 操作任意IO 密集(HTTP/DB)危险!阻塞公共池
简单 map 转换10 万+CPU 密集略快,收益有限
复杂 reduce/聚合50 万+CPU 密集明显提升

为什么"快 2~8 倍"而不是"快 8 倍"?

阿姆达尔定律的朴素版:并行加速比受"串行部分"限制——拆分、合并、结果收集都是串行的;加上分段不均衡(有的段先干完空等)、cache 冷启动,实际加速比 = 核数 × 一个折扣系数。8 核机器上 2~4 倍是常态,宣传"线性加速"的都是没实测过。

第 10 站

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() 打点——活跃线程数持续贴近池大小、队列堆积,就是有人在占用它。
第 11 站

陷阱 1~2:共享可变状态 & 有状态操作

这是面试最高频的考点——知道怎么用不难,知道什么不能用才是功力。五个陷阱里,前两个最致命(直接数据错乱),后三个最阴险(性能陷阱,不报错只变慢)。

陷阱 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 内部处理了并发安全
背后的原理

为什么 collect(Collectors.toList()) 就安全了?因为 Collector 的契约把并发安全内置了:每个线程各自维护一个局部容器(supplier),处理完自己那一段,最后由框架按分段归并(combine)——和 Fork/Join 的 join 合并是同一套思想。而 forEach(list::add) 是"所有线程挤同一个容器",Collector 机制完全没参与。判断标准:结果收集必须走 Collector/终端操作,不许走外部共享集合

陷阱 2:有状态操作在并行下性能暴降
// ❌ 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)每个元素独立决策,天然并行。这个判断法比背清单更耐用。

第 12 站

陷阱 3~5:装箱开销、小数据集、嵌套并行

陷阱 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());

为什么嵌套并行不是"更并行"而是"互相饿死"?

内外层都从同一个 commonPool 取工作线程。外层任务占着线程执行到 flatMap 时,内层并行流又向同一个池要线程——池的线程全被外层任务占着,内层任务排队等线程,而"等线程"这件事本身就在占线程。8 线程池跑两层嵌套,很容易出现全员互相等待、吞吐跌到比纯串行还低。规则铁板一块:一条流水线里只允许一个并行级

第 13 站

顺序的真相: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)。

第 14 站

红线场景:这四类情况 parallelStream 绝不用

前几站的陷阱是"性能账",这一站的四类是"安全账"——踩了不是慢,是错:

红线场景后果正确替代
lambda 里做 IO / 阻塞调用
(HTTP、DB、锁、sleep)
阻塞 commonPool 工作线程,拖垮全 JVM 的并行任务 CompletableFuture + 独立 IO 线程池;或批量 IO 后串行处理
操作有副作用 / 改共享状态
(累加外部变量、add 进外部集合、改对象字段)
数据错乱、丢失、偶发异常——且并发下极难复现 无状态化:结果全部通过 Collector 归并;必须写则改串行或显式同步
操作不可重入 / 依赖全局状态
(调 SimpleDateFormat、读全局单例的可变字段)
偶发时间解析错、脏读——线上"每周出一次"的玄学 bug ThreadLocal 化或不可变对象(DateTimeFormatter);依赖参数化传入
在请求处理路径里处理小数据
(几百条记录也 parallelStream)
比串行慢 5~10 倍 + 占用共享池的线程做无用调度 串行 stream;确需并行且数据大,用自建池隔离
一句话红线

parallelStream 只干一件事:对大份数据做无副作用的纯计算。凡是需要"等外面"(IO)、"改外面"(副作用)、"记外面"(全局状态)的,都该走别的方案。

第 15 站

选型:ForkJoinPool 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();
}

为什么 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 · 评论