← 返回题目列表

并行流 parallelStream 是怎么工作的?什么时候该用、什么时候会踩坑?

高频 困难 第 14 / 24 题 更新于 2026/08/03
并行流parallelStreamForkJoinPool并发

简化版

**并行流(parallelStream()stream().parallel())能让 Stream 的操作「多线程并行执行」——把数据分成多段,用多个线程同时处理,最后合并结果,从而利用多核 CPU 加速。**它底层用 ForkJoinPool(fork/join 框架的工作窃取线程池):把流的数据递归「分而治之」拆成小块,分给线程池的多个线程并行处理,再把各块结果「合并(reduce/combine)」。关键坑(比好处更需要记)① 默认用「公共 ForkJoinPool(commonPool)」,全 JVM 共享——如果并行流里有慢操作(IO、阻塞),会占满公共池、拖累其他也用它的地方;② 不是所有场景都更快——数据量小、每个元素处理很轻时,并行的「拆分 + 线程调度 + 合并」开销可能大于收益,反而更慢;③ 有状态/有副作用的操作会出错——并行时多线程同时跑,如果 lambda 里改共享变量(非线程安全)会数据错乱;④ 顺序不保证(除非用 forEachOrdered)。结论:并行流适合「数据量大 + 每个元素处理重(CPU 密集)+ 无状态无副作用 + 可拆分的数据源」;有 IO/阻塞、数据量小、有共享状态时别用。

详细版

串行流 vs 并行流

维度串行流 stream()并行流 parallelStream()
执行单线程顺序多线程并行
底层当前线程ForkJoinPool(默认公共池)
顺序保证不保证(forEachOrdered 才保证)
适合通用大数据 + CPU 密集 + 无状态
风险共享状态错乱、公共池被占、可能更慢
// 并行流:多线程并行处理
long count = list.parallelStream()
    .filter(x -> isPrime(x))     // CPU 密集计算,并行有收益
    .count();

// ❌ 坑一:并行流里做 IO/阻塞 → 占满公共池,拖累全局
urls.parallelStream().forEach(url -> httpGet(url));  // 阻塞操作,别这么用

// ❌ 坑二:有共享状态 → 数据错乱
List<Integer> result = new ArrayList<>();   // ArrayList 非线程安全
nums.parallelStream().forEach(result::add); // 多线程 add → 错乱/丢数据/异常

// ✅ 正确:用收集器(线程安全的归约)
List<Integer> result2 = nums.parallelStream()
    .filter(x -> x > 0)
    .collect(Collectors.toList());          // collect 是线程安全的归约

// 想用自己的线程池(避免公共池):
ForkJoinPool customPool = new ForkJoinPool(4);
customPool.submit(() -> list.parallelStream().forEach(...)).get();

⚠️ 并行流最大的隐藏坑是「默认共享全 JVM 的公共 ForkJoinPool(ForkJoinPool.commonPool())」——这个池的线程数默认是「CPU 核数 - 1」,且整个 JVM 里所有并行流、CompletableFuture 的默认执行都共用它。所以如果你在并行流里做「慢操作」(HTTP 请求、数据库查询、加锁等待),这些线程会被长时间占用,拖累 JVM 里其他依赖公共池的并行任务(甚至互相死锁)。这就是「并行流不要做 IO/阻塞操作」的根本原因——它的设计前提是「每个任务都是快速的 CPU 计算、很快就还回线程」。如果确实要并行做阻塞操作,应该用自己的线程池(ForkJoinPool 或普通 ExecutorService + CompletableFuture),别用并行流的公共池。

完整版教学

一、并行流是什么:分而治之

先理解并行流「多线程并行」的本质:

串行流:一个线程,从头到尾一个个处理元素
  stream().filter().map().collect()  → 单线程顺序跑

并行流:把数据分段,多线程并行处理
  parallelStream() 的执行(分而治之):
    1. 把数据源"拆分"成多个小块(Spliterator 负责拆分)
    2. 每个小块分给 ForkJoinPool 的一个线程处理
    3. 各线程并行执行 filter/map 等操作
    4. 把各块的结果"合并"(combine/reduce)成最终结果

  → 利用多核 CPU,多个线程同时干活,理论上能加速

fork/join 的思想:
  fork(分):大任务递归拆成小任务
  join(合):小任务结果合并成大结果
  → 并行流底层就是 fork/join 框架

前提:
  ① 数据源能高效拆分(数组、ArrayList 好拆;LinkedList 难拆)
  ② 操作能并行(无状态、无副作用、可结合的合并)

