1. 回顾 Stream API
你已经熟悉了 Stream API——这是一种处理集合的便捷方式,它允许你编写简洁且易读的数据处理代码:过滤、排序、计数等。
来看一个经典示例:
List<Integer> numbers = Arrays.asList(1, 2, 3, 4, 5);
int sum = numbers.stream()
.filter(n -> n % 2 == 0)
.mapToInt(n -> n)
.sum();
System.out.println(sum); // 6 (2 + 4)
在这个例子中,集合被转换为流(stream()),从中筛选出偶数,然后将其转换为 int,并通过调用 sum() 汇总结果。
Stream API 让代码更短、更具表达力:与其逐步描述如何实现,不如直接声明你想要什么。此外,如有需要,只需一行就能轻松切换到并行处理。
2. 并行流:语法与工作原理
如何让流并行?
很简单:把 stream() 换成 parallelStream()。或者对已有的流调用 .parallel()。
List<Integer> numbers = ...;
int sum = numbers.parallelStream()
.filter(n -> n % 2 == 0)
.mapToInt(n -> n)
.sum();
或者这样:
numbers.stream()
.parallel() // 转换为并行流
.filter(...)
.map(...)
.sum();
底层发生了什么?
- 集合会被自动拆分为若干部分。
- 每一部分在单独的线程中处理(使用 ForkJoinPool——专用的线程池)。
- 结果合并为一个最终值。
也就是说,如果你有多核处理器,处理会真正并行进行——例如,过滤和求和可以在多个内核上同时执行。
在哪些场景特别有用?
- 处理大集合(数万甚至更多元素)。
- 每个元素的计算较复杂。
- 不需要严格保持处理顺序。
示例:顺序流 vs 并行流
我们来看一个处理大数组的简单示例。
import java.util.*;
import java.util.stream.*;
public class ParallelStreamDemo {
public static void main(String[] args) {
List<Integer> numbers = IntStream.rangeClosed(1, 10_000_000)
.boxed()
.collect(Collectors.toList());
// 顺序流
long time1 = System.currentTimeMillis();
long count1 = numbers.stream()
.filter(n -> n % 2 == 0)
.count();
long time2 = System.currentTimeMillis();
System.out.println("顺序:" + (time2 - time1) + " 毫秒,偶数:" + count1);
// 并行流
long time3 = System.currentTimeMillis();
long count2 = numbers.parallelStream()
.filter(n -> n % 2 == 0)
.count();
long time4 = System.currentTimeMillis();
System.out.println("并行:" + (time4 - time3) + " 毫秒,偶数:" + count2);
}
}
在你的计算机上试试这段代码——很可能并行流会更快(尤其在多核处理器上)。但并非总是如此!细节见下文。
3. 工作机制:ForkJoinPool 与自动拆分
并行流在底层使用 ForkJoinPool.commonPool(),它会自动管理线程数量(通常等于可用处理器内核数)。
示意:
+-----------------------------+
| 你的集合 |
+-----------------------------+
| 1 | 2 | 3 | ... | 1000万 |
+----+----+----+-----+--------+
| | | |
v v v v
[线程1][线程2]...[线程N]
| | | |
+----+----+-----------+
|
[结果合并]
每个线程处理自己的那一部分,然后合并结果。
4. 限制与陷阱
并行流不是一键“全部加速”的魔法。 有时它甚至会让执行变慢!
不适合并行的情况:
- 集合很小(约 1000 个元素以内)。
- 对每个元素的操作非常快(例如只是 n * 2)。
- 你需要严格的处理顺序(例如顺序写入文件)。
为什么? 创建与同步线程本身也需要时间。如果任务本身很“轻”,并行化的开销可能超过收益。
副作用——并行的头号敌人
如果你的流操作会修改外部变量,请务必小心!
反例:
List<Integer> numbers = Arrays.asList(1, 2, 3, 4, 5);
int[] sum = {0};
numbers.parallelStream().forEach(n -> sum[0] += n);
System.out.println(sum[0]); // ???(你期望 15,但可能得到任意结果)
为什么?因为多个线程同时修改同一个变量——出现了 race condition(竞争状态)。最终结果可能不正确。
正确做法——使用返回结果的流方法:
int sum = numbers.parallelStream().mapToInt(n -> n).sum();
并不是所有集合都同样适合并行
有些集合(例如普通的 ArrayList)很容易拆分。而 LinkedList 或无限流(例如 Stream.generate(...))就不太适合。
5. 实践:性能对比
示例:查找最大值
import java.util.*;
import java.util.stream.*;
public class ParallelMaxDemo {
public static void main(String[] args) {
List<Integer> numbers = IntStream.rangeClosed(1, 30_000_000)
.boxed()
.collect(Collectors.toList());
// 顺序
long t1 = System.currentTimeMillis();
int max1 = numbers.stream().max(Integer::compareTo).get();
long t2 = System.currentTimeMillis();
System.out.println("顺序:" + (t2 - t1) + " 毫秒,max = " + max1);
// 并行
long t3 = System.currentTimeMillis();
int max2 = numbers.parallelStream().max(Integer::compareTo).get();
long t4 = System.currentTimeMillis();
System.out.println("并行:" + (t4 - t3) + " 毫秒,max = " + max2);
}
}
会看到什么? 在现代多核处理器上,并行流通常更快。但如果把 30_000_000 改为 1000,差别可能不明显——有时并行甚至更慢!
6. 使用示例:过滤、聚合、排序
过滤与计数
List<String> names = Arrays.asList("Anya", "Boris", "Basil", "Gregory", "Daria", "Egor", "Eugene");
long count = names.parallelStream()
.filter(name -> name.length() == 4)
.count();
System.out.println("长度为 4 的名字数量: " + count);
分组
List<String> words = Arrays.asList("cat", "whale", "cat", "dog", "whale", "cat");
Map<String, Long> freq = words.parallelStream()
.collect(Collectors.groupingBy(
w -> w,
Collectors.counting()
));
System.out.println(freq); // {dog=1, whale=2, cat=3}
排序(但在这里并行不一定带来提升!)
List<Integer> bigList = IntStream.rangeClosed(1, 5_000_000)
.boxed()
.collect(Collectors.toList());
long t1 = System.currentTimeMillis();
List<Integer> sorted = bigList.parallelStream()
.sorted()
.collect(Collectors.toList());
long t2 = System.currentTimeMillis();
System.out.println("并行排序耗时:" + (t2 - t1) + " 毫秒");
7. 重要细节与建议
何时应该使用 parallelStream()
- 集合很大(数万元素及以上)。
- 对元素的操作“很重”(复杂计算、文件/网络 IO)。
- 不依赖元素顺序。
- 没有副作用(不修改外部变量)。
何时不应该使用 parallelStream()
- 集合很小。
- 操作很快。
- 需要严格保持顺序。
- 存在对共享变量的访问(请考虑线程安全集合或其他方法)。
如何知道用了多少线程?
默认等于处理器内核数:Runtime.getRuntime().availableProcessors()。可以通过系统属性修改该行为:
System.setProperty("java.util.concurrent.ForkJoinPool.common.parallelism", "8");
只有在理解其影响的情况下再修改——否则可能“打满”CPU,反而变慢。
8. 使用并行流时的常见错误
错误一:在 forEach 中产生副作用
很多人会想:“我现在要并行往列表里塞数据!”
List<Integer> result = new ArrayList<>();
IntStream.range(0, 1_000)
.parallel()
.forEach(result::add); // 危险!
System.out.println(result.size()); // 结果是随机的!
为什么不好? ArrayList 不是线程安全的,同时从多个线程添加会导致不可预测的结果:可能丢失、重复、抛异常。
解决方案: 使用流的收集方法(collect),它们自身能保证安全,或使用专用集合。
List<Integer> result = IntStream.range(0, 1_000)
.parallel()
.boxed()
.collect(Collectors.toList());
错误二:在小任务上期待提速
并行不是免费的!当集合很小时,并行流可能因为调度与同步的开销而更慢。
错误三:破坏顺序
如果你很在意元素顺序(例如写文件),不要使用并行流——要么无法保证顺序,要么会更慢。
错误四:使用“不合适”的集合
某些集合(例如 LinkedList、非常规结构)不易拆分——并行效率会降低。
错误五:收集结果时忽视线程安全
如果你手工收集结果(例如往列表中 add),请使用线程安全集合(CopyOnWriteArrayList、ConcurrentLinkedQueue)或使用流的收集方法。
GO TO FULL VERSION