Stream / Spliterator / 流水线源码精读
概述
Stream 的优雅背后是一套惰性流水线(Lazy Pipeline):中间操作只组装 Sink 链,终端操作才触发遍历执行。Spliterator 是流的"数据源迭代器",支持拆分以便并行。理解 AbstractPipeline、Sink、Spliterator 三者的配合,就掌握了 Stream 的求值模型。本文基于 OpenJDK 21 源码拆解。
一、流水线的创建
1.1 Stream.of(T...)
java
// java.util.stream.Stream
@SafeVarargs
public static<T> Stream<T> of(T... values) {
return Arrays.stream(values); // 委托给数组流
}
// java.util.Arrays.stream
public static <T> Stream<T> stream(T[] array) {
return stream(array, 0, array.length);
}
public static <T> Stream<T> stream(T[] array, int startInclusive, int endExclusive) {
return StreamSupport.stream(spliterator(array, startInclusive, endExclusive), false);
}
// StreamSupport.stream → ReferencePipeline.Head
public static <T> Stream<T> stream(Spliterator<T> spliterator, boolean parallel) {
return new ReferencePipeline.Head<>(spliterator,
StreamOpFlag.fromCharacteristics(spliterator), parallel);
}流水线创建链路:
Stream.of(values)
→ Arrays.stream(array)
→ StreamSupport.stream(spliterator)
→ ReferencePipeline.Head(流水线头节点,持有数据源 Spliterator)1.2 惰性求值
中间操作(filter/map/sorted...):
只往流水线追加 stage(Sink 链节点),不遍历数据
返回值仍是 Stream → 可继续链式
终端操作(collect/forEach/count...):
触发 evaluate → 从源头 Spliterator 开始遍历
数据一次性流过整个 Sink 链java
// 中间操作示例:只是构建 stage,不执行
stream.filter(x -> x > 0) // 添加 filter stage
.map(x -> x * 2) // 添加 map stage
.collect(toList()); // 终端操作才真正遍历二、Spliterator 的四个核心方法
java
// java.util.Spliterator
public interface Spliterator<T> {
boolean tryAdvance(Consumer<? super T> action); // ① 尝试取下一个
default void forEachRemaining(Consumer<? super T> action); // ② 遍历剩余全部
Spliterator<T> trySplit(); // ③ 拆分子迭代器
long estimateSize(); // ④ 预估剩余数量
int characteristics(); // 特性位标志
}| 方法 | 语义 |
|---|---|
tryAdvance(action) | 取一个元素交给 action,有则返回 true,耗尽返回 false |
forEachRemaining(action) | 把剩余元素全部交给 action(顺序执行) |
trySplit() | 拆分出子迭代器(并行执行),返回 null 表示不可再分 |
estimateSize() | 预估剩余元素数(配合 SIZED 特性可精确) |
java
// 串行遍历:tryAdvance 循环
Spliterator<T> sp = stream.spliterator();
while (sp.tryAdvance(e -> result.add(e))) { }三、ArrayListSpliterator 的拆分实现
java
// java.util.ArrayList.ArrayListSpliterator
static final class ArrayListSpliterator<E> implements Spliterator<E> {
private final ArrayList<E> list;
private int index; // 当前位置
private int fence; // 结束位置(-1 = 未知,惰性取 size)
private int expectedModCount; // 结构性修改检测
// 拆分:二分取中
public ArrayListSpliterator<E> trySplit() {
int hi = getFence(), lo = index, mid = (lo + hi) >>> 1;
// 有足够元素才拆
if (lo >= mid) return null;
return new ArrayListSpliterator<>(list, lo, index = mid);
// 新迭代器 [lo, mid),原迭代器前移到 mid
}
// 遍历单个
public boolean tryAdvance(Consumer<? super E> action) {
if (action == null) throw new NullPointerException();
int hi = getFence(), i = index;
if (i < hi) {
index = i + 1;
action.accept(list.get(i)); // 取元素
return true;
}
return false; // 耗尽
}
public long estimateSize() { return getFence() - index; } // 剩余数量
}trySplit 二分拆分:
(lo + hi) >>> 1 → 中点
拆出 [lo, mid) 给新迭代器;自己继续 [mid, hi)
并行流 → 递归拆分到足够小粒度(ForkJoinPool 并行处理)
SIZED 特性 → 可精确切分(均衡负载)四、AbstractPipeline 链式结构
java
// java.util.stream.AbstractPipeline
abstract class AbstractPipeline<E_IN, E_OUT, S extends BaseStream<E_OUT, S>>
implements PipelineHelper<E_OUT> {
private final AbstractPipeline sourceStage; // 头节点(数据源)
private final long sourceOrOpFlags; // 源或操作的标志
private AbstractPipeline previousStage; // 前驱
private AbstractPipeline nextStage; // 后继
private int depth; // 深度(距头节点距离)
private Spliterator<?> sourceSpliterator; // 数据源迭代器
}流水线节点关系:
Head(头)→ stage1 → stage2 → ... → TerminalOp
sourceStage:恒指向头节点(所有节点共享)
previousStage/nextStage:双向链表
depth:头=0,每加一级 +1(用于分流策略)
创建中间操作 stage 时:
depth 递增加一 → 新 stage 接入链表
sink 链在终端操作时才 wrap 出来java
// ReferencePipeline.map 创建 stage(惰性)
public final <R> Stream<R> map(Function<? super P_OUT, ? extends R> mapper) {
return new StatelessOp<>(this, StreamShape.REFERENCE, StreamOpFlag.NOT_SORTED) {
// opWrapSink:返回一个 Sink,在 accept 时应用 mapper
Sink<P_OUT> opWrapSink(int flags, Sink<R> downstream) {
return new Sink.ChainedReference<>(downstream) {
public void accept(P_OUT u) {
downstream.accept(mapper.apply(u)); // 转换后传给下游
}
};
}
};
}五、终端操作 evaluate 三步
java
// AbstractPipeline.evaluate
final <R> R evaluate(TerminalOp<E_OUT, R> terminalOp) {
// ① 并行分支
if (isParallel()) { ... parallel processing ... }
// ② 串行执行
return terminalOp.evaluateSequential(this, sourceSpliterator());
}
// ReduceOp.evaluateSequential 核心
public <P_IN> R evaluateSequential(PipelineHelper<T> helper,
Spliterator<P_IN> spliterator) {
return helper.wrapAndCopyInto(makeSink(), spliterator)
.get(); // 最终获取结果
}三步(源码中 wrapAndCopyInto 展开):
evaluate 三步骤:
① makeSink():创建终端 Sink(收集结果/计数/求和的载体)
② wrapSink(sink):从流水线末尾往前,把每个 stage 的 opWrapSink
包装成一条链(head → ... → terminal sink)
③ copyInto(wrappedSink, spliterator):用 spliterator 遍历源数据,
逐个元素灌入 sink 链(tryAdvance 驱动)java
// AbstractPipeline.wrapAndCopyInto
final <P_IN, S extends Sink<E_OUT>> S wrapAndCopyInto(S sink, Spliterator<P_IN> spliterator) {
copyInto(wrapSink(sink), spliterator);
return sink;
}
final <P_IN> void copyInto(Sink<P_IN> wrappedSink, Spliterator<P_IN> spliterator) {
// 批量遍历(forEachRemaining)比单元素 tryAdvance 更高效
spliterator.forEachRemaining(wrappedSink);
}求值顺序:
filter → map 等 stage 只是"组装"
终端操作才 makeSink + wrapSink + copyInto
→ 所以中间操作惰性、终端操作热性六、Sink 链的执行
java
// java.util.stream.Sink
interface Sink<T> extends Consumer<T> {
default void begin(long size) { } // ① 链开始时调用(如预分配容器)
void accept(T t); // ② 处理每个元素(核心)
default void end() { } // ③ 链结束时调用(如排序合并)
default boolean cancellationRequested() { return false; } // 短路支持
}Sink 链执行周期(每个 stage 实现):
begin(size) → 准备(分配容器/初始化状态)
accept(value) → 逐个元素处理(转换/过滤/累积)
end() → 收尾(合并结果/关闭资源)
链式传递:
上游 accept 处理完 → 调下游 accept
数据从源头单向流到终端 sink
filter 短路:不满足条件 → 不下传(相当于跳过)
limit/anyMatch 等 → cancellationRequested 提前终止遍历java
// 终端 sink 示例:ArrayList 收集
class ListSink implements Sink<T> {
ArrayList<T> list;
public void begin(long size) { list = new ArrayList<>((int) size); }
public void accept(T t) { list.add(t); }
}七、collect(Collector) 的 ReduceOp
java
// java.util.stream.ReduceOp
static final class ReduceOp<T, R, I> implements TerminalOp<T, R> {
private final Collector<? super T, I, R> collector;
public <P_IN> R evaluateSequential(PipelineHelper<T> helper,
Spliterator<P_IN> spliterator) {
// ① makeSink
// ② 遍历累积到可变容器(如 ArrayList/Integer[1])
// ③ finisher 转换结果
return collector.finisher().apply(
helper.wrapAndCopyInto(makeSink(), spliterator).get());
}
private static final class ReducingSink extends Box<I> implements Sink<T> {
public void begin(long size) { state = collector.supplier().get(); } // 容器
public void accept(T t) {
collector.accumulator().accept(state, t); // 累积元素
}
public I get() {
return state; // 返回容器
}
}
}collect 流程(串行):
① supplier.get():创建可变结果容器(如 new ArrayList)
② accumulator.accept(container, t):把每个元素累积进容器
③ combiner:并行时合并两个容器(串行不调用)
④ finisher.apply(container):转换为最终结果(可选)
并行 collect:
每个线程一个容器 → 各自累积
最后 combiner 两两合并 → finisher 转换
CONCURRENT 特性 → 所有线程共享同一容器(并发累积)java
// 自定义 Collector 结构
Collector<T, A, R> c = Collector.of(
() -> new ArrayList<>(), // supplier:容器
(list, t) -> list.add(t), // accumulator:累积
(l1, l2) -> { l1.addAll(l2); return l1; }, // combiner:合并
list -> list, // finisher:转换
Collector.Characteristics.IDENTITY_FINISH // 特性
);八、StreamOpFlag 位标志
java
// java.util.stream.StreamOpFlag
static final int DISTINCT = 0x00000001; // 去重
static final int SORTED = 0x00000004; // 已排序
static final int ORDERED = 0x00000010; // 保序
static final int SIZED = 0x00000040; // 大小已知
static final int SHORT_CIRCUIT = 0x00040000; // 可短路(limit/anyMatch)位标志的作用:
source 特性(Spliterator.characteristics)→ 传入流水线标志
stage 可能改变标志:map 破坏 SORTED、filter 保留 SIZED(只增不减)
终端操作读取标志做优化:
SIZED → 预分配容器大小(避免扩容)
SORTED → sorted() 可跳过(已排序)
DISTINCT → distinct() 可跳过(已去重)
SHORT_CIRCUIT → 遍历可提前终止特性沿流水线传播:
sourceStage 记录源特性
每层 stage 用 opFlags 声明"改变了什么"(如 NOT_SORTED)
最终 evaluate 用最终 flags 决定优化策略九、实现要点
Stream 流水线核心:
ReferencePipeline.Head:流水线头节点(数据源)
Spliterator:tryAdvance/forEachRemaining/trySplit/estimateSize
ArrayListSpliterator:二分拆分(lo+hi)>>>1),支持并行
AbstractPipeline:sourceStage + previousStage/nextStage + depth
惰性求值:中间操作组装,终端操作才遍历
evaluate 三步:makeSink → wrapSink → copyInto
Sink 链:begin → accept → end,filter 短路不下传
collect:supplier/accumulator/combiner/finisher 四函数
StreamOpFlag:SIZED/SORTED/DISTINCT 等位标志驱动优化
常见陷阱:
无限流 + 无限中间操作 → 内存耗尽(配合 limit)
在流水线中做副作用操作 → 并行时顺序不确定
大集合上 sorted/distinct → 有状态操作会缓冲全部数据
重复使用已消费的流 → IllegalStateException
parallelStream 的容器线程安全 → 收集器不保证并发安全(无 CONCURRENT)