并行流的本质是「分而治之的多线程并行」——把数据源拆分成多个小块(Spliterator 负责)、每块分给 ForkJoinPool 的线程并行处理、各块结果合并(combine/reduce)。底层是 fork/join 框架(fork 递归拆小任务、join 合并结果),利用多核加速。前提:数据源能高效拆分(数组/ArrayList 好拆、LinkedList 难拆)、操作能并行(无状态无副作用)。理解「并行流=分而治之多线程并行:拆分(Spliterator)→各块并行处理→合并、底层 fork/join、前提是数据源好拆且操作可并行」,就理解了并行流的工作原理。

二、底层:公共 ForkJoinPool

并行流用「公共 ForkJoinPool」——这是最关键也最容易踩坑的点:

并行流的线程从哪来?
  默认用 ForkJoinPool.commonPool()(公共池)
  线程数默认 = Runtime.getRuntime().availableProcessors() - 1
    (CPU 核数 - 1,加上提交任务的主线程,正好用满核数)

★ 公共池是"全 JVM 共享"的:
  - 所有并行流共用它
  - CompletableFuture 的默认异步执行也用它
  - 整个 JVM 就这一个默认池

后果(关键):
  如果一个并行流长时间占用公共池的线程(做慢操作)
  → 其他并行流、CompletableFuture 都拿不到线程 → 被拖慢/饿死
  → 甚至互相死锁

这就是"并行流别做阻塞操作"的根本原因:
  公共池线程少(核数-1),且全局共享
  一个任务占着不放,全局受影响

并行流的线程来自「公共 ForkJoinPoolcommonPool」——线程数默认「CPU 核数 - 1」,且全 JVM 共享(所有并行流、CompletableFuture 的默认执行都用它,整个 JVM 就这一个默认池)。后果:一个并行流长时间占用公共池线程(慢操作),其他并行流、CompletableFuture 都拿不到线程被拖慢甚至死锁。这就是「并行流别做阻塞操作」的根本原因(池线程少、全局共享,一个任务占着不放全局受影响)。理解「并行流用公共 ForkJoinPool(核数-1 个线程,全 JVM 共享)、一个任务占着不放会拖累全局的并行流和 CompletableFuture、这是别做阻塞操作的根本原因」,就理解了并行流最关键的机制和坑。

三、坑一:不是都更快

并行流不一定更快——「并行开销」可能大于收益:

并行不是免费的,有额外开销:
  ① 拆分数据(Spliterator)
  ② 分发任务到线程、线程调度
  ③ 合并各块结果
  ④ 线程间协调

什么时候并行更快(收益 > 开销):
  ✓ 数据量大(几万几十万以上)
  ✓ 每个元素处理"重"(CPU 密集计算,如复杂计算、质数判断)
  → 每个任务干活多,摊薄了并行开销

什么时候并行更慢(开销 > 收益):
  ✗ 数据量小(几百几千)→ 拆分调度开销 > 处理时间
  ✗ 每个元素处理很轻(如简单加法)→ 并行开销占比大
  ✗ 数据源难拆(LinkedList、IO 流)→ 拆分本身就慢

经验法则(N × Q 判断):
  N = 数据量,Q = 每个元素处理成本
  N × Q 足够大(几万以上)才考虑并行
  小数据、轻处理 → 老老实实串行

★ 别盲目 parallelStream,多数场景串行就够且更快

并行不一定更快——并行有额外开销(拆分、线程调度、合并、协调)。并行更快:数据量大 + 每个元素处理重(CPU 密集),每个任务干活多摊薄开销。并行更慢:数据量小(拆分调度开销 > 处理)、每个元素处理很轻、数据源难拆(LinkedList/IO 流)。经验法则「N(数据量)× Q(每元素成本)足够大(几万以上)才考虑并行」。别盲目 parallelStream,多数场景串行更快。理解「并行有开销(拆分/调度/合并)、数据量大+处理重才更快、数据量小或处理轻或难拆反而更慢、N×Q 大才用并行、别盲目并行」,就避开了「以为并行一定快」的误区。

四、坑二:共享状态与副作用

并行流最危险的坑是「有状态/有副作用的操作」:

并行 = 多线程同时跑,如果操作里有"共享可变状态"→ 线程安全问题

❌ 错误示例:
  List<Integer> result = new ArrayList<>();  // 非线程安全
  nums.parallelStream().forEach(result::add);
  → 多线程同时 add → 数组越界异常、丢数据、结果错乱

❌ 改共享变量:
  int[] sum = {0};
  nums.parallelStream().forEach(n -> sum[0] += n);  // 竞态
  → sum[0] += n 非原子,多线程覆盖,结果错

