Java Stream Gatherer 组合短窗口事件的实现步骤
来源:17golang原创
时间:2026-09-28 20:05:39 495浏览 收藏
组合短窗口事件时,先判断“窗口”是按元素数量还是按事件时间定义:每 50 条组成一批可直接使用 Gatherers.windowFixed(50);“同一秒内的事件组成一组”则需要自定义带状态的 Gatherer。下面的实现要求输入已经按 occurredAt 升序排列,只保留当前活动窗口,并在流结束时补发不足一窗的尾部事件。
官方 API:https://docs.oracle.com/en/java/javase/26/docs/api/java.base/java/util/stream/Gatherer.html
本文使用 Java 24 及以上版本的标准 Stream Gatherer API。短窗口状态依赖事件顺序,因此示例明确创建顺序 Gatherer;即使上游调用
parallel(),也不会凭空获得正确的并行时间窗口。
先用基线指标明确窗口目标
不要先写状态机,再猜窗口是否正确。先为示例固定一组可核对的数据:8 条事件,窗口宽度 1 秒,事件时间有序,窗口边界采用左闭右开区间 [start, start + width)。时间正好落在右边界的事件进入下一窗,时间跨度中没有事件时不生成空窗口。
| 指标 | 目标值 | 用途 |
|---|---|---|
| 输入事件数 | 8 | 确认没有重复或漏消费 |
| 窗口宽度 | 1 秒 | 定义事件时间边界 |
| 预期窗口数 | 3 | 确认切窗与尾窗刷新 |
| 预期窗口大小 | 3、2、3 | 确认边界事件归属 |
| 预期 value 合计 | 9、6、15 | 确认窗口内容未被污染 |
这些值来自示例输入,是确定性验收数据,不是性能压测结果。性能部分单独使用吞吐、分配量和峰值活动窗口大小衡量。
区分固定数量窗口与事件时间窗口
Oracle 的 Gatherers 已内置固定窗口和滑动窗口。windowFixed(n) 产生互不重叠的定长分组,最后一组可以少于 n;windowSliding(n) 产生重叠窗口。两者返回的窗口列表都不可修改,n 会抛出 IllegalArgumentException。
import java.util.List; import java.util.stream.Gatherers; import java.util.stream.Stream; List> batches = Stream.of("E1", "E2", "E3", "E4", "E5") // 每两个元素组成一个固定窗口,最后一个窗口允许不足两个元素 .gather(Gatherers.windowFixed(2)) .toList(); System.out.println(batches); // [[E1, E2], [E3, E4], [E5]]
固定数量适合批量写入、分页请求和限量聚合,却不能表达“1 秒内到达的事件”。时间窗口需要读取每条事件的时间戳,并把窗口起点保存在私有状态中。
用状态、积分器和完成器实现短时间窗口
Gatherer 由初始化器、积分器、可选组合器和完成器构成。本例选择 ofSequential:初始化器创建一个窗口状态;积分器接收事件并判断是否越界;完成器在上游结束后刷新尾窗。输出前必须复制缓冲区,避免后续 clear() 改写已经发给下游的窗口。

import java.time.Duration;
import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.stream.Gatherer;
record Event(String type, Instant occurredAt, int value) {}
record EventWindow(
Instant start,
Instant exclusiveEnd,
List events,
int totalValue
) {}
final class WindowState {
Instant start;
final List buffer = new ArrayList();
}
static Gatherer shortTimeWindows(Duration width) {
Objects.requireNonNull(width, "width");
if (width.isZero() || width.isNegative()) {
throw new IllegalArgumentException("窗口宽度必须大于 0");
}
return Gatherer.ofSequential(
WindowState::new,
Gatherer.Integrator.ofGreedy((state, event, downstream) -> {
Objects.requireNonNull(event, "event");
Objects.requireNonNull(event.occurredAt(), "event.occurredAt");
if (state.start == null) {
// 第一条事件确定首个窗口起点
state.start = event.occurredAt();
} else if (event.occurredAt().isBefore(state.start)) {
// 本实现要求输入按事件时间升序,乱序数据必须先处理
throw new IllegalArgumentException("事件时间未按升序排列");
}
Instant boundary = state.start.plus(width);
if (!state.buffer.isEmpty() && !event.occurredAt().isBefore(boundary)) {
// 右边界是开区间;先发出旧窗口,再从当前事件开启新窗口
if (!downstream.push(snapshot(state, width))) {
return false;
}
state.buffer.clear();
state.start = event.occurredAt();
}
state.buffer.add(event);
return true;
}),
(state, downstream) -> {
// 上游结束后补发不足一个完整窗口的尾部事件
if (!state.buffer.isEmpty()) {
downstream.push(snapshot(state, width));
}
}
);
}
static EventWindow snapshot(WindowState state, Duration width) {
// 复制列表,防止状态缓冲区复用时污染已经输出的窗口
List immutableEvents = List.copyOf(state.buffer);
int total = immutableEvents.stream().mapToInt(Event::value).sum();
return new EventWindow(
state.start,
state.start.plus(width),
immutableEvents,
total
);
}
这里没有组合器,因为基于相邻事件和窗口起点的状态不容易在任意分区边界上无损合并。顺序语义比“表面上调用了并行流”更重要。
接入 Stream.gather 并核对结果
下面的输入故意包含窗口边界和时间空洞。第三个窗口从 3.3 秒开始,中间没有事件的区间不会生成空列表。
import java.time.Duration;
import java.time.Instant;
import java.util.List;
Instant base = Instant.parse("2026-09-28T10:00:00Z");
List events = List.of(
new Event("view", base.plusMillis(0), 3),
new Event("click", base.plusMillis(200), 4),
new Event("view", base.plusMillis(700), 2),
new Event("pay", base.plusMillis(1100), 5),
new Event("view", base.plusMillis(1700), 1),
new Event("click", base.plusMillis(3300), 7),
new Event("view", base.plusMillis(3500), 2),
new Event("pay", base.plusMillis(3900), 6)
);
List windows = events.stream()
// 一秒短窗口;输入已经按 occurredAt 升序排列
.gather(shortTimeWindows(Duration.ofSeconds(1)))
.toList();
windows.forEach(window -> System.out.printf(
"start=%s, size=%d, total=%d%n",
window.start(), window.events().size(), window.totalValue()
));
结果应有 3 个窗口,事件数分别为 3、2、3,聚合值分别为 9、6、15。若窗口数少了,先检查完成器是否刷新尾窗;若边界事件落错组,检查条件是否保持“时间早于右边界”而不是“早于或等于”。
按内存和吞吐指标比较改动
Gatherer 的直接价值不是保证所有场景都更快,而是让中间状态的边界更清晰。先把全部事件收进列表再分组,中间工作集约为 O(N);逐窗向下游消费时,Gatherer 的私有状态约为 O(W),其中 W 是峰值窗口事件数。

