Java 基础体系 · 第 67/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java ForkJoinPool 与并行流:工作窃取、阻塞和性能边界

ForkJoinPool 和并行流经常被放在一起讨论,但它们解决的问题并不相同:

  • ForkJoinPool 是一种执行器,核心是把任务拆分成子任务,并通过工作窃取提高处理器利用率。
  • 并行流是 Stream API 的一种执行模式,负责把数据处理管道划分为多个任务,通常依赖 ForkJoinPool.commonPool() 执行。
  • 工作窃取适合大量可拆分、计算密集、运行时间相对均衡的任务。
  • 阻塞会让工作线程在等待期间无法执行其他任务;如果阻塞发生在 Fork/Join 任务内部,还可能造成吞吐下降甚至线程饥饿。
  • 并行并不等于更快。任务拆分成本、合并成本、内存带宽、同步、数据倾斜和阻塞时间,都会构成性能边界。

本文以 Java 25 LTS 的 API 和语言语义为范围,重点解释这几个机制之间的因果关系,而不是把“并行流适合大数据量”作为未经推导的经验结论。

一、先建立几个必要概念

1. 任务、线程和并行度不是同一个概念

任务是待执行的工作单元,例如:

long sum = 1_000_000L;

这个任务可以继续拆成:

sum(1..250000)
sum(250001..500000)
sum(500001..750000)
sum(750001..1000000)

线程是执行任务的运行实体。一个线程在同一时刻通常只能执行一个 Java 代码路径。

并行度是同时参与执行任务的工作线程数量,通常记为 p。在 ForkJoinPool 中,并行度是池希望维持的工作线程数量目标,而不是“池中永远只有这么多线程”的绝对保证。阻塞补偿、线程启动失败以及实现细节都可能影响实际线程数。

如果任务总量是 W,理想情况下有 p 个工作线程,则计算时间至少接近:

TpWpT_p \geq \frac{W}{p}

但这只是把工作平均分配且忽略拆分、调度、合并和内存访问成本后的下界。真实执行时间还要受到最长依赖链、也就是 spancritical path 的限制。

设:

  • T1:单线程完成全部工作所需时间;
  • T∞:无限处理器下,受任务依赖关系限制的最短时间;
  • p:并行度;
  • Tp:使用 p 个工作线程时的执行时间。

则常用的理论下界是:

Tpmax(T1p,T)T_p \geq \max\left(\frac{T_1}{p}, T_\infty\right)

这解释了为什么增加线程数最终会失效:即使可分配的工作很多,只要合并、串行阶段或最长依赖链占据时间,执行时间也不能无限下降。

2. Fork/Join 的“Join”不是普通线程等待

Fork/Join 模型把一个大任务拆分成多个子任务:

大任务
├── 子任务 A
│   ├── A1
│   └── A2
└── 子任务 B
    ├── B1
    └── B2

fork 表示提交或安排一个子任务,join 表示等待并取得该子任务结果。

关键点是:Fork/Join 工作线程在调用 join 时,通常不会像普通线程那样简单地停在那里。它可以继续执行其他可运行的 Fork/Join 任务,或者帮助完成被等待的任务。这种行为称为 helping,是减少任务等待空转的重要机制。

但这种帮助主要针对 Fork/Join 任务之间的依赖。如果线程是在等待网络响应、磁盘、数据库连接、锁或 Thread.sleep,Fork/Join 框架并不能凭空把外部等待转换成可执行的计算任务。

二、ForkJoinPool 的核心结构:工作窃取

1. 每个工作线程通常拥有自己的任务队列

ForkJoinPool 的核心思想不是让所有任务都进入一个全局队列,而是让工作线程拥有局部队列:

Worker-1 deque: [任务 1, 任务 2, 任务 3]
Worker-2 deque: [任务 4]
Worker-3 deque: []
Worker-4 deque: [任务 5, 任务 6]

一个工作线程生成子任务时,通常优先把子任务放入自己的队列。这样做有两个好处:

  1. 本地提交和本地获取可以减少全局锁竞争;
  2. 递归拆分产生的子任务具有局部性,缓存行为通常比全局随机调度更好。

ForkJoinPool 的具体队列布局和调度细节属于实现层面,不是应用程序可以依赖的 API 契约。常见实现中,工作线程优先处理自己的任务,而空闲线程从其他工作线程的队列中“偷”任务。

2. 为什么要从另一端窃取

可以把双端队列抽象成:

本地工作线程处理端  <---- [A, B, C, D] ---->  窃取端

拥有该队列的线程从一端取任务,窃取者从另一端取任务。这样本地线程可以快速处理自己最近生成的任务,而窃取者获取较老的任务,减少双方争用同一个队列位置。

默认模式通常偏向让工作线程使用栈式的本地处理顺序,这有利于递归 Fork/Join 任务的局部性。ForkJoinPoolasyncMode 可以让未被 Join 的异步任务采用 FIFO 风格,但它更适用于事件式任务,而不是所有递归计算都应该打开的“性能开关”。

3. 工作窃取解决的是负载不均,不是所有性能问题

假设四个任务的执行时间分别为:

任务 A:100 ms
任务 B:100 ms
任务 C:100 ms
任务 D:1000 ms

四个线程各拿一个任务时,整体完成时间接近 1000 ms,因为最长的任务决定了结束时间。

如果任务 D 可以继续拆分:

D
├── D1:250 ms
├── D2:250 ms
├── D3:250 ms
└── D4:250 ms

空闲线程就可以窃取 D 的子任务,整体负载更均衡。

但如果 D 是不可拆分的黑盒调用,工作窃取无法把它切开。此时池再增加线程,也不能缩短 D 的执行时间。这是工作窃取的第一个边界:

必须存在足够多、粒度合适且相互独立的任务,工作窃取才有可发挥的空间。

4. join 如何与工作窃取配合

概念上的执行过程如下:

sequenceDiagram
    participant W1 as Worker-1
    participant Q1 as Worker-1队列
    participant W2 as Worker-2
    participant Q2 as Worker-2队列

    W1->>Q1: fork(A)
    W1->>Q1: fork(B)
    W1->>W1: 继续处理当前任务
    W2->>Q1: 窃取(A)
    W2->>W2: 执行 A
    W1->>W1: join(B)
    W1->>Q1: 处理或帮助完成 B
    W2->>Q2: 生成 A 的子任务
    W1->>W1: 获取 B 的结果

这里的关键不是“一个线程启动另一个线程”,而是:

  1. 一个 Fork/Join 任务产生子任务;
  2. 子任务进入某个工作队列;
  3. 其他空闲工作线程可以窃取;
  4. 父任务在 Join 时,尽量继续参与可运行任务;
  5. 所有依赖完成后,父任务合并结果。

因此,递归任务中的 join 和普通线程池中的 Future.get() 不能简单等同。二者都可能等待结果,但 Fork/Join 的设计目标是让等待者继续帮助池完成工作。

三、并行流是如何进入 ForkJoinPool 的

1. 并行流的执行数据流

并行流并不是“把 for 循环自动复制到多个线程”。它需要同时处理数据源、拆分器、流水线和终端操作。

典型过程是:

flowchart LR
    S[数据源] --> SP[Spliterator]
    SP --> D[递归拆分]
    D --> T1[子任务 1]
    D --> T2[子任务 2]
    D --> T3[子任务 3]
    T1 --> M[中间操作]
    T2 --> M
    T3 --> M
    M --> R[终端操作与结果合并]

其中:

  • Spliterator 描述如何遍历数据,并可通过 trySplit() 拆出另一部分;
  • 拆分任务通常继续递归,直到任务足够小;
  • 每个叶子任务执行中间操作;
  • 终端操作产生局部结果;
  • 框架再把局部结果合并成最终结果。

并行流是否适合某个数据源,首先取决于它能否有效拆分。数组、ArrayList 等数据源通常容易按索引拆分;链表、逐个读取的输入流或拆分成本高的数据源,可能难以获得均衡并行。

2. 一个可运行的并行流示例

下面的程序分别使用顺序流和并行流计算平方和。它不把运行时间写死,因为实际结果会受到处理器、JIT 编译、数据规模和系统负载影响。

import java.util.stream.LongStream;