✅ 正确做法:用无副作用的归约
  ① 用 collect(线程安全的归约收集):
     nums.parallelStream().collect(Collectors.toList())
  ② 用 reduce(无状态归约):
     nums.parallelStream().reduce(0, Integer::sum)
  → 这些是"函数式、无副作用"的,天生适合并行

并行流的要求(函数式纯净):
  - 操作无状态(不依赖外部可变状态)
  - 无副作用(不修改外部变量)
  - 合并操作可结合(结合律:(a+b)+c = a+(b+c))
  → 满足这些,并行才安全正确

并行流最危险的坑是「共享状态/副作用」——并行是多线程同时跑,操作里有共享可变状态就有线程安全问题:forEach(result::add) 往非线程安全的 ArrayList 加会数组越界/丢数据/错乱、forEach(n -> sum[0] += n) 改共享变量有竞态。正确做法:用无副作用的归约——collect(Collectors.toList())(线程安全归约)、reduce(0, Integer::sum)(无状态归约)。并行流要求操作无状态、无副作用、合并可结合。理解「并行流危险坑=共享状态/副作用(forEach 往非线程安全集合加会错乱、改共享变量有竞态)、正确用无副作用归约(collect/reduce)、要求无状态无副作用可结合」,就避开了并行流最危险的坑。

五、坑三:顺序与 forEach

并行流的「顺序」问题也要注意:

并行流不保证处理/输出顺序:
  parallelStream().forEach(System.out::println)
  → 多线程并行,输出顺序是乱的(每次可能不同)

想保证顺序:
  forEachOrdered:按流的原始顺序处理(但会损失部分并行性)
  parallelStream().forEachOrdered(...)

有些操作本身保证顺序:
  collect(toList()):结果 List 保持原始顺序(即使并行)
  → 因为收集器内部会按顺序合并

sorted/limit/findFirst 在并行流里的表现:
  - sorted:并行排序(有额外开销)
  - findFirst:要保证是"第一个",并行下有额外协调成本
    → 无所谓第几个就用 findAny(并行下更快)
  - limit:并行下也有协调成本

所以:
  在意顺序 → 注意用 forEachOrdered,或用 collect(保序)
  不在意"第几个" → findAny 比 findFirst 更适合并行

并行流的顺序问题:不保证处理/输出顺序parallelStream().forEach 输出是乱的)。想保序用 forEachOrdered(但损失部分并行性);有些操作本身保序(collect(toList()) 结果保持原始顺序)。并行下 findAnyfindFirst 更快(不用协调「第一个」)。理解「并行流不保证顺序(forEach 输出乱)、保序用 forEachOrdered(损失并行性)或 collect(保序)、并行下 findAny 比 findFirst 快」,就掌握了并行流的顺序问题。

六、实践建议

总结并行流的实践建议和适用判断:

用并行流的条件(都满足才用):
  ① 数据量大(几万以上,N × Q 足够大)
  ② 每个元素处理重(CPU 密集计算,不是简单操作)
  ③ 无 IO/阻塞操作(不占公共池线程)
  ④ 无状态、无副作用(不改共享变量,用 collect/reduce)
  ⑤ 数据源易拆分(数组、ArrayList;避免 LinkedList)

不该用并行流:
  ✗ 数据量小、处理轻 → 串行更快
  ✗ 有 IO/阻塞 → 用自己的线程池 + CompletableFuture
  ✗ 有共享状态 → 会错乱
  ✗ 强顺序依赖 → 并行没优势

如果要并行做阻塞操作:
  别用并行流(公共池)
  用自己的 ExecutorService + CompletableFuture
  → 隔离线程池,不影响全局

一句话:并行流是"CPU 密集 + 大数据 + 无副作用"的加速器,
  不是万能提速,多数场景串行就够,盲目并行反而慢或出 bug

用并行流的条件(都满足才用):数据量大(N×Q 足够大)、每个元素处理重(CPU 密集)、无 IO/阻塞、无状态无副作用(用 collect/reduce)、数据源易拆分(数组/ArrayList)不该用:数据量小/处理轻(串行更快)、有 IO/阻塞(用自己的线程池+CompletableFuture)、有共享状态(会错乱)、强顺序依赖。要并行做阻塞操作用自己的 ExecutorService + CompletableFuture(隔离线程池)。理解「用并行流条件:大数据+CPU 密集+无阻塞+无副作用+易拆分都满足才用;不该用:小数据/有 IO/有共享状态/强顺序;并行阻塞操作用自己的线程池不用并行流」,就掌握了并行流的实践判断。

