大翔:飞翔科技的实时监控大屏,响应时间指标每 5 秒一个窗口做聚合,用 Stream 怎么搞? 白歌:
collect分组只能拿到最终结果,中间的窗口滑动做不了。JDK 22 的 Stream Gatherers 就是干这个的——自定义中间操作。 小崔:Gatherer 和 Collector 有啥本质区别? 孔蓝:我猜……Collector 是终止操作,Gatherer 是中间操作?
Stream Gatherers(JDK 22 预览 / JDK 24 第二预览)
核心定义
| 概念 | 定义 |
|---|---|
| Gatherer | Stream 的自定义中间操作,可逐元素处理、可短路、可有状态 |
| Collector | Stream 的终止操作,将元素归约为最终结果 |
| Integrator | 整合器,逐元素调用,决定是否推送到下游 |
| Initializer | 初始化器,创建中间状态对象 |
| Finisher | 完成器,所有元素处理完后执行收尾逻辑 |
| Combiner | 组合器,并行流时合并两个中间状态 |
Gatherer vs Collector
本质区别:
Collector消费流并产出最终值,是 终止操作;Gatherer** 在流中做变换再推送下游,是 **中间操作**,后面还能接filter、map、collect` 等。
Gatherer 四大组件
| 组件 | 作用 | 对应 Collector |
|---|---|---|
| Initializer | 创建中间状态(如窗口缓冲区) | supplier() |
| Integrator | 每个元素到达时执行逻辑 | accumulator() |
| Finisher | 流结束后收尾(如刷新残余窗口) | finisher() |
| Combiner | 并行流合并两个中间状态 | combiner() |
JDK 内置 Gatherers
| 方法 | 功能 |
|---|---|
Gatherers.windowFixed(n) | 将元素按固定大小 n 分组,不足一组时保留 |
Gatherers.windowSliding(n) | 滑动窗口,每次滑动 1 步,窗口大小 n |
Gatherers.fold(init, fn) | 有状态折叠,类似 reduce 但可推送多个中间结果 |
Gatherers.mapConcurrent(limit, fn) | 并发映射,最多 limit 个并发任务 |
windowFixed vs windowSliding
实战场景
场景一:飞翔科技实时指标聚合
// === 场景说明:每 5 个采样点一个窗口,统计平均响应时间 ===
import java.util.stream.Gatherers;
record Sample(long timestamp, int responseMs) {}
List<Sample> samples = List.of(
new Sample(1, 120), new Sample(2, 85),
new Sample(3, 200), new Sample(4, 90),
new Sample(5, 150), new Sample(6, 110),
new Sample(7, 300), new Sample(8, 75)
);
List<Double> avgPerWindow = samples.stream()
.gather(Gatherers.windowFixed(5))
.map(window -> window.stream()
.mapToInt(Sample::responseMs)
.average().orElse(0))
.toList();
System.out.println(avgPerWindow);
[129.0, 161.66666666666666]
场景二:连续重复日志去重
// === 场景说明:连续重复的日志只保留第一条 ===
import java.util.stream.Gatherer;
Gatherer<String, ?, String> dedupConsecutive = Gatherer.ofSequential(
() -> new Object() { String prev = null; },
(state, element, downstream) -> {
if (!element.equals(state.prev)) {
state.prev = element;
downstream.push(element);
}
return true;
}
);
List<String> logs = List.of("ERROR", "ERROR", "WARN", "ERROR", "ERROR", "ERROR", "INFO");
List<String> deduped = logs.stream().gather(dedupConsecutive).toList();
System.out.println(deduped);
[ERROR, WARN, ERROR, INFO]
场景三:批量写入订单
// === 场景说明:每 3 条订单打包为一批,模拟批量写入数据库 ===
List<String> orders = List.of("O1","O2","O3","O4","O5","O6","O7");
List<String> batchResults = orders.stream()
.gather(Gatherers.windowFixed(3))
.map(batch -> "写入DB: " + batch)
.toList();
System.out.println(batchResults);
[写入DB: [O1, O2, O3], 写入DB: [O4, O5, O6], 写入DB: [O7]]
与 collect 后手动处理 List 的对比
// === 场景说明:对比 windowFixed 与手动分组的写法 ===
// ❌ 旧方式:collect 后手动分片,失去流的惰性,必须等全部元素就绪
List<List<Sample>> manualWindows = new ArrayList<>();
List<Sample> buffer = new ArrayList<>();
for (Sample s : samples) {
buffer.add(s);
if (buffer.size() == 5) {
manualWindows.add(new ArrayList<>(buffer));
buffer.clear();
}
}
if (!buffer.isEmpty()) manualWindows.add(new ArrayList<>(buffer));
// ✅ Gatherer 方式:保持流式,惰性求值,可短路
List<List<Sample>> gatherWindows = samples.stream()
.gather(Gatherers.windowFixed(5))
.toList();
Gatherer 方式保持流的惰性求值链路,中间操作不触发计算,且支持短路——找到目标即可提前终止。手动分片则必须遍历全部数据。
短路操作(Short-Circuiting)
// === 场景说明:找到第一个平均响应时间 > 200ms 的窗口后立即停止 ===
Optional<Double> firstSlow = samples.stream()
.gather(Gatherers.windowFixed(5))
.map(window -> window.stream()
.mapToInt(Sample::responseMs)
.average().orElse(0))
.filter(avg -> avg > 200)
.findFirst();
System.out.println(firstSlow);
Optional.empty
Integrator 返回
false即可停止向下游推送,实现真正的短路,不会处理剩余元素。
启用预览特性
# 编译
javac --enable-preview --release 22 GathererDemo.java
# 运行
java --enable-preview GathererDemo
Maven 中在
maven-compiler-plugin配置<compilerArgs><arg>--enable-preview</arg></compilerArgs>。
易错场景
1. 忘记启用预览
// ❌ 直接编译运行,Gatherer 找不到或报错
java GathererDemo
// ✅ 必须加 --enable-preview
java --enable-preview GathererDemo
2. 混淆 Collector 与 Gatherer 的位置
// ❌ 把 Collector 当中间操作用
stream().collect(Collectors.toList()).filter(...) // 已终止,无法再流式操作
// ✅ Gatherer 是中间操作,后面还能继续链式
stream().gather(Gatherers.windowFixed(5)).map(...).collect(Collectors.toList())
3. windowFixed 残余窗口丢失
// ❌ 手动实现时忘记处理最后不足一组的元素
if (buffer.size() == windowSize) { emit(buffer); }
// 忘了: if (!buffer.isEmpty()) { emit(buffer); }
// ✅ Gatherers.windowFixed 自动处理残余窗口,无需手动关心
stream().gather(Gatherers.windowFixed(5))
面试考点
Q1:Gatherer 和 Collector 的本质区别是什么?
Gatherer 是中间操作,处理元素后推送到下游,后面可以继续
map、filter、collect等,支持惰性求值和短路;Collector 是终止操作,消费整个流后产出最终结果,无法再接流操作。Gatherer 让 Stream 终于拥有了可自定义的有状态中间操作能力。
Q2:Gatherer 的 Integrator 返回 false 意味着什么?
返回
false表示不再接收后续元素,实现 短路(short-circuiting)。这是 Gatherer 相比手动collect后处理的关键优势——可以在找到目标后立即停止,无需遍历全部数据。
Q3:windowSliding 和 windowFixed 有什么区别?
windowFixed(n)是 不重叠分组,元素分完即止,最后一组可能不足n个;windowSliding(n)是 滑动窗口,每次步进 1,窗口间有n-1个重叠元素。前者适合批量分片,后者适合趋势分析。
"Stream 有了 Gatherer,就像管道装了阀门——不再是只能一冲到底,而是随时可以截流、分流、再接流。"