public class ParallelSum {
    public static void main(String[] args) {
        long n = 10_000_000L;

        long sequential = LongStream.rangeClosed(1, n)
                .map(x -> x * x)
                .sum();

        long parallel = LongStream.rangeClosed(1, n)
                .parallel()
                .map(x -> x * x)
                .sum();

        System.out.println(sequential);
        System.out.println(parallel);
        System.out.println(sequential == parallel);
    }
}

输入是 110_000_000 的整数。对于这个表达式,map 是无状态转换,sum 使用加法合并局部结果,满足并行归约所需的基本结构,因此顺序和并行结果应相同,最后一行应输出:

true

这里的正确性依赖于操作的性质,而不是依赖线程恰好按某种顺序执行。

3. 并行流的正确性条件

一个并行流终端操作通常需要满足以下条件:

无干扰

流执行期间不能修改数据源的结构。例如:

List<Integer> values = new ArrayList<>(List.of(1, 2, 3));

values.parallelStream()
        .forEach(values::add); // 错误:修改正在遍历的数据源

这可能抛出异常,也可能产生不可预测结果。即使换成线程安全集合,也不代表这种设计正确,因为数据源在遍历过程中被修改仍然会破坏处理语义。

无状态或状态隔离

下面的代码有竞争问题:

long[] total = {0};

LongStream.rangeClosed(1, 1_000_000)
        .parallel()
        .forEach(x -> total[0] += x);

total[0] += x 不是原子操作,多个线程会读到相同旧值并互相覆盖,最终结果可能小于正确值。

应优先使用支持并行归约的终端操作:

long total = LongStream.rangeClosed(1, 1_000_000)
        .parallel()
        .sum();

归约操作要满足结合性

并行归约会先计算局部结果,再合并局部结果。合并顺序通常不同于顺序执行,因此运算应满足:

(ab)c=a(bc)(a \mathbin{\oplus} b) \mathbin{\oplus} c = a \mathbin{\oplus} (b \mathbin{\oplus} c)

例如整数加法满足结合性,字符串拼接在固定遇到顺序下也可以保持语义;但浮点加法不严格满足结合性:

double sequential = java.util.stream.DoubleStream.of(
        1e16, 1.0, -1e16
).sum();

double parallel = java.util.stream.DoubleStream.of(
        1e16, 1.0, -1e16
).parallel().sum();

浮点舍入会导致不同的分组方式产生不同结果。并行流不承诺为了数值复现而强制使用顺序求和。

4. 顺序流、并行流和 encounter order

并行并不自动意味着结果无序。流可能具有 encounter order,即数据源定义的遇到顺序。

例如:

List<Integer> result = java.util.stream.IntStream.range(0, 20)
        .parallel()
        .boxed()
        .toList();

对于有序流,toList() 仍然应按照流的遇到顺序形成结果。为了保持结果顺序,框架可能需要额外协调,这会增加成本。

而:

java.util.stream.IntStream.range(0, 20)
        .parallel()
        .forEach(System.out::println);

不应假设打印顺序是 019。如果确实需要保持顺序,可以使用 forEachOrdered,但这通常会削弱并行执行的自由度。

四、并行流使用哪个 ForkJoinPool

1. common pool 是默认共享资源

在常见的 JDK 实现中,并行流使用 ForkJoinPool.commonPool() 执行。common pool 是 JVM 进程内共享的公共池,许多默认异步机制也可能使用它。

可以查看其目标并行度:

import java.util.concurrent.ForkJoinPool;

public class CommonPoolInfo {
    public static void main(String[] args) {
        ForkJoinPool pool = ForkJoinPool.commonPool();

        System.out.println("parallelism = "
                + pool.getParallelism());
        System.out.println("common parallelism = "
                + ForkJoinPool.getCommonPoolParallelism());
    }
}

常见实现中,common pool 的默认并行度与可用处理器数相关,通常接近:

availableProcessors() - 1

但应用程序不应把这个表达式当作 Java 规范对所有运行环境的固定承诺。容器 CPU 配额、JVM 版本、启动参数和实现策略都可能影响可用处理器数及公共池配置。

2. 不要把 common pool 当成应用专属线程池

common pool 是共享资源,因此一个组件提交大量并行流任务,可能影响同一进程中依赖该池的其他任务。

例如,某个请求处理器执行:

orders.parallelStream()
        .map(this::calculate)
        .toList();

如果 calculate 内部还执行数据库访问或 HTTP 调用,工作线程可能长时间阻塞。此时其他使用 common pool 的异步任务也可能排队。

这不是“并行流线程不够”这么简单,而是共享执行资源发生了错误的任务混合:

CPU 计算任务 ─┐
              ├── commonPool ── 工作线程被阻塞
HTTP 等待任务 ┘

3. 使用自定义 ForkJoinPool 的边界

可以显式创建 ForkJoinPool 并提交一个并行流计算:

import java.util.List;
import java.util.concurrent.ForkJoinPool;

public class CustomPoolStream {
    public static void main(String[] args) throws Exception {
        ForkJoinPool pool = new ForkJoinPool(4);

        try {
            List<Integer> result = pool.submit(() ->
                    java.util.stream.IntStream.range(0, 100)
                            .parallel()
                            .map(x -> x * x)
                            .boxed()
                            .toList()
            ).get();

            System.out.println(result.size());
            System.out.println(result.get(10));
        } finally {
            pool.shutdown();
        }
    }
}

预期输出类似:

100
100

这里的 pool.submit 负责把外层任务放入自定义池。当前 JDK 实现能够让该并行流任务使用执行它的 Fork/Join 池,但 Stream API 并没有提供“为这个并行流绑定指定 ForkJoinPool”的通用流级 API。生产代码不能把这种行为当成一个明确的 Stream 规范保证,尤其不应依赖复杂嵌套场景下的池选择细节。

如果必须严格隔离线程资源,更明确的方式通常是:

  • 直接向专用执行器提交明确的子任务;
  • 使用专用 ForkJoinPool 编写 RecursiveTask
  • 对外部 I/O 使用专门的执行模型,而不是把阻塞操作塞入并行流。

另外,创建自定义池并不意味着线程数越大越好。CPU 密集任务的并行度过高,会增加上下文切换、缓存失效和内存带宽竞争。

五、阻塞为什么会破坏工作窃取

1. 计算等待和外部阻塞的区别

Fork/Join 适合这样的依赖:

任务 A 拆成 A1、A2
A 等待 A1、A2 的计算结果

当 A 等待时,工作线程可以继续帮助完成 A1、A2 或其他 Fork/Join 任务。

但下面的等待不属于可帮助的 Fork/Join 依赖:

线程调用 read()
线程等待数据库连接
线程等待 HTTP 响应
线程等待锁
线程执行 Thread.sleep()

如果池的并行度为 4,四个工作线程都进入外部等待:

Worker-1:等待 HTTP
Worker-2:等待数据库
Worker-3:等待锁
Worker-4:sleep

此时即使队列中还有大量计算任务,也没有空闲工作线程执行它们。工作窃取只能在“有线程可窃取”时发挥作用,不能从阻塞调用中窃取出计算能力。

2. 并行流不会自动识别阻塞

下面的代码语法正确,但执行模型可能完全不适合:

List<String> bodies = urls.parallelStream()
        .map(this::httpGet)
        .toList();

parallelStream() 只知道这是一个并行任务,不知道 httpGet 会阻塞多久,也不会自动为每个阻塞调用创建一个可替代的计算线程。

如果请求数量超过 common pool 的并行度,部分任务会等待;如果已有其他任务占满公共池,影响会进一步扩大。结果可能表现为:

  • 吞吐下降;
  • 延迟突然升高;
  • 任务在队列中长时间等待;
  • 应用看起来“线程很多”,但真正可运行的计算线程很少;
  • 多层异步任务之间形成间接饥饿。

3. ManagedBlocker 的作用

ForkJoinPool.ManagedBlocker 用于向 Fork/Join 池声明:

当前任务即将进入可能长时间阻塞的区域,池可以考虑补充或补偿工作线程,以维持目标并行度。

下面是一个完整的阻塞示例:

import java.util.concurrent.ForkJoinPool;
import java.util.concurrent.ForkJoinPool.ManagedBlocker;

public class ManagedBlockExample {

    static final class SleepBlocker implements ManagedBlocker {
        private final long millis;
        private volatile boolean done;

