跳到主要内容

Stream 流水线与惰性求值

Java Stream 用一条流水线描述数据如何被筛选、转换和聚合。中间操作只登记处理步骤,终止操作才会触发遍历;理解这条执行链,比记住一组方法名更重要。

1. Stream 描述一次数据处理过程

一个 Stream 流水线由三部分组成:

  1. 数据源:集合、数组、文件或生成函数。
  2. 中间操作filtermapsorted 等转换步骤。
  3. 终止操作toListcollectreducefindFirst 等结果出口。
List<String> names = users.stream()
.filter(User::active)
.map(User::displayName)
.toList();

users 保存数据,Stream 不保存另一份用户列表。它保存的是“从这个数据源开始,依次执行过滤和映射,最后收集结果”的处理描述。

Stream 使用内部迭代:调用方提供每个元素的处理规则,遍历时机和执行方式由 Stream 实现控制。这使框架能够合并处理步骤、提前结束遍历,也能在满足条件时并行执行。

2. 中间操作在终止操作前不会执行

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

Stream<String> pipeline = users.stream()
.filter(user -> {
System.out.println("filter " + user.id());
return user.active();
})
.map(User::displayName);

filtermap 都是中间操作。它们返回新的 Stream,并记录当前阶段需要执行的函数。加入终止操作后,数据才开始从源头进入流水线:

List<String> names = pipeline.toList();

这就是惰性求值。它带来两个直接结果:

  • 没有终止操作,流水线不会遍历数据源。
  • 终止操作不需要全部元素时,前面的操作也可以提前停止。

2.1 元素通常按纵向穿过流水线

对无状态中间操作,Stream 通常不会先生成完整的过滤结果,再生成完整的映射结果。一个元素可以依次经过 filtermap,然后才处理下一个元素。

Optional<String> first = users.stream()
.filter(User::active)
.map(User::displayName)
.findFirst();

找到第一个满足条件的结果后,findFirst 可以请求取消后续遍历。这里不需要先构造一份包含所有活跃用户的临时列表。

惰性求值不保证每条 Stream 都更快。排序、去重和复杂 Lambda 仍然有真实成本,性能要看数据源、操作顺序和终止条件。

3. 无状态操作与有状态操作

无状态操作只依赖当前元素:

  • filter
  • map
  • flatMap
  • peek

有状态操作需要记住已经处理过的元素,或者先看到更多输入:

  • distinct 需要记录已出现的值。
  • sorted 通常需要取得全部输入后才能按顺序输出。
  • limitskip 在有序并行流中需要维护位置关系。
List<String> topNames = users.stream()
.filter(User::active)
.map(User::displayName)
.sorted()
.limit(10)
.toList();

这里先过滤再排序,可以减少需要排序的数据量。把 sorted 放在 filter 前面,结果可能相同,但通常会处理更多元素。

3.1 短路终止操作

anyMatchallMatchnoneMatchfindFirstfindAny 可以在结论确定后停止读取。limit 是短路中间操作,它限制后续阶段最多接收多少元素。

无限流必须配合可以终止的边界:

List<Integer> powers = Stream.iterate(1, value -> value * 2)
.limit(10)
.toList();

删除 limit 后,toList 永远等不到数据结束,并且会持续消耗内存。

4. Stream 只能消费一次

终止操作执行后,Stream 就被消费:

Stream<User> activeUsers = users.stream().filter(User::active);

long count = activeUsers.count();
// List<User> result = activeUsers.toList();
// 运行时抛出 IllegalStateException

Stream 可能持有遍历位置、短路状态和资源,第二次执行无法保证仍从相同状态开始。需要重复计算时,重新从可重复读取的数据源创建 Stream:

Supplier<Stream<User>> activeUsers =
() -> users.stream().filter(User::active);

long count = activeUsers.get().count();
List<User> result = activeUsers.get().toList();

Supplier 每次返回的是新 Stream,不是复用已经消费的对象。

4.1 I/O 数据源需要关闭

集合创建的 Stream 通常没有需要释放的资源。Files.lines 等 I/O 方法返回的 Stream 持有打开的文件,应该使用 try-with-resources:

try (Stream<String> lines = Files.lines(path, StandardCharsets.UTF_8)) {
long nonBlankLines = lines.filter(line -> !line.isBlank()).count();
}

终止操作完成不等于所有 Stream 都会自动关闭。资源所有权仍由创建它的调用方负责。

5. map、flatMap 与归约

map 让一个输入对应一个输出:

List<Long> ids = users.stream()
.map(User::id)
.toList();

如果映射结果本身也是容器或 Stream,flatMap 会把多层结果展开:

List<String> roles = users.stream()
.flatMap(user -> user.roles().stream())
.distinct()
.toList();

