Java 基础体系 · 第 11/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java Stream API:惰性、Collector、并行流、性能和使用边界

Java Stream API 是对“数据处理过程”的抽象,而不是一种新的集合。集合保存数据,Stream 描述从数据源读取元素、转换元素、筛选元素并产生结果的计算过程。

一个 Stream 通常由三部分组成:

数据源 → 中间操作* → 终止操作

例如:

List<String> names = List.of("alice", "bob", "anna");

List<String> result = names.stream()
        .filter(name -> name.startsWith("a"))
        .map(String::toUpperCase)
        .toList();

System.out.println(result); // [ALICE, ANNA]

这里:

  • names 是数据源;
  • filtermap 是中间操作;
  • toList 是终止操作;
  • Stream 本身不负责保存结果,结果由终止操作产生。

Stream API 位于 java.util.stream 包中,Java 25 的相关核心抽象仍然是 StreamIntStreamLongStreamDoubleStreamCollectorSpliterator


一、先区分集合、迭代器和 Stream

理解 Stream 的边界,需要先区分三个概念。

集合保存数据

List<Integer> numbers = new ArrayList<>(List.of(1, 2, 3));
numbers.add(4);

List 的职责是保存元素,并提供按位置访问、遍历和修改等能力。集合通常允许多次遍历:

for (int n : numbers) {
    System.out.println(n);
}

for (int n : numbers) {
    System.out.println(n);
}

迭代器表示一次遍历状态

Iterator<Integer> iterator = numbers.iterator();

while (iterator.hasNext()) {
    System.out.println(iterator.next());
}

迭代器内部有当前位置。遍历结束后,再调用 next() 会失败,除非重新获得一个新的迭代器。

Stream 表示一次计算管道

Stream<Integer> stream = numbers.stream();

Stream 类似于一个更高层的、可组合的遍历计算。它不是集合,也不是可以反复消费的容器:

Stream<Integer> stream = numbers.stream();

long count = stream.count();

// IllegalStateException:同一个 Stream 不能再次消费
long sum = stream.mapToInt(Integer::intValue).sum();

如果需要再次处理集合,应该重新创建 Stream:

long count = numbers.stream().count();
long sum = numbers.stream().mapToInt(Integer::intValue).sum();

这个限制不是偶然的实现细节。Stream 可能来自文件、网络、生成器或无限序列,这些数据源未必能回退,也未必适合保存全部数据。


二、惰性求值:中间操作只是构建管道

1. 中间操作通常不会立即执行

以下代码不会打印任何内容:

List<Integer> numbers = List.of(1, 2, 3);

numbers.stream()
        .filter(n -> {
            System.out.println("filter: " + n);
            return n % 2 == 1;
        })
        .map(n -> {
            System.out.println("map: " + n);
            return n * 10;
        });

原因是 filtermap 都是中间操作。它们返回了新的 Stream 描述,但没有终止整个管道的操作。

加入终止操作后,计算才开始:

List<Integer> result = numbers.stream()
        .filter(n -> {
            System.out.println("filter: " + n);
            return n % 2 == 1;
        })
        .map(n -> {
            System.out.println("map: " + n);
            return n * 10;
        })
        .toList();

一次可能的输出顺序是:

filter: 1
map: 1
filter: 2
filter: 3
map: 3

这里可以观察到一个重要事实:Stream 管道通常按元素推进,而不是先完整执行所有 filter,再完整执行所有 map

对元素 1,它先通过 filter,然后立即进入 map;元素 2filter 排除,不会进入 map

2. 惰性带来短路

短路终止操作只需要得到足够结果,就可以停止遍历。例如:

boolean found = Stream.of(3, 5, 8, 10, 12)
        .peek(n -> System.out.println("visit: " + n))
        .anyMatch(n -> n % 2 == 0);

System.out.println(found);

典型输出:

visit: 3
visit: 5
visit: 8
true

遇到 8 后,anyMatch 已经可以确定结果为 true,后面的 1012 不需要访问。

常见短路操作包括:

  • anyMatch
  • allMatch
  • noneMatch
  • findFirst
  • findAny
  • 某些情况下的 limit

但“支持短路”不代表所有元素在所有并行执行场景下都会立刻停止。并行流中可能已经有多个任务在处理数据,取消通常只能阻止尚未开始或尚未完成的部分,不能撤销已经发生的外部副作用。

3. 有状态中间操作可能需要缓存或全局信息

并非所有中间操作都能只看当前元素:

stream.sorted()
stream.distinct()
stream.limit(n)
stream.skip(n)

例如 sorted() 必须知道足够多的输入元素,才能确定最小元素和最大元素。因此它通常需要缓存输入并进行排序。

这与下面的 filter 不同:

stream.filter(x -> x > 10)

filter 只需要检查当前元素,通常可以边读边输出。

可以将操作粗略分为:

类型 示例 是否需要全局信息
无状态 mapfilterpeek 通常不需要
有状态 sorteddistinctskip 可能需要
短路 findFirstanyMatchlimit 可能提前结束

“无状态”并不等于“没有副作用”。它主要描述操作是否需要记住之前元素。一个把数据写入共享 Listmap 仍然可能有严重副作用。


三、Stream 管道的正确性:非干扰和无状态

Stream API 的计算通常要求行为满足两个重要条件。

1. 不要在遍历过程中修改数据源

错误示例:

List<Integer> numbers = new ArrayList<>(List.of(1, 2, 3));

numbers.stream().forEach(n -> {
    if (n == 2) {
        numbers.remove(n);
    }
});

这可能抛出 ConcurrentModificationException,也可能在某些数据源和实现下表现为不可预测的遍历结果。

正确做法是让 Stream 产生新结果:

List<Integer> remaining = numbers.stream()
        .filter(n -> n != 2)
        .toList();

这称为避免对数据源的干扰。对于普通集合,遍历期间不应通过集合的结构修改方法改变集合;并发集合有额外的弱一致性或并发语义,但不能因此推断所有 Stream 操作都自动变成线程安全。

2. Lambda 应尽量是无状态函数

以下代码依赖外部可变状态:

List<Integer> result = new ArrayList<>();

numbers.parallelStream()
        .map(n -> n * 2)
        .forEach(result::add);

ArrayList 不是线程安全的,结果可能丢失元素或产生其他并发问题。

可以改成由 Stream 负责收集:

List<Integer> result = numbers.parallelStream()
        .map(n -> n * 2)
        .toList();

或者使用正确的并发容器,但这通常改变了问题性质,并且可能引入锁竞争:

List<Integer> result = new CopyOnWriteArrayList<>();
numbers.parallelStream()
        .map(n -> n * 2)
        .forEach(result::add);

更值得注意的是,forEach 的共享可变状态问题并不只存在于并行流。即使当前使用顺序流,未来把 stream() 改成 parallelStream() 也可能立即暴露问题,因此不应把“当前恰好串行”当作正确性的基础。


四、终止操作:遍历、归约和收集

终止操作触发管道执行,并产生结果或副作用。

常见终止操作可以分为三组。

遍历操作

numbers.stream().forEach(System.out::println);

forEach 不保证并行流中的遇到顺序。若必须按照流的遇到顺序处理,可以使用:

numbers.parallelStream()
        .forEachOrdered(System.out::println);

但顺序约束会限制并行执行的自由度,不能无成本地获得。

查找和匹配操作

boolean hasEven = numbers.stream().anyMatch(n -> n % 2 == 0);
Optional<Integer> first = numbers.stream().findFirst();
Optional<Integer> arbitrary = numbers.parallelStream().findAny();
  • findFirst 需要遵守有序流的第一个元素语义;
  • findAny 允许返回任意匹配元素,更适合不关心顺序的并行场景;
  • Optional 表示可能没有结果,避免用 null 表示缺失。

归约操作

归约把多个元素逐步合并为一个结果:

int sum = Stream.of(1, 2, 3, 4)
        .reduce(0, Integer::sum);

可以形式化为:

r0=identityr_0 = identity

ri+1=accumulator(ri,xi)r_{i+1} = accumulator(r_i, x_i)

其中:

  • identity 是初始值;
  • x_i 是第 i 个元素;
  • accumulator 把当前结果和一个元素合并。

对于并行执行,还需要一个合并器:

R=combiner(R1,R2)R = combiner(R_1, R_2)

要让并行结果与顺序结果一致,通常需要满足:

  1. 恒等元条件:

combiner(identity,r)=rcombiner(identity, r) = r

  1. 结合律:

combiner(combiner(a,b),c)=combiner(a,combiner(b,c))combiner(combiner(a,b),c) = combiner(a,combiner(b,c))

  1. 累加器和合并器语义一致:

combiner(accumulator(r,x),s)combiner(accumulator(r,x), s)

必须与把对应元素放入同一逻辑归约过程的结果相容。

例如整数加法满足结合律,因此适合并行归约:

int sum = numbers.parallelStream()
        .reduce(0, Integer::sum);

浮点数加法在数学上看似满足结合律,但计算机浮点运算会发生舍入:

double a = (1e16 + -1e16) + 1.0; // 1.0
double b = 1e16 + (-1e16 + 1.0); // 0.0,可能出现不同结果

因此并行浮点归约可能与顺序归约出现微小差异。金融金额通常不应直接使用 double 归约,而应考虑整数最小单位或 BigDecimal,并明确舍入规则。

错误归约示例:

int wrong = numbers.parallelStream()
        .reduce(0, (total, n) -> {
            total += n;
            return total;
        });

这个例子本身因为 int 不可变、累加器没有共享状态,语义仍然可能正确;真正危险的是把共享可变对象当作归约结果,并在多个任务之间共同修改它。

例如:

StringBuilder builder = new StringBuilder();

numbers.parallelStream()
        .forEach(builder::append); // 错误:共享可变对象

归约和收集应尽量使用 Stream 提供的组合机制,而不是手动共享累加器。


五、Collector:可组合的可变归约

Collector<T, A, R> 描述如何把类型为 T 的元素收集成类型为 R 的结果,中间可能使用类型为 A 的可变容器。

它包含四个核心概念:

supplier    创建空的中间容器 A
accumulator 把一个 T 加入 A
combiner    合并两个 A
finisher    把 A 转换成最终结果 R

可以写成:

T 元素
  │
  ▼
supplier → A
  │
  ├─ accumulator(A, T) → A
  │
  ├─ accumulator(A, T) → A
  │
  ├─ ...
  │
  └─ combiner(A1, A2) → A
              │
              ▼
         finisher(A) → R

1. toList:最简单的收集

List<String> names = Stream.of("a", "b", "c")
        .collect(Collectors.toList());

在 Java 16 及以后,也可以使用:

List<String> names = Stream.of("a", "b", "c").toList();

两者不要简单视为完全相同:

  • Stream.toList() 返回不可修改的 List;
  • Collectors.toList() 不保证返回的具体 List 类型,也不保证结果可修改性。

如果需要明确的可修改 ArrayList

List<String> names = Stream.of("a", "b", "c")
        .collect(Collectors.toCollection(ArrayList::new));

如果需要明确的不可修改结果,也可以:

List<String> names = Stream.of("a", "b", "c")
        .collect(Collectors.toUnmodifiableList());

2. toMap 必须处理重复键

record User(long id, String name) {}

List<User> users = List.of(
        new User(1, "Alice"),
        new User(1, "Alicia")
);

Map<Long, String> map = users.stream()
        .collect(Collectors.toMap(User::id, User::name));

这段代码会因为键 1 重复而抛出 IllegalStateException

必须显式指定冲突策略:

Map<Long, String> map = users.stream()
        .collect(Collectors.toMap(
                User::id,
                User::name,
                (oldValue, newValue) -> newValue
        ));

此时重复键保留后出现的值:

{1=Alicia}

如果业务上重复键本身就是错误,那么让异常暴露通常比静默覆盖更安全。

3. 分组和下游 Collector

record Order(String customer, int amount) {}

List<Order> orders = List.of(
        new Order("Alice", 100),
        new Order("Bob", 50),
        new Order("Alice", 30)
);

Map<String, Integer> totalByCustomer = orders.stream()
        .collect(Collectors.groupingBy(
                Order::customer,
                Collectors.summingInt(Order::amount)
        ));

System.out.println(totalByCustomer);
// {Alice=130, Bob=50}

这里不是先把每个分组收集成 List 再手工求和,而是把 summingInt 作为下游 Collector,直接在分组内部累计金额。

还可以组合多个下游操作:

Map<String, List<String>> namesByInitial = Stream.of("Alice", "Bob", "Anna")
        .collect(Collectors.groupingBy(
                name -> name.substring(0, 1),
                Collectors.toList()
        ));

groupingBy 的键数量和分布会影响内存。它需要为每个分组维护容器;高基数键可能导致大量对象和较高内存占用。

4. Collector 的并行正确性

一个 Collector 不能因为使用了线程安全容器就自动正确。它的 combiner 必须能合并并行任务产生的部分结果,并保持 Collector 的语义。

Collector.Characteristics 可能声明:

  • CONCURRENT:多个线程可以并发累加到同一个结果容器;
  • UNORDERED:结果不依赖遇到顺序;
  • IDENTITY_FINISH:中间容器 A 可以直接作为结果 R,不需要额外的 finisher

groupingByConcurrent 适用于允许并发收集且不要求结果按遇到顺序组织的场景:

ConcurrentMap<String, List<String>> grouped =
        Stream.of("Alice", "Anna", "Bob")
                .parallel()
                .collect(Collectors.groupingByConcurrent(
                        name -> name.substring(0, 1)
                ));

但“并发收集”不等于结果有序。若业务需要严格顺序,应使用有序 Collector 或先使用顺序流明确建立顺序。


六、顺序流和并行流的执行模型

1. 顺序流不等于“每个方法调用立刻执行”

顺序流表示逻辑上的串行处理。典型路径是:

一个线程:
读取元素 → filter → map → 下一个元素

这不代表实现必须为每个中间操作创建一个完整的中间集合。相反,惰性管道通常可以融合多个操作,减少中间对象。

2. 并行流如何拆分数据

并行流通常通过数据源的 Spliterator 拆分任务:

源数据
  │
  ├─ 分片 A → 任务 1
  ├─ 分片 B → 任务 2
  ├─ 分片 C → 任务 3
  └─ 分片 D → 任务 4
                │
                ▼
          合并部分结果

Spliterator 不仅能遍历元素,还能通过 trySplit() 尝试拆分剩余数据。它的特征可能包括:

  • ORDERED
  • DISTINCT
  • SORTED
  • SIZED
  • SUBSIZED
  • NONNULL
  • IMMUTABLE
  • CONCURRENT

这些特征帮助 Stream 实现判断数据源是否有顺序、大小是否已知以及能否有效拆分。

ArrayList 的随机访问和已知大小通常有利于拆分;链表、基于迭代器的外部数据源或每次读取成本差异很大的数据源,拆分可能更昂贵或不均衡。

3. 并行流使用哪个线程池

默认情况下,并行流使用由 ForkJoinPool.commonPool() 支持的公共 ForkJoinPool。它不是为当前业务单独创建的线程池。

因此下面的代码可能与应用中的其他公共 ForkJoinPool 使用者互相影响:

List<Result> results = inputs.parallelStream()
        .map(this::compute)
        .toList();

如果 compute 内部执行阻塞 I/O,例如等待数据库、HTTP 或文件系统操作,公共池中的工作线程可能长期占用,影响其他使用该公共池的任务。

可以手动调用 stream.parallel(),但这不会改变默认并行执行器。使用自定义 ForkJoinPool 包装任务在某些实现和场景下可以隔离执行器,但它不是 Stream API 的通用“传入线程池”参数,不能把它误认为所有情况下都有稳定、独立的调度保证。

4. 并行流的顺序语义

对有序源:

List<Integer> result = numbers.parallelStream()
        .map(n -> n * 2)
        .toList();

收集到 List 时,通常应保留流的遇到顺序。也就是说,输入 [1, 2, 3] 的逻辑结果是 [2, 4, 6],而不是任意排列。

但:

numbers.parallelStream()
        .forEach(System.out::println);

不保证打印顺序。

如果使用:

numbers.parallelStream()
        .forEachOrdered(System.out::println);

则需要维护遇到顺序,可能削弱并行收益。

如果业务不关心顺序,可以使用:

Set<Integer> result = numbers.parallelStream()
        .unordered()
        .collect(Collectors.toSet());

unordered() 不是随机打乱数据,而是声明后续计算不需要维护遇到顺序。这样实现可以减少排序或顺序协调成本,但结果容器本身仍要符合其自己的语义。


七、并行归约为什么要求结合律

假设输入是:

1, 2, 3, 4

顺序归约可以计算为:

(((0 + 1) + 2) + 3) + 4 = 10

并行执行可能拆成:

左分片:1, 2 → (0 + 1) + 2 = 3
右分片:3, 4 → (0 + 3) + 4 = 7
合并:3 + 7 = 10

只要运算满足结合律,两种分组都得到相同结果。

但字符串拼接通常依赖顺序:

String result = Stream.of("A", "B", "C")
        .parallel()
        .reduce("", String::concat);

对有序流,规范语义仍需与顺序拼接相容;但如果使用共享 StringBuilder,则会因为共享可变状态和并发修改而错误。

对于减法:

((0 - 1) - 2) - 3 = -6

另一种分组:

(0 - 1) - (2 - 3) = 0

减法不满足结合律,因此不能把普通减法当作任意拆分后仍保持相同语义的并行归约。


八、性能:Stream 不是自动优化器

Stream 可以改善表达能力,也可能帮助实现进行管道融合,但它不会自动消除算法复杂度、对象分配或外部系统瓶颈。

1. 先看算法复杂度

以下代码是平方级复杂度:

List<User> users = ...;
List<Long> ids = ...;

List<User> result = users.stream()
        .filter(user -> ids.contains(user.id()))
        .toList();

如果 idsArrayListcontains 平均需要扫描 O(n) 个元素,整体接近:

O(users×ids)O(|users| \times |ids|)

若只需要判断成员资格,可以预先使用 HashSet

Set<Long> idSet = new HashSet<>(ids);

List<User> result = users.stream()
        .filter(user -> idSet.contains(user.id()))
        .toList();

平均情况下,HashSet.contains 为接近 O(1),整体通常接近:

O(ids+users)O(|ids| + |users|)

这里真正带来性能变化的是数据结构和算法,而不是把 for 改成 stream()

2. 装箱和原始类型流

int sum = numbers.stream()
        .reduce(0, Integer::sum);

如果 numbersStream<Integer>,元素已经是对象类型。涉及大量数值计算时,可以使用原始类型流:

int sum = numbers.stream()
        .mapToInt(Integer::intValue)
        .sum();

IntStreamLongStreamDoubleStream 可以减少部分装箱和拆箱开销。它们仍不是保证更快的承诺;实际成本取决于数据规模、操作链、缓存局部性和终止操作。

3. 中间对象和边界操作

下面的链式代码可能产生中间对象:

List<String> result = users.stream()
        .map(User::name)
        .filter(name -> name.length() > 3)
        .map(String::toUpperCase)
        .toList();

这通常具有良好的可读性,不能仅凭“有多个 map”就断言性能差。现代 JVM 可能通过内联、逃逸分析等优化部分开销,但这些是实现和运行时优化,不是 Stream API 的规范保证。

真正需要关注的是:

  • 是否反复创建大集合;
  • 是否调用了昂贵的 sorteddistinct
  • 是否在循环中重复构建 Stream;
  • 是否产生大量临时对象;
  • 是否把数据库查询、网络调用放进每个元素的处理函数;
  • 是否因为并行同步和任务拆分抵消了计算收益。

4. 短路可以改变实际工作量

Optional<User> admin = users.stream()
        .filter(User::enabled)
        .filter(User::admin)
        .findFirst();

如果第一个元素就满足条件,后面的元素不会被访问。与之相对:

List<User> allAdmins = users.stream()
        .filter(User::enabled)
        .filter(User::admin)
        .toList();

必须检查全部元素。终止操作的选择直接决定了工作量。

5. 不要用一次 System.nanoTime 判断性能

简单基准容易被以下因素干扰:

  • JIT 预热;
  • 死代码消除;
  • 垃圾回收;
  • CPU 频率变化;
  • 数据分布;
  • 运行时线程竞争。

严肃比较应使用 JMH,至少要设置预热、测量迭代,并消费结果。例如基准应比较同一算法的:

for 循环 vs 顺序 Stream vs 并行 Stream

而不是比较不同数据结构、不同过滤条件或不同结果处理方式。

parallelStream() 是否更快,没有脱离输入规模、任务粒度、拆分成本、内存带宽和线程竞争的统一答案。每个元素只做一次简单加法时,任务调度和合并成本很可能超过计算本身;每个元素执行足够独立且 CPU 密集的复杂计算时,并行才可能有收益。


九、并行流和阻塞任务的边界

Stream 的并行化适合“数据并行”:把一批相互独立的数据分片交给多个任务处理。

它不适合直接替代以下机制:

  • 异步 I/O;
  • 请求级并发控制;
  • 超时和重试调度;
  • 结构化任务生命周期;
  • 大量阻塞操作的隔离线程池。

例如:

List<Response> responses = urls.parallelStream()
        .map(httpClient::fetch)
        .toList();

这段代码可能同时发起多个请求,但它没有表达:

  • 最大并发请求数;
  • 单个请求超时;
  • 失败后是否取消其他请求;
  • 一个请求失败时整体如何处理;
  • 请求任务由谁负责关闭和回收。

Java 25 的虚拟线程可以降低大量阻塞任务的线程占用成本,但虚拟线程不会改变 Stream 的终止、收集和并行流执行语义,也不会自动让 parallelStream() 使用虚拟线程。结构化并发则强调任务的生命周期、取消和错误传播;它与“对集合元素做并行计算”是不同抽象。

如果任务本质是受容量约束的外部 I/O,通常应显式设计执行器、信号量、超时和取消策略,而不是把外部调用隐藏在并行流的 map 中。


十、何时使用 Stream,何时使用普通循环

Stream 适合表达:

筛选 → 转换 → 聚合

例如:

Map<String, Long> activeCountByRole =
        users.stream()
                .filter(User::enabled)
                .collect(Collectors.groupingBy(
                        User::role,
                        Collectors.counting()
                ));

这种代码把业务意图直接表达为“按角色统计启用用户”。

普通循环更适合以下情况。

1. 需要复杂状态机

for (Event event : events) {
    if (state == State.CLOSED && event.type() == OPEN) {
        state = State.OPEN;
    } else if (...) {
        ...
    }
}

强行用多个 reduce 或可变 Collector 表达状态机,可能比循环更难验证。

2. 需要精确的异常恢复

for (Item item : items) {
    try {
        process(item);
    } catch (RetryableException e) {
        retry(item);
    } catch (PermanentException e) {
        recordFailure(item, e);
    }
}

Stream 中的异常仍会沿调用栈传播,但逐元素重试、跳过、补偿和错误分类会迅速让 Lambda 失去清晰度。

3. 需要多个退出条件

breakcontinue、标签跳转、阶段性提交在循环中通常更直接。Stream 的短路操作可以表达部分退出条件,但不能自然表达所有流程控制。

4. 需要调试中间状态

peek 主要用于诊断,不应作为业务副作用机制:

users.stream()
        .peek(user -> logger.debug("user={}", user))
        .filter(User::enabled)
        .toList();

由于 Stream 具有惰性,如果没有终止操作,peek 不会执行;并行流中日志顺序也可能变化。调试时还要避免把敏感信息写入日志。


十一、常见失败表现与诊断路径

同一个 Stream 被复用

表现:

java.lang.IllegalStateException: stream has already been operated upon or closed

检查是否把 Stream 保存为字段、返回后多次消费,或在一次终止操作后继续使用。

toMap 因重复键失败

表现:

IllegalStateException: Duplicate key ...

检查键是否真的唯一;若允许重复,提供合并函数;若不允许重复,保留异常并记录业务上下文。

并行流结果顺序异常

表现:

  • forEach 输出顺序变化;
  • 共享 ArrayList 结果数量不稳定;
  • 依赖执行顺序的副作用发生乱序。

检查是否误用了 forEach、共享可变对象或依赖顺序的外部操作。需要顺序时使用有序终止操作或改用顺序流,而不是事后排序所有副作用。

并行后反而变慢

诊断步骤应先拆分成本:

  1. 测量顺序流和并行流的端到端耗时;
  2. 确认数据规模是否足够;
  3. 检查每个元素的计算是否足够重;
  4. 检查数据源是否容易拆分;
  5. 检查是否有锁、I/O、日志或共享容器;
  6. 检查 GC、CPU 使用率和公共 ForkJoinPool 是否拥塞;
  7. 通过 JMH 或生产级观测验证,而不是凭一次请求判断。

十二、一个完整示例:订单统计流水线

下面的程序展示数据源、惰性处理、原始类型统计、分组 Collector 和重复键处理。

import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;

public class StreamDemo {
    record Order(
            long id,
            String customer,
            String status,
            int amount
    ) {}

    public static void main(String[] args) {
        List<Order> orders = List.of(
                new Order(1, "Alice", "PAID", 100),
                new Order(2, "Bob", "CANCELLED", 80),
                new Order(3, "Alice", "PAID", 50),
                new Order(4, "Carol", "PAID", 120),
                new Order(5, "Bob", "PAID", 70)
        );

        Set<String> paidCustomers = orders.stream()
                .filter(order -> order.status().equals("PAID"))
                .map(Order::customer)
                .collect(Collectors.toUnmodifiableSet());

        Map<String, Integer> paidAmountByCustomer = orders.stream()
                .filter(order -> order.status().equals("PAID"))
                .collect(Collectors.groupingBy(
                        Order::customer,
                        Collectors.summingInt(Order::amount)
                ));

        int totalPaid = orders.stream()
                .filter(order -> order.status().equals("PAID"))
                .mapToInt(Order::amount)
                .sum();

        Map<Long, Order> orderById = orders.stream()
                .collect(Collectors.toMap(
                        Order::id,
                        order -> order
                ));

        System.out.println(paidCustomers);
        System.out.println(paidAmountByCustomer);
        System.out.println(totalPaid);
        System.out.println(orderById.size());
    }
}

输入中的已支付订单是:

Alice: 100
Alice: 50
Carol: 120
Bob: 70

因此:

paidCustomers = [Alice, Carol, Bob]
paidAmountByCustomer = {Alice=150, Bob=70, Carol=120}
totalPaid = 340
orderById.size() = 5

其中:

  • filter 只保留 PAID 订单;
  • map(Order::customer) 把订单转换为客户名;
  • toUnmodifiableSet 去重并返回不可修改 Set;
  • groupingBy 创建按客户分组的归约;
  • summingInt 使用 IntStream 风格的整数累加,避免把金额先收集成列表;
  • toMap 使用订单 ID 作为键,因此要求 ID 唯一。

如果存在两个相同 ID,最后一个 Collector 会抛出重复键异常,而不是默认静默覆盖。


十三、使用边界的最终判断

Stream 的核心价值是把“遍历控制”转换为“数据变换和归约描述”。它最适合无副作用、可组合、边界清晰的数据处理。

判断一段代码是否适合 Stream,可以依次问:

  1. 数据源是否适合一次性遍历?
  2. 中间操作是否主要是无状态转换和筛选?
  3. 终止操作是否明确?
  4. 是否需要顺序保证?
  5. 是否存在共享可变状态?
  6. Collector 是否正确处理重复键、下游分组和并行合并?
  7. 是否真的需要并行?
  8. 每个元素的工作是否独立、足够计算密集且可拆分?
  9. 是否把阻塞 I/O、重试、取消和容量控制误塞进并行流?
  10. 用普通循环是否更容易表达状态、异常和资源生命周期?

顺序流解决的是可组合的数据处理;Collector 解决的是可验证的结果归约;并行流解决的是部分数据并行问题。它们都不能替代集合选择、算法设计、线程池隔离、I/O 并发控制或结构化任务管理。理解这些边界后,Stream 才能从语法简化工具变成可靠的计算抽象。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。