        SleepBlocker(long millis) {
            this.millis = millis;
        }

        @Override
        public boolean isReleasable() {
            return done;
        }

        @Override
        public boolean block() throws InterruptedException {
            if (!done) {
                try {
                    Thread.sleep(millis);
                } finally {
                    done = true;
                }
            }
            return true;
        }
    }

    static void managedSleep(long millis) {
        try {
            ForkJoinPool.managedBlock(new SleepBlocker(millis));
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("Interrupted while blocking", e);
        }
    }

    public static void main(String[] args) {
        ForkJoinPool pool = new ForkJoinPool(2);

        try {
            pool.submit(() -> {
                managedSleep(500);
                System.out.println("task-1 finished");
            });

            pool.submit(() -> {
                managedSleep(500);
                System.out.println("task-2 finished");
            });

            pool.shutdown();
            pool.awaitQuiescence(2, java.util.concurrent.TimeUnit.SECONDS);
        } finally {
            pool.shutdownNow();
        }
    }
}

isReleasable() 用于快速检查是否已经解除阻塞;block() 在确实需要等待时执行阻塞操作。managedBlock 允许池采取补偿行动,但它不保证:

  • 阻塞操作一定能快速结束;
  • 一定立即创建新的线程;
  • 一定突破操作系统、容器或资源限制;
  • 外部资源一定有足够容量;
  • 应用整体吞吐一定提高。

补偿线程本身也需要 CPU 和栈空间。如果阻塞资源是数据库连接池,增加工作线程可能只会让更多线程排队等待数据库连接。

4. ManagedBlocker 不能自动修复并行流

ManagedBlocker 需要在阻塞发生的位置显式调用。可以把某些同步等待封装进去:

String value = java.util.stream.Stream.of("a", "b", "c")
        .parallel()
        .map(key -> {
            managedSleep(100);
            return key.toUpperCase();
        })
        .findFirst()
        .orElseThrow();

但这并不意味着“任何并行流 I/O 都应该用 ManagedBlocker”。如果实际操作是 HTTP、数据库或文件 I/O,还需要处理:

  • 超时;
  • 连接池耗尽;
  • 中断;
  • 重试;
  • 失败传播;
  • 下游服务限流;
  • 取消已经提交但尚未完成的任务。

在应用层,使用专用 I/O 执行器、异步客户端或虚拟线程,通常比把阻塞 I/O 隐藏在并行流 lambda 中更容易控制。Java 25 的虚拟线程适合大量等待型任务,但它们并不会让 CPU 计算无限并行,也不会自动解决数据库连接数和远端服务容量限制。

六、一个完整的性能推导

1. 理想加速不可能超过处理器并行度

假设一个纯计算任务:

  • 单线程工作量为 T1 = 800 ms
  • 处理器并行度为 p = 8
  • 任务可以均匀拆分;
  • 任务之间无依赖;
  • 暂时忽略调度和合并成本。

则理想计算下界为:

T1p=8008=100 ms\frac{T_1}{p} = \frac{800}{8} = 100\text{ ms}

但如果任务存在一段无法并行的串行路径,例如 T∞ = 120 ms,那么:

T8max(100,120)=120 msT_8 \geq \max(100,120)=120\text{ ms}

因此,理论最大加速比不超过:

T1T=8001206.67\frac{T_1}{T_\infty} = \frac{800}{120} \approx 6.67

即使机器有 8 个并行工作线程,也不可能达到 8 倍加速。

2. 用 Amdahl 定律看串行比例

如果任务中有 s 的时间必须串行执行,剩余 1-s 可以并行,则加速比上限为:

S(p)1s+1spS(p) \leq \frac{1}{s+\frac{1-s}{p}}

令:

  • s = 0.1
  • p = 8

则:

S(8)10.1+0.98=10.21254.71S(8) \leq \frac{1}{0.1+\frac{0.9}{8}} = \frac{1}{0.2125} \approx 4.71

所以一个包含 10% 串行部分的任务,即使拥有 8 个工作线程,理论加速上限也只有约 4.71 倍。

3. 任务过小时,调度成本占主导

设每个元素只执行一次非常简单的加法:

.map(x -> x + 1)

如果单个元素的实际计算时间为 c,并行拆分、队列调度和局部结果合并的平均成本为 o,当:

coc \ll o

时,并行化的收益会被调度开销吞掉。

这也是为什么下面的代码不一定比普通循环快:

List<Integer> result = values.parallelStream()
        .map(x -> x + 1)
        .toList();

如果 values 只有几十或几百个元素,或者 map 只是一个极轻量操作,顺序执行通常更简单。这里不能根据元素数量给出统一阈值,因为阈值取决于数据源、处理函数、CPU、内存层级和 JIT 优化结果。

4. 内存带宽可能成为上限

假设每个任务只是遍历一个大数组并读取数据:

long total = java.util.Arrays.stream(array)
        .parallel()
        .asLongStream()
        .sum();

当多个核心同时读取内存时,瓶颈可能从 CPU 算术单元转移到内存带宽。此时增加并行度不会线性提升速度,甚至可能因为缓存争用和内存控制器饱和而变慢。

因此“CPU 使用率高”也不能直接推出“还应该增加并行度”。需要区分:

  • CPU 计算受限;
  • 内存带宽受限;
  • 锁竞争受限;
  • I/O 等待受限;
  • GC 或分配受限。

七、哪些数据源和操作适合并行

1. 适合拆分的数据源

并行流通常更容易从以下特征中获益:

  • 数据规模足够大;
  • 数据源可以低成本、近似均匀地拆分;
  • 每个元素的处理成本明显高于调度成本;
  • 元素之间相互独立;
  • 终端操作可以高效合并局部结果;
  • 结果不要求额外的全局同步。

数组和基于索引访问的列表通常容易拆分。相反,链式结构、外部游标或每次拆分都需要大量扫描的数据源,可能在拆分阶段就付出高成本。

2. 有状态操作会破坏并行性

例如:

List<Integer> output = new ArrayList<>();

values.parallelStream()
        .map(this::transform)
        .forEach(output::add);

ArrayList 不是并发集合,多线程写入会造成数据竞争和结果错误。

即使改为:

List<Integer> output =
        java.util.Collections.synchronizedList(new ArrayList<>());

values.parallelStream()
        .map(this::transform)
        .forEach(output::add);

也只是让单次 add 的并发访问受保护,所有线程仍然在竞争同一个共享状态,锁可能成为串行瓶颈,而且结果顺序不一定符合要求。

更适合并行流的是让每个任务产生局部结果,再由框架负责合并,例如:

List<String> output = values.parallelStream()
        .map(this::transform)
        .toList();

3. limitfindFirst 和有序约束可能降低并行收益

短路操作需要在满足条件后尽快停止,但并行任务可能已经被提交或正在运行。为了保证有序语义,框架还可能需要协调多个分区的结果。

例如:

var first = values.parallelStream()
        .filter(this::isMatch)
        .findFirst();

它的语义是找到遇到顺序中的第一个匹配项,而不是任意一个最快找到的匹配项。findAny() 在允许任意匹配结果时通常有更大的并行自由度,但仍不能保证必然更快。

八、阻塞、嵌套和池饥饿的典型失败路径

1. 池饥饿的形成过程

考虑一个并行度为 4 的池:

外层任务 1:等待内层任务 1
外层任务 2:等待内层任务 2
外层任务 3:等待内层任务 3
外层任务 4:等待内层任务 4

如果这些外层任务占满了工作线程,而内层任务由于错误的池隔离策略进入另一个无法及时调度的队列,就可能出现互相等待。

简化的依赖图是:

外层任务占用全部工作线程
        │
        ▼
等待内层任务完成
        │
        ▼
内层任务没有可运行线程
        │
        └──── 外层任务无法结束

Fork/Join 对同池内的 Join 依赖有帮助机制,但它不能保证跨执行器、跨锁、跨 I/O 资源的依赖都能被打破。

2. 嵌套并行流不是免费的二级并行

outer.parallelStream()
        .map(x -> inner.parallelStream()
                .map(this::calculate)
                .toList())
        .toList();

嵌套并行流可能导致:

  • 外层和内层争用同一个 common pool;
  • 任务数量快速膨胀;
  • 外层任务等待内层任务;
  • 额外拆分和合并;
  • 缓存局部性变差。