使用 map 会得到 Stream<List<String>>,使用 flatMap 得到连续的 Stream<String>

reduce 把元素归约成一个值:

int total = orders.stream()
.mapToInt(Order::quantity)
.sum();

求和、计数和最值优先使用专门的终止操作。把结果收集到可变容器时使用 collect,不要用 reduce 反复复制或修改同一个容器。

6. 副作用会破坏流水线的可推理性

下面的代码把结果写入外部列表:

List<String> names = new ArrayList<>();

users.stream()
.filter(User::active)
.map(User::displayName)
.forEach(names::add);

串行执行时它看起来可用,改成并行流后却会并发修改 ArrayList。直接使用终止操作更清楚,也能让实现选择正确的合并方式:

List<String> names = users.stream()
.filter(User::active)
.map(User::displayName)
.toList();

peek 主要用于观察流水线,不适合承担保存数据库、发送消息或修改业务状态等必要逻辑。短路操作可能让一部分元素根本不会到达 peek,实现也可以在不影响结果时优化遍历。

7. Stream 的适用边界

Stream 适合一条方向明确的数据转换链,例如筛选、映射、分组和聚合。出现以下情况时,普通循环可能更容易维护:

  • 每一步包含多处分支和提前返回。
  • 处理过程需要频繁修改多个外部状态。
  • 异常恢复与重试逻辑比数据转换更重要。
  • 调试时必须观察细粒度的状态变化。

不要以“行数更少”作为改写依据。可以先判断每个阶段是否能用一个明确动词命名,再检查副作用、异常和数据规模。

性能敏感时至少验证:

  1. 数据源是否容易遍历或拆分。
  2. sorteddistinct 等有状态操作放在什么位置。
  3. 是否因装箱产生大量临时对象。
  4. 短路操作能否实际提前结束。
  5. 普通循环、串行 Stream 和并行 Stream 的基准结果。

8. 常见问题

8.1 Stream 会修改原集合吗

Stream API 本身不会因为 filtermap 修改数据源,但 Lambda 可以主动修改元素或外部对象。流水线执行期间修改源集合,还可能触发 ConcurrentModificationException 或得到未定义的遍历结果。

8.2 findFirstfindAny 有什么区别

findFirst 在有顺序的数据源上保留遇到顺序。findAny 只要求返回任意一个元素,并行执行时有更大的调度空间。业务不关心具体哪一个结果时,可以使用 findAny

8.3 Stream.toList() 返回的列表能修改吗

不能。Stream.toList() 返回的列表是不可修改的,调用 addremove 等方法会抛出 UnsupportedOperationException。需要可变列表时,可以使用 collect(Collectors.toCollection(ArrayList::new))

8.4 为什么 peek 有时没有执行

peek 是惰性中间操作。没有终止操作时整条流水线都不执行;短路或实现优化也可能让部分元素不经过它。因此不能把必须发生的业务动作放进 peek

9. 面试题

9.1 Stream 的中间操作和终止操作有什么区别,惰性求值怎样工作

出现公司:阿里巴巴、腾讯、京东、美团、经纬恒润

考察重点

  • 流水线由数据源、中间操作和终止操作组成。
  • 终止操作怎样触发遍历,短路操作怎样提前停止。
  • 有状态操作为什么会限制流水线融合和并行效率。

相关内容:第 1 节“Stream 描述一次数据处理过程”、第 2 节“中间操作在终止操作前不会执行”、第 3 节“无状态操作与有状态操作”。

参考回答

中间操作返回新的 Stream,只登记过滤、映射或排序等阶段;终止操作产生列表、数值或查找结果,并触发数据从源头进入流水线。无状态操作通常可以让一个元素连续经过多个阶段,减少中间容器。

findFirstanyMatch 等短路操作在结果确定后可以停止读取。sorteddistinct 等有状态操作需要保存更多上下文,可能必须先看到全部或大量输入。惰性求值提供了优化空间,但不代表每条流水线都比循环快。

9.2 如何把 List<Person> 转成以姓名为键的 Map,姓名重复时怎么办

出现公司:飞猪

考察重点

  • mapcollect 分别负责什么。
  • Collectors.toMap 遇到重复键时的行为。
  • 合并规则应该来自业务语义。

相关内容:第 5 节“map、flatMap 与归约”。

参考回答

可以使用 Collectors.toMap(Person::name, Function.identity(), mergeFunction)。如果不提供合并函数,重复姓名会抛出 IllegalStateException。合并函数可以保留较新的记录、拒绝冲突或聚合成列表,但具体规则必须由姓名是否唯一以及业务怎样处理冲突决定,不能随意选择“保留第一个”。