Stream 惰性求值:filter → map → collect 到底什么时候执行?
一句话结论(30s)
Stream 的中间操作只构建 Pipeline 链表而不执行,因为真正触发计算的是 collect/forEach/reduce 这类终止操作调用的 evaluate()。关键设计是惰性求值 + Sink 链,让多步操作只遍历一次数据、短路操作提前终止;代价是并行流默认共用全 JVM 的 ForkJoinPool.commonPool,IO 阻塞任务会占满线程拖垮所有并行流。
核心原理(2min)
每个中间操作返回持有上游引用的 ReferencePipeline 形成链表,终止操作从链尾向链头逐级 opWrapSink() 构建 Sink 链后触发数据源遍历,每个元素一次性穿越所有中间操作的 Sink(避免多次全量遍历);findFirst/limit 命中即终止遍历,parallelStream() 依赖 ForkJoinPool 工作窃取,仅适合 CPU 密集任务。
底层深入(5-10min)
中间操作只构建管道,不执行
List<String> result = list.stream()
.filter(s -> s.length() > 3) // ← 不执行!
.map(String::toUpperCase) // ← 不执行!
.collect(Collectors.toList()); // ← 这里触发全部执行
每个中间操作返回一个新的 ReferencePipeline(持有对上游 Pipeline 的引用),形成链表。只有终止操作(collect/forEach/reduce)调用 evaluate() 才触发整链执行。
💭 思考:为什么中间操作不执行、要等终止操作?因为只有终止操作知道「你要什么结果」(收集成 List 还是求和还是找第一个)。中间操作只是记录「先 filter 再 map」这个意图,攒成一条管道,终止操作才一次性驱动数据流过管道。
Sink 链:每个元素一次性穿越所有中间操作
// filter 的 Sink
class FilterSink<T> implements Sink<T> {
Sink<T> downstream;
Predicate<T> predicate;
public void accept(T t) {
if (predicate.test(t)) downstream.accept(t); // 通过才往下传
}
}
终止操作触发后,Pipeline 从链尾向链头逐级调用 opWrapSink() 构建完整 Sink 链,再触发数据源遍历。每个元素从数据源被取出后,一次性穿越所有中间操作的 Sink——避免了多次全量遍历。
💭 思考:Sink 链为什么能避免多次遍历?如果 filter 完再 map,就要遍历两遍数据。Sink 链把每个操作包装成 accept 方法串起来,每个元素从源头取出后一次穿过所有 Sink,一遍遍历完成所有操作——数据像流水线一样流过。
短路操作:findFirst/limit 提前终止
list.stream().filter(x -> x > 10).findFirst();
findFirst 找到第一个满足 x > 10 的元素 → 立即终止遍历。后续元素不会被处理。惰性求值让 findFirst().limit(1) 本身就是单次遍历——一旦找到目标立刻停止,不需要先算出所有满足条件的再取第一个。
💭 思考:短路操作和惰性求值什么关系?正因为惰性,findFirst 可以边遍历边判断、找到就停,不需要先算出全部结果。如果像普通集合操作那样先 filter 全量再取第一个,短路就没意义了。惰性让短路成为可能。
parallelStream() 的陷阱
list.parallelStream().map(this::slowIO).collect(toList());
// slowIO() 是阻塞 IO → commonPool 线程被占满 → 整个 JVM 所有并行流卡住
parallelStream() 默认用 ForkJoinPool.commonPool()——全 JVM 共享的线程池。CPU 密集计算(如大规模数字求和)可以用,IO 阻塞任务必须用自定义线程池。
💭 思考:parallelStream 为什么默认用 commonPool 会出问题?因为 commonPool 是全局唯一、线程数默认等于 CPU 核数,全 JVM 的并行流都抢它;一个 IO 阻塞任务占住线程,别的并行流就饿死。所以 IO 密集任务必须用自定义线程池。
为什么ForkJoinPool适合CPU密集而非IO?
ForkJoinPool 内置工作窃取(Work Stealing)——空闲线程从繁忙线程的双端队列尾部偷任务。这需要任务是纯 CPU 计算(可在任意线程执行),IO 阻塞任务窃取后也无法推进(等待 IO 返回)。CPU 密集 + 可窃取 = 所有核心满载。IO 密集 + 阻塞 = 窃取也没用。
💭 思考:工作窃取为什么只适合 CPU 密集?因为窃取的前提是任务在任意线程都能推进;CPU 任务是纯计算、可窃取,能满载所有核;IO 任务窃走后还得等 IO 返回,阻塞着不推进,窃取也白窃。
总结
Stream = 声明式数据处理管道。惰性求值 + Sink 链的组合让多步操作只遍历一次数据,短路操作提前终止,是函数式 Java 最精妙的设计之一。
章末提问
追问 1:Stream 的中间操作为什么不立即执行?终止操作扮演什么角色?
结论先行:中间操作只构建持有上游引用的 ReferencePipeline 链表记录意图,真正触发计算的是终止操作调用的 evaluate()。
因为:中间操作不知道「最终要什么结果」,只能把 filter/map 攒成管道;终止操作(collect/forEach/reduce)知道目标,触发时从链尾向链头 opWrapSink 构建 Sink 链,再驱动数据源遍历。这是典型的惰性求值。
追问 2:Sink 链是怎么做到「多步操作只遍历一次」的?
结论先行:每个中间操作都包装成一个 Sink,各 Sink 通过 downstream 引用串联成链,每个元素从源头取出后一次穿过整条链。
因为:FilterSink.accept 里 if (predicate.test(t)) downstream.accept(t),通过就传给下游,依次往下;这样每个元素只被取出一次、在一条链上完成全部中间操作,避免了「filter 遍历一遍、map 再遍历一遍」的多次全量遍历。
追问 3:parallelStream 为什么 IO 密集场景会踩坑?
结论先行:因为 parallelStream 默认用全 JVM 共享的 ForkJoinPool.commonPool,线程数等于 CPU 核数,IO 阻塞会占满这些线程拖垮所有并行流。
因为:commonPool 是全局唯一的,IO 任务在等待返回时阻塞线程、又不释放,其他并行流拿不到线程;且 ForkJoinPool 的工作窃取只对 CPU 密集任务有效(任务可窃取、可满载),IO 阻塞任务窃走后也无法推进。所以 IO 密集必须换自定义线程池。