当前实现可能在嵌套场景中复用或感知当前 Fork/Join 执行环境,但应用不应把这种行为当成“自动形成两级独立线程池”。嵌套并行只有在经过测量并确认任务结构确实需要时才有理由存在。

3. 锁阻塞比普通 I/O 更危险

下面的代码把并行计算强行串行化:

Object lock = new Object();

values.parallelStream()
        .forEach(value -> {
            synchronized (lock) {
                updateSharedState(value);
            }
        });

如果 updateSharedState 很快,锁竞争可能已经抵消并行计算的收益;如果它内部还执行 I/O,则工作线程会在持有锁时阻塞,后果更严重:

  1. 一个线程持锁并等待 I/O;
  2. 其他工作线程全部等待同一把锁;
  3. 池中没有可执行的计算任务;
  4. 并行度退化为接近 1,甚至出现级联等待。

九、错误处理、取消和资源生命周期

1. 并行流中的异常传播

List<Integer> result = values.parallelStream()
        .map(value -> {
            if (value < 0) {
                throw new IllegalArgumentException("negative value");
            }
            return value * 2;
        })
        .toList();

如果某个任务抛出异常,终端操作通常会向调用方传播异常。但其他已经开始运行的任务不一定能瞬间停止,已经产生的副作用也不会自动回滚。

因此,lambda 中不应进行不可逆的外部副作用,例如:

values.parallelStream()
        .forEach(this::sendPayment);

如果中途失败,部分支付可能已经成功,而流本身只会报告一个异常。需要事务、幂等键、补偿操作或明确的批处理协议,而不是把并行流当成事务协调器。

2. 自定义 ForkJoinPool 必须管理生命周期

创建自定义池后,应明确关闭:

ForkJoinPool pool = new ForkJoinPool(4);
try {
    var future = pool.submit(() -> doParallelWork());
    var result = future.join();
} finally {
    pool.shutdown();
}

shutdown() 表示不再接受新任务;已有任务会继续执行。需要尽快停止时可以使用 shutdownNow(),但它主要通过中断尝试停止任务,任务是否响应中断取决于任务自身。

对于阻塞任务,必须检查中断并恢复中断标记:

try {
    blockingCall();
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    throw new RuntimeException(e);
}

吞掉中断异常会破坏上层取消和关闭逻辑。

3. parallelStream() 不提供统一的超时控制

并行流本身没有一个类似“整个流水线最多执行 2 秒”的终端操作参数。把并行流包在 Future.get(timeout, unit) 外面只能限制调用方等待时间,不能保证已经运行的元素处理立即停止:

Future<List<Result>> future = executor.submit(() ->
        values.parallelStream()
                .map(this::process)
                .toList()
);

try {
    return future.get(2, java.util.concurrent.TimeUnit.SECONDS);
} catch (java.util.concurrent.TimeoutException e) {
    future.cancel(true);
    throw e;
}

这里的取消是否有效,取决于 process 是否检查中断,以及底层 I/O 是否支持取消。超时不是自动回滚机制。

十、如何诊断并行流和 ForkJoinPool 的问题

1. 先区分“排队”与“阻塞”

可以在自定义池上观察一些运行时指标:

ForkJoinPool pool = new ForkJoinPool(4);

System.out.println(pool.getParallelism());
System.out.println(pool.getActiveThreadCount());
System.out.println(pool.getQueuedTaskCount());
System.out.println(pool.getQueuedSubmissionCount());
System.out.println(pool.getStealCount());

这些指标只能提供线索:

  • getActiveThreadCount() 较低但任务队列很长,可能存在阻塞或任务未被及时调度;
  • 队列很短但运行时间很长,可能是单个任务太大、I/O 慢或锁竞争;
  • getStealCount() 持续增长说明发生了窃取,但不代表一定获得了良好性能;
  • 提交任务数很多而工作线程活跃度低,可能是外部资源或锁成为瓶颈。

这些值是近似监控信息,不应当被当成精确的事务计数器。

2. 观察线程栈

线程转储中常见的阻塞位置包括:

java.lang.Thread.State: WAITING
    at java.util.concurrent.ForkJoinPool.awaitWork(...)

这通常意味着线程当前没有可执行工作,未必是故障。

