CodeGym /课程 /JAVA 25 SELF /Spliterator 与并行流

Spliterator 与并行流

JAVA 25 SELF
第 33 级 , 课程 3
可用

1. 认识 Spliterator

如果你认为在 Java 中集合只能通过 Iterator 来遍历,那么在 Java 8 之前你完全正确。但随着 Stream API 的到来以及并行化的流行,出现了一个新角色 — Spliterator

Spliterator 是一个接口,它不仅允许遍历集合元素,还能拆分数据源以进行并行处理。名称来自 splititerator 的合成。

想象一个大蛋糕。普通的 Iterator 一块块按顺序吃。Spliterator 可以把蛋糕切成两半,把一半给朋友——你们俩可以同时开吃。朋友多就继续拆分!

Spliterator 接口 — 核心方法

public interface Spliterator<T> {
    boolean tryAdvance(java.util.function.Consumer<? super T> action);
    Spliterator<T> trySplit();
    long estimateSize();
    int characteristics();
    // ... 还有几个方法,但这些 — 最重要
}
  • tryAdvance — 对下一个元素执行操作(相当于 next() + 动作)。
  • trySplit — 尝试把源拆成两部分,并返回一个用于“被分离部分”的新 Spliterator
  • estimateSize — 估算还剩多少元素。
  • characteristics — 返回特性的位掩码(有序、唯一、不可变等)。

2. 使用 Spliterator:手动遍历与拆分

从集合获取 Spliterator

任何实现了 Collection 的集合都可以提供自己的 Spliterator

import java.util.List;
import java.util.Spliterator;

List<String> names = List.of("Vasya", "Petya", "Masha", "Lena");
Spliterator<String> spliterator = names.spliterator();

手动遍历元素

Spliterator<String> spliterator = names.spliterator();
while (spliterator.tryAdvance(name -> System.out.println("Name: " + name))) {
    // 一切都在 tryAdvance 内部完成
}

拆分集合

最有意思的是方法 trySplit()

Spliterator<String> spliterator1 = names.spliterator();
Spliterator<String> spliterator2 = spliterator1.trySplit();

System.out.println("第一部分:");
spliterator1.forEachRemaining(System.out::println);

System.out.println("第二部分:");
if (spliterator2 != null) {
    spliterator2.forEachRemaining(System.out::println);
}

会发生什么:Spliterator 会尝试将集合拆成两部分(不一定严格对半——取决于实现)。现在你可以独立处理两部分——甚至在不同的线程中!

3. 并行流:为什么以及如何工作

并行流(parallelStream())是在多个线程中同时处理元素的流。对数据量大且多核处理器的场景尤其有用。

import java.util.List;

List<String> names = List.of("Vasya", "Petya", "Masha", "Lena");
// 普通流:
names.stream().forEach(System.out::println);
// 并行流:
names.parallelStream().forEach(System.out::println);

关键点是什么?
在普通流中,元素在一个线程里顺序处理。在并行流中——数据源被拆分成多段(借助 Spliterator),每段在单独的线程中处理。

内部是如何工作的?

  1. Spliterator 将集合拆分成多段——通常按可用内核数(或略多)来划分。
  2. 每一段在自己的线程中处理——使用公共的 ForkJoinPool
  3. 结果被汇总——合并为最终的集合或数值。

并行流的工作示意

flowchart LR
    A[集合] --> B{Spliterator}
    B --> C1[部分 1] --> D1[线程 1]
    B --> C2[部分 2] --> D2[线程 2]
    B --> C3[部分 3] --> D3[线程 3]
    D1 & D2 & D3 --> E[结果汇总]

4. 并行流的优势与限制

优势

  • 加速处理大集合:在重计算场景下,并行流能显著缩短执行时间。
  • 简单:无需手写多线程代码——把 stream() 换成 parallelStream() 即可。

限制与坑点

  • 不一定更快:对小集合,额外开销可能会“吃掉”收益。
  • 顺序不保证:在 forEach/map/filter 中顺序可能不同。若需要顺序——使用 forEachOrdered
  • 线程安全问题:带副作用的操作(修改外部集合/变量)会导致数据竞争。
  • 并非所有操作都适用:存在依赖的计算(例如顺序累加)可能不会按预期工作。

何时使用并行流?

  • 大集合(数万元素及以上)。
  • 每个元素上的操作很重。
  • 对严格顺序不敏感。
  • 无副作用(纯函数)。

何时不要使用?

  • 元素很少。
  • 代码会修改外部变量或集合。
  • 需要保持处理顺序。
  • 数据源不易拆分(例如 LinkedList)。

5. 实践示例

示例 1:执行时间对比

import java.util.*;
import java.util.stream.*;