注意:示例最后调用 toList(),因此所有输出窗口仍会被保留。若要体现流式内存优势,应让下游逐窗写数据库、发送消息或计算摘要,而不是再次整体收集。可按同一份输入分别测量基线和 Gatherer 实现:
# 使用 JMH 的 GC 分析器记录吞吐、每次操作分配字节数和 GC 次数 java -jar benchmarks.jar WindowBenchmark -prof gc # 固定堆大小,避免两组测试因堆配置不同而失去可比性 java -Xms512m -Xmx512m -jar benchmarks.jar WindowBenchmark -prof gc
| 对比项 | 基线:先整体收集 | Gatherer:逐窗消费 |
|---|---|---|
| 中间工作集 | 随全部事件数 N 增长 | 随峰值窗口 W 增长 |
| 尾窗 | 由分组代码额外处理 | 由 finisher 统一刷新 |
| 顺序保证 | 取决于手写循环 | 由顺序 Gatherer 明确约束 |
| 吞吐与分配量 | 必须在相同 JDK、数据分布、堆大小和消费端条件下用 JMH 实测 | |
如果需要容量估算,可以用“峰值窗口事件数 × 单事件近似占用”估算活动缓冲区,但 Java 对象头、引用压缩和字段布局会影响实际字节数,最终仍以 JMH 的 gc.alloc.rate.norm 与 JFR/GC 数据为准。
处理乱序、并行和空窗口边界
- 乱序事件:示例直接拒绝早于当前窗口起点的事件。若允许迟到数据,需要先按时间排序,或使用带水位线、允许迟到策略和状态存储的流处理框架。
- 并行流:
ofSequential明确关闭任意分区合并。时间窗口要并行化,必须设计能处理分区边界的 combiner,不能只添加parallel()。 - 空窗口:当前实现只输出含事件的窗口。报表若要求补齐空时间段,应在窗口结果之后按时间范围补齐。
- 窗口锚点:当前窗口以每个活动段的第一条事件为起点。若业务要求按整秒、整分钟对齐,需要先把时间戳截断到统一边界。
- 下游停止:积分器转发
downstream.push的返回值,使limit、findFirst等短路操作可以停止继续消费。
结果检查清单
- 固定数量分组优先使用
Gatherers.windowFixed,不要重复造状态机。 - 时间窗口先明确排序、区间开闭、窗口锚点和空窗口策略。
- 输出使用
List.copyOf,不要把可复用缓冲区直接交给下游。 - 完成器必须刷新尾窗,输入为空时不输出窗口。
- 性能结论同时观察吞吐、标准化分配量、GC 次数和峰值窗口大小。
- 想降低整体内存时,让下游逐窗消费;末尾
toList()会重新保留全部输出。
相关问题
windowFixed 与 windowSliding 有什么区别?
windowFixed 的窗口互不重叠;windowSliding 每次移出最早元素并加入新元素,会产生重叠窗口。两者都按元素数量,而不是事件时间切分。
为什么一定要写 finisher?
最后一个窗口可能在上游结束前一直没有遇到下一条越界事件。没有 finisher,这部分事件就不会被输出。
这个 Gatherer 能处理无限流吗?
Stream Gatherer 可以逐步处理元素,但普通 Java Stream 本身不是带检查点、水位线和持久状态的分布式流系统。持续运行、迟到数据和故障恢复需要更完整的流处理设计。
-
297 收藏
-
Golang · Go问答 | 1个月前 | 并发 · golang · csv · 数据处理 · Go问答 · 数据导入 encoding/csv csv.Reader Go问答 ReuseRecord 切片复用465 收藏
-
146 收藏
-
367 收藏
-
238 收藏
-
244 收藏
-
361 收藏
-
342 收藏
-
182 收藏
-
291 收藏
-
145 收藏
-
419 收藏
-
494 收藏
-
223 收藏
-
152 收藏
-
394 收藏
-
199 收藏
-
- 前端进阶之JavaScript设计模式
- 设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
- 立即学习 543次学习
-
- GO语言核心编程课程
- 本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
- 立即学习 516次学习
-
- 简单聊聊mysql8与网络通信
- 如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
- 立即学习 500次学习
-
- JavaScript正则表达式基础与实战
- 在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
- 立即学习 487次学习
-
- 从零制作响应式网站—Grid布局
- 本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
- 立即学习 485次学习