更值得关注的是:

WAITING / TIMED_WAITING
    at java.net.SocketInputStream...
    at java.sql...
    at java.lang.Object.wait(...)
    at java.util.concurrent.locks.LockSupport.park(...)

如果大量 Fork/Join 工作线程同时停在网络、数据库、锁或睡眠调用上,问题通常不是“窃取算法不够好”,而是把阻塞型任务放进了不合适的池。

3. 用基准测试而不是单次 System.nanoTime

并行性能容易受到以下因素影响:

  • JIT 预热;
  • CPU 频率变化;
  • 垃圾回收;
  • 操作系统调度;
  • 数据缓存状态;
  • 其他进程负载;
  • 第一次拆分和类初始化成本。

比较顺序流和并行流时,至少要保证:

  1. 使用相同的数据;
  2. 结果被消费,避免无效计算被优化或逻辑遗漏;
  3. 进行预热;
  4. 多次测量;
  5. 观察不同数据规模;
  6. 分别测试 CPU 密集和阻塞场景。

生产级基准测试应优先使用 JMH,而不是依赖一次程序运行的墙上时钟时间。

十一、规范保证、实现行为和工程建议要分开

规范或 API 层面的保证

  • ForkJoinPool 提供 Fork/Join 任务执行、提交、等待和工作窃取相关能力。
  • ManagedBlocker 是用于配合可能阻塞操作的 API。
  • 流操作必须遵守其定义的顺序、归约和非干扰语义。
  • 并行流不允许通过数据竞争和共享可变状态来获得正确结果。
  • 有序流的顺序语义不能因为并行执行而被任意破坏。

常见 OpenJDK 实现行为

  • 并行流通常使用 common pool;
  • Fork/Join 工作线程使用局部双端队列并支持窃取;
  • Fork/Join 任务在 Join 等待期间可能帮助执行其他任务;
  • 自定义池外层提交并行流,当前实现通常会让计算在该池中运行。

最后一类行为不应被误写成 Stream API 的通用“指定线程池”能力。

工程取舍

  • CPU 密集、可拆分、无共享可变状态的任务,可以考虑并行流或专用 Fork/Join 任务。
  • 大量阻塞 I/O 不应直接塞进 common pool。
  • 需要隔离资源、控制队列、配置线程名、限制并发或实现超时时,显式执行器通常更容易管理。
  • 需要精确的任务依赖、递归拆分和合并策略时,RecursiveTask 比并行流更直接。
  • 需要处理大量等待型独立任务时,可以评估虚拟线程或异步 I/O,但仍需单独限制数据库、连接池和远端服务的并发容量。

十二、最终边界

ForkJoinPool 的性能来自三个条件共同成立:

收益足够的可拆分工作+较低的共享协调成本+工作线程保持可运行\text{收益} \approx \text{足够的可拆分工作} + \text{较低的共享协调成本} + \text{工作线程保持可运行}

工作窃取解决的是任务分布不均:某些线程空闲时,去其他线程的队列中获取任务。它不解决不可拆分的大任务,不消除内存带宽上限,也不会让阻塞的线程继续执行计算。

并行流解决的是数据处理管道的任务化和结果合并。它要求操作满足无干扰、适当的无状态性和归约结合性;它也不负责事务回滚、外部 I/O 限流、超时管理或共享池隔离。

阻塞则改变了整个模型:线程不再执行可窃取的任务,而是在等待外部条件。ManagedBlocker 可以帮助 Fork/Join 池识别这类等待并尝试补偿,但补偿不是无限扩容,且不能替代对外部资源容量、取消、超时和故障传播的设计。

因此,判断是否使用并行流或 ForkJoinPool 的正确问题不是“数据量大不大”,而是:

  1. 工作能否低成本、均匀拆分;
  2. 单个任务是否足够重;
  3. 任务是否主要消耗 CPU;
  4. 结果是否可以安全、高效地合并;
  5. 是否会占用共享 common pool;
  6. 阻塞、取消、异常和资源生命周期是否可控。

只有这些条件经过分析或测量后成立,并行才可能把额外的调度复杂度转化为实际吞吐,而不是把一个简单的顺序问题变成更难诊断的并发问题。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。