public class ParallelStreamDemo {
    public static void main(String[] args) {
        List<Integer> numbers = IntStream.range(0, 10_000_000)
                                         .boxed()
                                         .collect(Collectors.toList());

        long start = System.currentTimeMillis();
        long count = numbers.stream()
                .filter(n -> isPrime(n))
                .count();
        long time = System.currentTimeMillis() - start;
        System.out.println("普通流: " + time + " 毫秒,找到的质数: " + count);

        start = System.currentTimeMillis();
        count = numbers.parallelStream()
                .filter(n -> isPrime(n))
                .count();
        time = System.currentTimeMillis() - start;
        System.out.println("并行流: " + time + " 毫秒,找到的质数: " + count);
    }

    // 最简单的素数检查(示例用)
    public static boolean isPrime(int n) {
        if (n < 2) return false;
        for (int i = 2, sqrt = (int)Math.sqrt(n); i <= sqrt; i++)
            if (n % i == 0) return false;
        return true;
    }
}

会得到什么:在大数据量上,并行流通常更快(尤其在多核处理器上)。在小数据量上——可能没有差异,甚至并行版本更慢。

示例 2:顺序问题

import java.util.List;

List<String> names = List.of("Vasya", "Petya", "Masha", "Lena");
System.out.println("普通流:");
names.stream().forEach(System.out::println);

System.out.println("并行流:");
names.parallelStream().forEach(System.out::println);

System.out.println("使用 forEachOrdered 的并行流:");
names.parallelStream().forEachOrdered(System.out::println);

结论:在普通流和使用 forEachOrdered 的情况下顺序会被保留,而未使用它的并行流则不保证顺序。

示例 3:副作用的危险

import java.util.*;
import java.util.stream.*;

List<Integer> numbers = IntStream.range(1, 1000).boxed().collect(Collectors.toList());
List<Integer> results = new ArrayList<>();

// 危险!不要这样做!
numbers.parallelStream().forEach(n -> results.add(n * n));

System.out.println("列表大小: " + results.size());

可能发生什么?列表大小可能小于预期,有时还会出现 ConcurrentModificationException。原因是 ArrayList 不是线程安全的,而并行流会同时启动多个线程。

6. Spliterator:特性与特点

Spliterator 的特性

Spliterator 通过位掩码描述自身属性:

  • ORDERED — 元素具有确定顺序(例如列表)。
  • DISTINCT — 所有元素唯一(例如集合)。
  • SORTED — 元素已排序。
  • SIZED — 已知大小。
  • IMMUTABLE — 集合不可变。
  • CONCURRENT — 集合是线程安全的。
  • SUBSIZED — 经过 trySplit() 后的所有 spliterator 也知道自己的大小。
Spliterator<String> spliterator = names.spliterator();
int characteristics = spliterator.characteristics();
System.out.println(Integer.toBinaryString(characteristics));

为什么要了解这些?Stream API 和并行流会利用这些标志进行优化。比如,如果源是不可变且已排序,就能更安全、更高效地拆分并汇总结果。

7. 何时以及如何直接使用 Spliterator?

在日常开发中很少需要自己编写 Spliterator:标准集合已经实现好了。但如果你创建自己的数据源,或希望精细控制遍历/拆分,Spliterator 就会派上用场。

示例:使用 tryAdvance 手动遍历

import java.util.List;
import java.util.Spliterator;

List<String> names = List.of("Vasya", "Petya", "Masha", "Lena");
Spliterator<String> spliterator = names.spliterator();
spliterator.tryAdvance(name -> System.out.println("第一个元素: " + name));
spliterator.forEachRemaining(name -> System.out.println("其余: " + name));

示例:拆分集合

Spliterator<String> spliterator1 = names.spliterator();
Spliterator<String> spliterator2 = spliterator1.trySplit();

if (spliterator2 != null) {
    spliterator2.forEachRemaining(name -> System.out.println("第 2 部分: " + name));
}
spliterator1.forEachRemaining(name -> System.out.println("第 1 部分: " + name));

8. 使用 Spliterator 和并行流的常见错误

错误 1:将并行流用于小集合。 与其加速,你会得到变慢——拆分与任务调度的开销会超过收益。

错误 2:期望保持元素顺序。 并行流不保证顺序。如果顺序重要——使用 forEachOrdered,但会损失部分并行效率。

错误 3:在 lambda 表达式中产生副作用。 在并行流内部无法安全地修改外部变量/集合——会出现数据竞争和难以复现的 bug。

错误 4:在并行流中使用不安全的集合。 从多个线程往普通 ArrayList 中添加——很容易出现诸如 ConcurrentModificationException 的错误。

错误 5:期望立竿见影的加速。 并行流不是魔法。请进行性能分析:如果数据量小或操作很轻——普通流更快。

错误 6:将并行流用于不易拆分的源。 例如,LinkedList 往往拆分效率低——并行反而可能变慢。

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