跳到主要内容

并行流的执行模型与陷阱

并行 Stream 会把数据拆成多个分区,在 Fork/Join 任务中处理并合并结果。它适合可拆分、计算量足够且没有共享副作用的数据运算;在请求线程中加入阻塞 I/O,往往会增加资源竞争和尾延迟。

1. 从 Spliterator 开始拆分数据

Stream 通过 Spliterator 遍历和拆分数据源。并行流水线会尝试调用 trySplit(),把一个数据范围分成更小的范围:

[0 ................ 15]
/ \
[0 ...... 7] [8 ...... 15]
/ \ / \
[0..3] [4..7] [8..11] [12..15]

数组和 ArrayList 可以按索引近似均匀拆分。链表、迭代器、I/O 数据源或分布不均的数据,拆分成本可能更高,任务也更难均衡。

Spliterator 的特征会告诉 Stream 数据是否有顺序、大小是否已知、元素是否互异等。执行器可以据此选择拆分和合并策略,但不会凭空消除数据源本身的限制。

2. ForkJoinPool 执行分治任务

并行流通常使用 ForkJoinPool.commonPool() 执行任务。Fork/Join 的核心流程是:

  1. 大任务递归拆成较小任务。
  2. Worker 优先处理自己的任务队列。
  3. 空闲 Worker 从其他队列窃取任务。
  4. 子任务结果按终止操作的规则合并。

工作窃取可以缓解任务耗时不均,但不能修复无法拆分的数据源,也不能解决一个子任务长期阻塞的问题。

long total = orders.parallelStream()
.filter(Order::confirmed)
.mapToLong(Order::amountInCents)
.sum();

这段代码适合并行的前提是:数据量足够、过滤和金额计算无共享副作用、计算成本能够覆盖拆分调度开销,而且当前进程还有可用 CPU。

3. 并行不会自动保持执行顺序

有遇到顺序的数据源,不等于每个操作都按这个顺序执行:

List.of(1, 2, 3, 4)
.parallelStream()
.forEach(System.out::println);

输出顺序不确定。forEachOrdered 可以按遇到顺序提交结果,但会增加协调并限制并行空间。

有些终止操作在语义上保留顺序,例如有序流的 findFirstfindAny 则允许返回任意分区的结果。业务不要求顺序时,可以调用 unordered() 放宽约束,但必须确认结果语义确实不依赖顺序。

4. 外部可变状态会产生数据竞争

下面的代码不安全:

List<String> names = new ArrayList<>();

users.parallelStream()
.map(User::displayName)
.forEach(names::add);

多个任务会同时修改 ArrayList。结果可能丢失、损坏或抛出异常。换成同步列表虽然能避免容器结构损坏,却会让所有写入竞争同一把锁:

List<String> names = users.parallelStream()
.map(User::displayName)
.toList();

让 Collector 为分区维护独立结果,再合并,是更符合并行归约模型的做法。

Lambda 还应避免依赖:

  • 请求线程中的 ThreadLocal 上下文。
  • 绑定当前线程的事务或安全上下文。
  • 非线程安全客户端和格式化器。
  • 需要严格调用顺序的外部系统。

并行任务可能在主线程和不同 Worker 上执行,线程本地状态不会自动传播。

5. 阻塞 I/O 会占住共享 Worker

List<Profile> profiles = userIds.parallelStream()
.map(profileClient::load)
.toList();

如果 load 发起阻塞网络请求,每个等待中的任务都会占用 Worker。其他并行 Stream 和使用公共 ForkJoinPool 的任务可能因此排队。请求量增大后,单次处理看起来更并行,整个进程的尾延迟却可能更差。

需要并发 I/O 时,应该显式控制:

  • 最大并发数。
  • 连接池容量。
  • 每次调用与端到端超时。
  • 取消和异常传播。
  • 不同请求或业务的资源隔离。

虚拟线程、带界限的执行器或异步客户端都比“把 Stream 改成 parallel”更容易表达这些约束。具体选择取决于调用模型,不是统一替代关系。

6. commonPool 会形成进程级竞争

公共 ForkJoinPool 是 JVM 进程内的共享资源。一个模块提交大量长任务,会影响其他使用者。只调整 java.util.concurrent.ForkJoinPool.common.parallelism 也是进程级配置,无法为每个接口建立独立预算。

如果业务必须隔离执行资源,使用显式执行器和明确的任务编排。把并行 Stream 包在自建 ForkJoinPool.submit 中的行为不应作为跨实现的资源隔离契约;即使当前 JDK 达到预期,也要通过升级测试验证。

