乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • 第1章 Java概述与环境搭建

    • Java概述与环境搭建
    • Java语言概述
    • 解释型语言与编译型语言对比
    • JDK安装与配置
    • JDK、JRE、JVM 详解
    • HelloWorld程序详解
    • IDE 介绍
  • 第2章 标识符与基本数据类型

    • 章节导读
    • 变量概述
    • 常量概述
    • 基本类型与包装类
    • 字节型 byte
    • 短整型 short
    • 整型 int
    • 长整型 long
    • 单精度浮点型 float
    • 双精度浮点型 double
    • 字符型 char
    • 布尔型 boolean
    • 类型转换
  • 第3章 运算符与表达式

    • 章节导读
    • 算术运算符
    • 赋值运算符
    • 关系运算符
    • 逻辑运算符
    • 位运算符
    • 条件运算符
    • 运算符优先级
    • 表达式
  • 第4章 流程控制

    • 章节导读
    • 常见的程序运行流程
    • if-else 选择结构
    • switch 多分支选择
    • while 循环
    • do-while 循环
    • for 循环
    • break 与 continue
  • 第5章 数组

    • 章节导读
    • 一维数组
    • 多维数组
    • Arrays 工具类
  • 第6章 类与对象

    • 章节导读
    • 类与对象
    • 方法定义与调用
    • 构造方法
    • 封装
    • 访问修饰符
    • package 与 import
    • static 关键字
    • this 关键字
    • 参数传递 详解
    • 枚举
    • 成员内部类
    • 局部内部类
    • 静态内部类
    • 匿名内部类
  • 第7章 接口与继承

    • 章节导读
    • 继承
    • super 关键字
    • final 关键字
    • 多态
    • 向上转型与向下转型
    • 抽象类
    • 接口
    • 抽象类与接口对比
  • 第8章 注解

    • 章节导读
    • 注解基础
    • 元注解详解
    • 自定义注解
  • 第9章 常用类

    • 章节导读:Java 常用类
    • Object 类:万类之祖
    • 包装类:基本类型的对象化
    • String:不可变的字符串
    • StringBuffer:线程安全的可变字符串
    • StringBuilder:可变的字符串构建器
    • Math:数学运算工具类
    • Random:伪随机数生成器
    • 大数值运算 详解
    • 日期时间API 详解
  • 第10章 异常机制

    • 章节导读
    • 异常体系与分类
    • try-catch-finally
    • try-with-resources
    • throws 与 throw
    • 自定义异常
  • 第11章 泛型

    • 章节导读
    • 泛型基础
    • 通配符与PECS原则
    • 类型擦除
  • 第12章 集合框架

    • 章节导读
    • 集合框架概述
    • ArrayList
    • LinkedList
    • HashMap 详解
    • LinkedHashMap 详解
    • TreeMap 详解
    • HashSet
    • TreeSet 详解
    • TreeSet 与 Comparable
    • Collections 工具类详解
  • 第13章 IO流

    • 章节导读
    • IO流概述
    • 字节流
    • 字符流
    • 缓冲流
    • 转换流 详解
    • 序列化 详解
    • NIO与Files 详解
    • NIO与Files工具类
  • 第14章 多线程与并发

    • 第十六章 多线程与并发 —— 章节导读
    • 线程基础详解
    • synchronized 详解
    • Lock 与显式锁详解
    • volatile 详解
    • wait 与 notify 详解
    • ThreadLocal详解
    • 原子类详解
    • 并发工具类详解
    • 线程池详解
  • 第15章 反射

    • 章节导读
    • 反射概述与 Class 对象
    • Constructor 与对象创建
    • Field与Method详解
    • 反射应用详解
  • 第16章 JDK8新特性

    • 章节导读
    • Lambda 表达式
    • Stream API 基础
    • Stream API 高级详解
    • Optional 详解
    • 新日期时间API详解
  • 第17章 JDK9-11新特性

    • 章节导读
    • 模块化系统 — Project Jigsaw(JDK 9)
    • var 局部变量类型推断(JDK 10)
    • 集合工厂方法与增强(JDK 9 / 10 / 11)
    • 接口增强:private 方法(JDK 9)
    • Stream API 增强(JDK 9)
    • Optional 增强(JDK 9 / 10 / 11)
    • String 新增方法(JDK 11)
    • HTTP Client 与 Files 增强(JDK 11)
    • 直接运行 Java 源文件 — JEP 330(JDK 11)
  • 第18章 JDK12-17新特性

    • 章节导读
    • Switch 表达式(JDK 12 预览 / JDK 14 正式)
    • 文本块 Text Blocks(JDK 13 预览 / JDK 15 正式)
    • Records 记录类(JDK 14 预览 / JDK 16 正式)
    • 密封类 Sealed Classes(JDK 15 预览 / JDK 17 正式)
    • instanceof 模式匹配(JDK 14 预览 / JDK 16 正式)
    • Switch 模式匹配 — Pattern Matching for switch(JDK 17 预览 / JDK 21 正式)
    • Helpful NPE 与 String 增强(JDK 12 / JDK 14 / JDK 15)
    • Stream 增强(JDK 12 / JDK 16)
    • 日期时间增强 — Day Period 支持(JDK 16)
  • 第19章 JDK18-21新特性

    • 章节导读
    • 虚拟线程(JDK 19 预览 / JDK 20 第二预览 / JDK 21 正式)
    • 序列集合(JDK 21 正式)
    • Switch 模式匹配(JDK 17 预览 / JDK 18 第二预览 / JDK 20 第四预览 / JDK 21 正式)
    • Record 模式匹配(JDK 19 预览 / JDK 20 第二预览 / JDK 21 正式)
    • 未命名模式与变量(JDK 21 预览 / JDK 22 正式)
  • 第20章 JDK 22-25 新特性

    • 章节导读
    • 字符串模板(JDK 22 预览 / JDK 23 第二预览 / JDK 24 第三预览)
    • Stream Gatherers(JDK 22 预览 / JDK 24 第二预览)
    • 隐式声明类与实例方法(JDK 23 预览 / JDK 24 第二预览)
    • 原始类型模式匹配(JDK 24 预览)
  • 附录

    • Java 核心知识点
    • Java SE 专业术语
    • Java特性索引(JDK 8 → 25)

大翔:飞翔科技的实时监控大屏,响应时间指标每 5 秒一个窗口做聚合,用 Stream 怎么搞? 白歌:collect 分组只能拿到最终结果,中间的窗口滑动做不了。JDK 22 的 Stream Gatherers 就是干这个的——自定义中间操作。 小崔:Gatherer 和 Collector 有啥本质区别? 孔蓝:我猜……Collector 是终止操作,Gatherer 是中间操作?

Stream Gatherers(JDK 22 预览 / JDK 24 第二预览)

核心定义

概念定义
GathererStream 的自定义中间操作,可逐元素处理、可短路、可有状态
CollectorStream 的终止操作,将元素归约为最终结果
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,就像管道装了阀门——不再是只能一冲到底,而是随时可以截流、分流、再接流。"

上一页
字符串模板(JDK 22 预览 / JDK 23 第二预览 / JDK 24 第三预览)
下一页
隐式声明类与实例方法(JDK 23 预览 / JDK 24 第二预览)