Stream 流水线与惰性求值
Java Stream 用一条流水线描述数据如何被筛选、转换和聚合。中间操作只登记处理步骤,终止操作才会触发遍历;理解这条执行链,比记住一组方法名更重要。
1. Stream 描述一次数据处理过程
一个 Stream 流水线由三部分组成:
- 数据源:集合、数组、文件或生成函数。
- 中间操作:
filter、map、sorted等转换步骤。 - 终止操作:
toList、collect、reduce、findFirst等结果出口。
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);
filter 和 map 都是中间操作。它们返回新的 Stream,并记录当前阶段需要执行的函数。加入终止操作后,数据才开始从源头进入流水线:
List<String> names = pipeline.toList();
这就是惰性求值。它带来两个直接结果:
- 没有终止操作,流水线不会遍历数据源。
- 终止操作不需要全部元素时,前面的操作也可以提前停止。
2.1 元素通常按纵向穿过流水线
对无状态中间操作,Stream 通常不会先生成完整的过滤结果,再生成完整的映射结果。一个元素可以依次经过 filter 和 map,然后才处理下一个元素。
Optional<String> first = users.stream()
.filter(User::active)
.map(User::displayName)
.findFirst();
找到第一个满足条件的结果后,findFirst 可以请求取消后续遍历。这里不需要先构造一份包含所有活跃用户的临时列表。
惰性求值不保证每条 Stream 都更快。排序、去重和复杂 Lambda 仍然有真实成本,性能要看数据源、操作顺序和终止条件。
3. 无状态操作与有状态操作
无状态操作只依赖当前元素:
filtermapflatMappeek
有状态操作需要记住已经处理过的元素,或者先看到更多输入:
distinct需要记录已出现的值。sorted通常需要取得全部输入后才能按顺序输出。limit、skip在有序并行流中需要维护位置关系。
List<String> topNames = users.stream()
.filter(User::active)
.map(User::displayName)
.sorted()
.limit(10)
.toList();
这里先过滤再排序,可以减少需要排序的数据量。把 sorted 放在 filter 前面,结果可能相同,但通常会处理更多元素。
3.1 短路终止操作
anyMatch、allMatch、noneMatch、findFirst 和 findAny 可以在结论确定后停止读取。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 适合一条方向明确的数据转换链,例如筛选、映射、分组和聚合。出现以下情况时,普通循环可能更容易维护:
- 每一步包含多处分支和提前返回。
- 处理过程需要频繁修改多个外部状态。
- 异常恢复与重试逻辑比数据转换更重要。
- 调试时必须观察细粒度的状态变化。
不要以“行数更少”作为改写依据。可以先判断每个阶段是否能用一个明确动词命名,再检查副作用、异常和数据规模。
性能敏感时至少验证:
- 数据源是否容易遍历或拆分。
sorted、distinct等有状态操作放在什么位置。- 是否因装箱产生大量临时对象。
- 短路操作能否实际提前结束。
- 普通循环、串行 Stream 和并行 Stream 的基准结果。
8. 常见问题
8.1 Stream 会修改原集合吗
Stream API 本身不会因为 filter 或 map 修改数据源,但 Lambda 可以主动修改元素或外部对象。流水线执行期间修改源集合,还可能触发 ConcurrentModificationException 或得到未定义的遍历结果。
8.2 findFirst 与 findAny 有什么区别
findFirst 在有顺序的数据源上保留遇到顺序。findAny 只要求返回任意一个元素,并行执行时有更大的调度空间。业务不关心具体哪一个结果时,可以使用 findAny。
8.3 Stream.toList() 返回的列表能修改吗
不能。Stream.toList() 返回的列表是不可修改的,调用 add、remove 等方法会抛出 UnsupportedOperationException。需要可变列表时,可以使用 collect(Collectors.toCollection(ArrayList::new))。
8.4 为什么 peek 有时没有执行
peek 是惰性中间操作。没有终止操作时整条流水线都不执行;短路或实现优化也可能让部分元素不经过它。因此不能把必须发生的业务动作放进 peek。
9. 面试题
9.1 Stream 的中间操作和终止操作有什么区别,惰性求值怎样工作
出现公司:阿里巴巴、腾讯、京东、美团、经纬恒润
考察重点
- 流水线由数据源、中间操作和终止操作组成。
- 终止操作怎样触发遍历,短路操作怎样提前停止。
- 有状态操作为什么会限制流水线融合和并行效率。
相关内容:第 1 节“Stream 描述一次数据处理过程”、第 2 节“中间操作在终止操作前不会执行”、第 3 节“无状态操作与有状态操作”。
参考回答
中间操作返回新的 Stream,只登记过滤、映射或排序等阶段;终止操作产生列表、数值或查找结果,并触发数据从源头进入流水线。无状态操作通常可以让一个元素连续经过多个阶段,减少中间容器。
findFirst、anyMatch 等短路操作在结果确定后可以停止读取。sorted、distinct 等有状态操作需要保存更多上下文,可能必须先看到全部或大量输入。惰性求值提供了优化空间,但不代表每条流水线都比循环快。
9.2 如何把 List<Person> 转成以姓名为键的 Map,姓名重复时怎么办
出现公司:飞猪
考察重点
map与collect分别负责什么。Collectors.toMap遇到重复键时的行为。- 合并规则应该来自业务语义。
相关内容:第 5 节“map、flatMap 与归约”。
参考回答
可以使用 Collectors.toMap(Person::name, Function.identity(), mergeFunction)。如果不提供合并函数,重复姓名会抛出 IllegalStateException。合并函数可以保留较新的记录、拒绝冲突或聚合成列表,但具体规则必须由姓名是否唯一以及业务怎样处理冲突决定,不能随意选择“保留第一个”。