登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  文章 >  java教程

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() 改写已经发给下游的窗口。

Event 输入、WindowState 状态、Integrator 边界判断、Finisher 尾窗刷新与 EventWindow 输出的关系图
图1:短时间窗口 Gatherer 的静态结构图;状态只保存当前窗口,Integrator 处理边界,Finisher 负责尾窗。
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 是峰值窗口事件数。

全部事件列表与 Gatherer 活动窗口缓冲的内存边界对比图
图2:两种内存边界的静态对比图;Gatherer 可把中间工作集收敛到活动窗口,但终端 toList 仍会保留全部窗口。

注意:示例最后调用 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 本身不是带检查点、水位线和持久状态的分布式流系统。持续运行、迟到数据和故障恢复需要更完整的流处理设计。

声明:本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>