服务端使用前先确认

并行流没有队列长度、拒绝策略、业务优先级和请求级并发上限这些接口。需要过载保护与隔离时,线程池或任务框架应当由业务显式拥有。

7. 什么情况下可能获得收益

并行流更可能适合:

  • 数据规模足够大,拆分近似均匀。
  • 每个元素执行独立的 CPU 计算。
  • 没有外部可变状态和线程绑定上下文。
  • 合并操作满足结合律,合并成本较低。
  • 机器有空闲核心,应用允许这项计算使用它们。

不利因素包括:

  • 元素很少或单个操作过轻,调度成本占主导。
  • 数据源难以拆分,例如链式结构或串行 I/O。
  • 大量装箱、分配和全局排序。
  • 有序短路操作限制任务协作。
  • 运行环境已经 CPU 饱和。
  • 流水线包含数据库、网络或文件阻塞。

性能判断需要在接近生产的数据规模、核心数和并发背景下用 JMH 或端到端压测验证。只测一次 System.currentTimeMillis() 会混入类加载、JIT 编译和 GC。

8. 异常与取消

并行流水线中的某个任务抛出异常时,终止操作会失败,但已经开始的其他任务不一定立刻停止。Stream 没有提供请求级超时参数,也没有一套面向远程调用的补偿协议。

需要超时、部分成功、逐项错误记录或主动取消时,应使用显式任务模型。不要在 Lambda 内统一吞掉异常并返回 null,这会让失败混入正常数据,还可能在后续归约阶段产生新的异常。

9. 常见问题

9.1 parallelStream 使用多少线程

通常受公共 ForkJoinPool 的并行度影响,调用线程也可能参与计算。并行度不是“同时执行的业务请求数”,阻塞、嵌套任务和运行环境都会影响实际活动线程。不要据此直接推导数据库或远程接口的并发上限。

9.2 并行流里使用线程安全集合就安全了吗

只能说明单次容器操作满足它的线程安全契约。检查后写入、跨多个容器更新、顺序要求和业务不变量仍可能出错;锁竞争也可能抵消并行收益。优先使用无副作用操作和并行 Collector。

9.3 parallel() 会创建一个新的线程池吗

不会。它只把 Stream 的并行标志设为 true,终止操作执行时再按并行模型处理。它也不会异步返回,调用终止操作的线程仍然等待最终结果。

9.4 嵌套 parallelStream 会更快吗

通常不会按层数增加并行能力。内外层可能竞争同一公共池,增加拆分、排队和合并成本。应先选择一个清晰的并行维度,并用基准验证。

10. 面试题

10.1 ForkJoinPool 怎样调度任务,工作窃取解决了什么问题

出现公司:阿里巴巴、去哪儿

考察重点

  • 分治、双端队列与工作窃取的协作关系。
  • 数据源拆分质量为什么影响并行流。
  • 阻塞任务为什么会降低 Fork/Join 的利用率。

相关内容:第 1 节“从 Spliterator 开始拆分数据”、第 2 节“ForkJoinPool 执行分治任务”。

参考回答

Fork/Join 把大任务递归拆成小任务,Worker 主要从自己的双端队列处理任务;空闲 Worker 会从其他队列窃取任务,从而减少某些线程空闲而另一些线程堆积的情况。子任务完成后再按归约规则合并结果。

它适合可以继续拆分、计算量相对均衡的 CPU 任务。数据源难以拆分或任务长期阻塞时,工作窃取无法创造更多 CPU 或连接资源,公共池中的其他任务也可能受到影响。

10.2 在服务端使用 parallelStream 需要检查哪些条件

出现公司:美团、阿里巴巴

考察重点

  • 计算密集与阻塞 I/O 的区别。
  • commonPool 共享、副作用、顺序与线程上下文。
  • 怎样用基准和容量约束判断是否采用并行流。

相关内容:第 4 节“外部可变状态会产生数据竞争”、第 5 节“阻塞 I/O 会占住共享 Worker”、第 7 节“什么情况下可能获得收益”。

参考回答

先确认流水线没有共享副作用,归约满足结合律,结果不依赖隐含线程上下文。然后判断数据源是否容易拆分、数据量与单项计算是否足以覆盖任务调度成本,以及进程是否有可用 CPU。

服务端还要检查公共 ForkJoinPool 的共享影响。数据库或网络阻塞不适合直接放进并行流,因为缺少业务并发上限、队列和拒绝策略。最终应在真实并发背景下比较串行与并行的吞吐和尾延迟,而不是根据核心数推断一定更快。