CodeGym /课程 /JAVA 25 SELF /并行流:语法与应用

并行流:语法与应用

JAVA 25 SELF
第 54 级 , 课程 2
可用

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),请使用线程安全集合(CopyOnWriteArrayListConcurrentLinkedQueue)或使用流的收集方法。

评论
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION