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是数据源;filter和map是中间操作;toList是终止操作;- Stream 本身不负责保存结果,结果由终止操作产生。
Stream API 位于 java.util.stream 包中,Java 25 的相关核心抽象仍然是 Stream、IntStream、LongStream、DoubleStream、Collector 和 Spliterator。
一、先区分集合、迭代器和 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;
});
原因是 filter 和 map 都是中间操作。它们返回了新的 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;元素 2 被 filter 排除,不会进入 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,后面的 10 和 12 不需要访问。
常见短路操作包括:
anyMatchallMatchnoneMatchfindFirstfindAny- 某些情况下的
limit
但“支持短路”不代表所有元素在所有并行执行场景下都会立刻停止。并行流中可能已经有多个任务在处理数据,取消通常只能阻止尚未开始或尚未完成的部分,不能撤销已经发生的外部副作用。
3. 有状态中间操作可能需要缓存或全局信息
并非所有中间操作都能只看当前元素:
stream.sorted()
stream.distinct()
stream.limit(n)
stream.skip(n)
例如 sorted() 必须知道足够多的输入元素,才能确定最小元素和最大元素。因此它通常需要缓存输入并进行排序。
这与下面的 filter 不同:
stream.filter(x -> x > 10)
filter 只需要检查当前元素,通常可以边读边输出。
可以将操作粗略分为:
| 类型 | 示例 | 是否需要全局信息 |
|---|---|---|
| 无状态 | map、filter、peek |
通常不需要 |
| 有状态 | sorted、distinct、skip |
可能需要 |
| 短路 | findFirst、anyMatch、limit |
可能提前结束 |
“无状态”并不等于“没有副作用”。它主要描述操作是否需要记住之前元素。一个把数据写入共享 List 的 map 仍然可能有严重副作用。
三、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);
可以形式化为:
其中:
identity是初始值;x_i是第i个元素;accumulator把当前结果和一个元素合并。
对于并行执行,还需要一个合并器:
要让并行结果与顺序结果一致,通常需要满足:
- 恒等元条件:
- 结合律:
- 累加器和合并器语义一致:
必须与把对应元素放入同一逻辑归约过程的结果相容。
例如整数加法满足结合律,因此适合并行归约:
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() 尝试拆分剩余数据。它的特征可能包括:
ORDEREDDISTINCTSORTEDSIZEDSUBSIZEDNONNULLIMMUTABLECONCURRENT
这些特征帮助 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();
如果 ids 是 ArrayList,contains 平均需要扫描 O(n) 个元素,整体接近:
若只需要判断成员资格,可以预先使用 HashSet:
Set<Long> idSet = new HashSet<>(ids);
List<User> result = users.stream()
.filter(user -> idSet.contains(user.id()))
.toList();
平均情况下,HashSet.contains 为接近 O(1),整体通常接近:
这里真正带来性能变化的是数据结构和算法,而不是把 for 改成 stream()。
2. 装箱和原始类型流
int sum = numbers.stream()
.reduce(0, Integer::sum);
如果 numbers 是 Stream<Integer>,元素已经是对象类型。涉及大量数值计算时,可以使用原始类型流:
int sum = numbers.stream()
.mapToInt(Integer::intValue)
.sum();
IntStream、LongStream 和 DoubleStream 可以减少部分装箱和拆箱开销。它们仍不是保证更快的承诺;实际成本取决于数据规模、操作链、缓存局部性和终止操作。
3. 中间对象和边界操作
下面的链式代码可能产生中间对象:
List<String> result = users.stream()
.map(User::name)
.filter(name -> name.length() > 3)
.map(String::toUpperCase)
.toList();
这通常具有良好的可读性,不能仅凭“有多个 map”就断言性能差。现代 JVM 可能通过内联、逃逸分析等优化部分开销,但这些是实现和运行时优化,不是 Stream API 的规范保证。
真正需要关注的是:
- 是否反复创建大集合;
- 是否调用了昂贵的
sorted或distinct; - 是否在循环中重复构建 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. 需要多个退出条件
break、continue、标签跳转、阶段性提交在循环中通常更直接。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、共享可变对象或依赖顺序的外部操作。需要顺序时使用有序终止操作或改用顺序流,而不是事后排序所有副作用。
并行后反而变慢
诊断步骤应先拆分成本:
- 测量顺序流和并行流的端到端耗时;
- 确认数据规模是否足够;
- 检查每个元素的计算是否足够重;
- 检查数据源是否容易拆分;
- 检查是否有锁、I/O、日志或共享容器;
- 检查 GC、CPU 使用率和公共 ForkJoinPool 是否拥塞;
- 通过 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,可以依次问:
- 数据源是否适合一次性遍历?
- 中间操作是否主要是无状态转换和筛选?
- 终止操作是否明确?
- 是否需要顺序保证?
- 是否存在共享可变状态?
- Collector 是否正确处理重复键、下游分组和并行合并?
- 是否真的需要并行?
- 每个元素的工作是否独立、足够计算密集且可拆分?
- 是否把阻塞 I/O、重试、取消和容量控制误塞进并行流?
- 用普通循环是否更容易表达状态、异常和资源生命周期?
顺序流解决的是可组合的数据处理;Collector 解决的是可验证的结果归约;并行流解决的是部分数据并行问题。它们都不能替代集合选择、算法设计、线程池隔离、I/O 并发控制或结构化任务管理。理解这些边界后,Stream 才能从语法简化工具变成可靠的计算抽象。
系列导航与关联阅读
- 系列入口:Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Java 反射与注解:Class、MethodHandle、元注解、处理器和边界
- 下一篇:Java 网络与 HTTP Client:连接、TLS、超时、流和错误恢复
- 延伸:Java 集合框架:List、Set、Map、Queue、迭代器和复杂度
- 延伸:Java 25 虚拟线程与结构化并发:调度、Pinning、取消和容量
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论