记忆钩子:「并行流 parallelStream=多线程并行处理(分而治之:Spliterator 拆分→各块并行→合并,底层 fork/join);★用公共 ForkJoinPool(核数-1 线程,全 JVM 共享)→一个任务占着不放拖累全局(这是别做 IO/阻塞的根本原因);坑:①不是都更快(并行有开销,数据量大+CPU 密集才划算,N×Q 大才用,小数据反而慢)②共享状态/副作用会错乱(forEach 往 ArrayList 加/改共享变量有竞态,要用 collect/reduce 无副作用归约)③不保证顺序(保序用 forEachOrdered/collect,findAny 比 findFirst 快);用并行流条件:大数据+CPU 密集+无阻塞+无副作用+易拆分都满足;并行阻塞用自己的线程池+CompletableFuture」

七、常见误区与追问

  • 误区:parallelStream 一定比 stream 快。 不一定——并行有拆分/调度/合并开销;数据量小、每个元素处理轻、数据源难拆(LinkedList)时,开销大于收益反而更慢;只有数据量大 + 每个元素处理重(CPU 密集)才划算;多数场景串行更快,别盲目并行。
  • 误区:并行流里可以随便改共享变量。 不行——并行是多线程同时跑,改共享的非线程安全变量(ArrayList.add、sum += n)会数据错乱、丢数据、抛异常;要用无副作用的归约(collect(Collectors.toList())、reduce(0, Integer::sum))代替 forEach 改共享状态。
  • 误区:并行流有自己独立的线程池。 默认用全 JVM 共享的公共 ForkJoinPool(commonPool,线程数=核数-1)——所有并行流和 CompletableFuture 的默认执行都共用它;一个并行流做慢操作会占满公共池、拖累全局其他并行任务。
  • 误区:并行流里做 HTTP 请求/数据库查询能加速。 不该这么做——IO/阻塞操作会长时间占用公共池的少量线程,拖累 JVM 里其他依赖公共池的并行任务甚至死锁;并行流的设计前提是「每个任务都是快速 CPU 计算」;要并行做阻塞操作用自己的 ExecutorService + CompletableFuture。
  • 追问:并行流为什么不能做阻塞操作? 因为它默认用全 JVM 共享的公共 ForkJoinPool,线程数只有「核数-1」个;阻塞操作(IO、加锁等待)会长时间占用这些宝贵的线程,导致其他并行流、CompletableFuture 拿不到线程被拖慢或饿死甚至死锁;公共池的设计前提是每个任务快速执行、很快还回线程。
  • 追问:怎么让并行流用自己的线程池而不是公共池? 把并行流操作提交到自己创建的 ForkJoinPool:myForkJoinPool.submit(() -> list.parallelStream().forEach(...)).get()——并行流会用提交它的那个 ForkJoinPool 的线程;或者更推荐直接用 ExecutorService + CompletableFuture 手动控制并行,避免并行流的公共池问题。
  • 追问:什么样的数据源适合并行流? 能高效「均匀拆分」的数据源——数组、ArrayList(基于数组,能按索引范围快速对半分)、IntStream.range 等;不适合的:LinkedList(拆分要遍历、慢)、IO 流(无法预知大小、难拆);数据源的 Spliterator 拆分效率直接影响并行收益。

八、加强记忆

并行流(parallelStream() / stream().parallel())让 Stream 操作多线程并行执行——分而治之Spliterator 把数据源拆成小块、ForkJoinPool 的多个线程并行处理各块、再合并(combine/reduce)结果,底层是 fork/join 框架,利用多核加速。最关键的机制和坑:默认用全 JVM 共享的公共 ForkJoinPoolcommonPool,线程数=核数-1)——所有并行流和 CompletableFuture 默认执行都共用它,一个任务长时间占用(做 IO/阻塞)会拖累全局甚至死锁(这是「并行流别做阻塞操作」的根本原因)。三大坑① 不是都更快(并行有拆分/调度/合并开销,数据量小或处理轻反而更慢,N×Q 足够大才用);② 共享状态/副作用会错乱forEach 往非线程安全集合加、改共享变量有竞态,要用 collect/reduce 无副作用归约);③ 不保证顺序(保序用 forEachOrdered/collect,并行下 findAnyfindFirst 快)。用并行流的条件(都满足才用):数据量大 + 每个元素处理重(CPU 密集)+ 无 IO/阻塞 + 无状态无副作用 + 数据源易拆分(数组/ArrayList);要并行做阻塞操作用自己的 ExecutorService + CompletableFuture。一句话「并行流=分而治之多线程并行(fork/join),用公共 ForkJoinPool(核数-1,全 JVM 共享)→做阻塞操作拖累全局(别做 IO);坑:不是都更快(N×Q 大才用)、共享状态会错乱(用 collect/reduce)、不保证顺序;条件:大数据+CPU 密集+无阻塞+无副作用+易拆分,阻塞操作用自己的线程池」。