Java——Stream Gatherers 自定义管道
《Stream API 常见误用》讲的是内建操作怎么用错;这篇讲内建操作不够用时怎么办。
过去要在流中间做”滑动窗口””按状态切分””限流并发”这类事,答案通常是”收集成 List 再用 for 循环”,或者硬塞一个Collector——后者只能在终结处用,中间管道依然插不进去。
JDK 24 转正的 JEP 485 补上了这块:Stream.gather(Gatherer)是中间操作,和map、filter处在同一个位置上。全部代码在 JDK 25 上编译运行,不需要预览开关。
实验环境
1 | $ java -version |
一、Gatherer 在管道里的位置
JEP 485 的定位说得很清楚:
Stream::gather(Gatherer)is to intermediate operations whatStream::collect(Collector)is to terminal operations.
也就是:Collector 之于终结操作,等于 Gatherer 之于中间操作。签名如下(用反射从 JDK 25 打出来的):
1 | Stream.gather: public default <R> Stream<R> java.util.stream.Stream.gather(Gatherer<? super T, ?, R>) |
一个 Gatherer 由四个部件构成,任何一个都可以省略:
| 部件 | 类型 | 作用 |
|---|---|---|
| initializer | Supplier<A> |
创建私有状态(每个流/每个分区一份) |
| integrator | Integrator<A, T, R> |
处理元素,决定是否继续接收 |
| combiner | BinaryOperator<A> |
并行时合并两个分区的状态 |
| finisher | BiConsumer<A, Downstream<R>> |
流结束时把残留状态推给下游 |
其中 integrator 是唯一必须的。Integrator 的形态是:
1 | boolean integrate(A state, T element, Downstream<? super R> downstream) |
返回值是是否愿意继续接收新元素——false 表示”我不要后面的了”,这是流式短路的基础。
二、五个内建 gatherer
java.util.stream.Gatherers 这个工厂类只提供五个方法(反射枚举,JDK 25):
1 | ### Gatherers 工厂类 |
用同一份输入跑一遍:
1 | var nums = List.of(1, 2, 3, 4, 5, 6, 7, 8, 9); |
1 | windowFixed(3): [[1, 2, 3], [4, 5, 6], [7, 8, 9]] |
几个容易搞混的点:
windowFixed与windowSliding的区别:前者不重叠、窗口数是⌈n/size⌉;后者滑动步长为 1、窗口数是n - size + 1。凑不满的尾部窗口会怎样?windowFixed(3)对 10 个元素会给出最后一组[10]——短窗口照给,不会丢弃。这一点在批量提交场景里要特别小心(见 4.3 节)。fold与scan的区别:fold把整个流塌缩成一个值(下游只收到一个元素),scan是”前缀和”,每个输入对应一个输出。fold常常配合findFirst()取值,因为它的下游只会有一个元素。
2.1 scan 可以配合 limit 做”无限流前缀”
1 | System.out.println("scan 配合 limit: " + |
1 | scan 配合 limit: [1, 3, 6, 10, 15] |
Stream.iterate 是无限的,但整条管道是惰性的:limit(5) 决定了下游只拉 5 个元素,scan 也只处理 5 个。
三、惰性与短路:两个可验证的行为
这两个性质经常被口头描述,但都能用几行代码验证。
3.1 惰性:不消费就不执行
1 | var pipeline = Stream.of("a", "b", "c").gather(Gatherer.<String, List<String>, List<String>>ofSequential( |
1 | === 惰性验证:下面构造了管道但不消费 === |
构造管道时一行 处理元素 都没打印,toList() 之后才逐个出现——gather 是中间操作,行为与 map 一致。
3.2 短路:integrator 的返回值决定还能收多少个
1 | final int[] n = {0}; |
1 | === 短路验证:拿到 3 个元素就停 === |
输入是无限流,limit(10) 允许到 10 个,但 integrator 第 3 次返回 false 之后就不再被调用了。注意这与 Integrator.ofGreedy 的区别:greedy 的 integrator 被明确假设”不会短路”,运行时可以据此优化;javadoc 原文是”Gatherers whose integrator is an instance of Gatherer.Integrator.Greedy can be assumed not to short-circuit”。要不要短路,由你用哪个工厂方法声明。
四、自己写 gatherer:三个真实用例
内建的五个覆盖不了的状态机,就得自己写。工厂方法有两组(反射枚举的完整清单):
1 | Gatherer.of(...) —— 可并行的版本(提供 combiner 时) |
有状态、且状态在分区之间无法独立合并的,用 ofSequential;能给出合法 combiner 的,用 of 换取并行能力。
4.1 相邻去重:丢弃连续重复
需求:日志事件流里,同一用户连续重复的同一动作只保留第一条。
1 | record Event(long ts, String user, String action) {} |
1 | 输入: alice:login, alice:login, alice:view, bob:login, bob:login, bob:logout, alice:login |
注意状态管理的小细节:状态里只保留”最后一条通过的记录”,所以内存是 O(1),与流的长度无关。如果写成把全部事件缓存进 List 再比对,就退化成了一次全量收集——gatherer 的价值在于让这种”只记一点点状态”的实现方式变得自然。
4.2 会话切分:遇到新用户就断开窗口
需求:把事件流按”连续同一个用户”切分成会话,每次切换就吐出一个完整的会话。
1 | var sessions = events.stream().gather(Gatherer.<Event, ArrayList<Event>, List<Event>>ofSequential( |
1 | 会话切分: [alice×3, bob×3, alice×1] |
这个例子展示了两个容易漏的点:
finisher不能忘。流结束时状态里还剩最后一段没有切换点,只能靠 finisher 推出去——不写它,最后一个会话就丢了。downstream.push(...)的返回值要检查。它返回false表示下游已经不再接受元素(例如后面接了limit),此时应该停止工作而不是继续跑。
4.3 批量提交:windowFixed 的尾部行为
1 | var batches = IntStream.rangeClosed(1, 10).boxed().gather(Gatherers.windowFixed(3)).toList(); |
1 | 每 3 条一批: [[1, 2, 3], [4, 5, 6], [7, 8, 9], [10]] |
最后一批只有 1 个元素。批量写库/批量发消息时要显式处理这种短批次——很多下游接口对”批量大小不固定”是不接受的,此时要么在 gatherer 里丢弃短尾,要么用 fold 之类的自定义逻辑补齐。
4.4 组合:andThen
多个 gatherer 可以拼起来,两种写法等价(JEP 485 原文示例的两种形态):
1 | source.gather(a).gather(b).gather(c).collect(...); // 串联三次 gather |
区别在于:andThen 组合出来的是一个 Gatherer 对象,可以存成常量、复用、测试;连续调用 gather 则更直观。组合的短路语义是从右往左传播的:下游不要了,最上游也会停,规范里有一句明确的话——“elements from earlier partitions may be discarded if processing an earlier partition short-circuits”。
五、并行流下的注意事项
gather 用在 parallelStream() 上时,合并逻辑由 combiner 决定,规则比较硬:
- 没有 combiner 的 gatherer 在并行流上会被当作顺序处理,不会报错,也不会有并行收益
- 有 combiner 时,每个分区各自持有独立状态(这正是 initializer 必须存在的原因),最后合并
Gatherers.windowFixed这类内建实现自带 combiner,但窗口不会跨越分区边界——如果你期望”全局每 3 个一批”,并行流下的分块顺序会让结果与顺序流不同(顺序流保证批次边界固定)
这就是选择 of 还是 ofSequential 的实际后果:顺序语义敏感的 gatherer 别用 of 硬凑并行,用 ofSequential 声明清楚,让运行时按顺序执行。
六、mapConcurrent:内建里最实用的一个
Gatherers.mapConcurrent(limit, fn) 的 javadoc 说它”executes a function concurrently with a configured level of max concurrency, using virtual threads”——内部用虚拟线程,所以并发度给到几十上百也不会像平台线程池那样爆内存。
实测:12 个下游调用、每个 200ms、并发度 4。
1 | var results = ids.stream() |
1 | === mapConcurrent(4):12 个下游调用,最多 4 个并发 === |
- 峰值并发精确停在 4(用
AtomicInteger在函数内部计数验证) - 总耗时 623ms,接近理论下限 600ms
- 输出顺序与输入一致:这是它比”自己开线程池提交”更好用的地方——并发执行但保序返回
6.1 与 parallelStream 的区别
同一批任务换个写法:
1 | === 对照:parallel() 与 mapConcurrent 的并发度 === |
本机可用核数 10、公共 ForkJoinPool 并行度 9,实测 parallel() 的峰值并发是 10(提交线程也会参与计算)。两者的差别不是”谁快”:
parallelStream().map() |
mapConcurrent(4, fn) |
|
|---|---|---|
| 并发度 | 由 CPU 核数决定(这里 10) | 由你指定(这里 4) |
| 适合 | CPU 密集计算 | IO 密集调用(下游限流、HTTP) |
| 底层线程 | 公共 ForkJoinPool(平台线程) | 虚拟线程 |
| 与其它并行流的关系 | 共用公共池,可能互相影响 | 独立,不干扰公共池 |
调用下游服务用 mapConcurrent 并显式给并发度,别用 parallelStream。前者不会污染公共 ForkJoinPool,也不会因为机器核数变化而改变对下游的压力。
七、版本与陷阱
| JDK | JEP | 状态 |
|---|---|---|
| 22 | JEP 461 | 预览 |
| 23 | JEP 473 | 第二次预览 |
| 24 | JEP 485 | 正式(Status: Closed / Delivered、Release: 24) |
写代码时实测遇到的两类坑:
坑一:Integrator 的泛型推断要显式指定。 下面这行编译不过:
1 | // 编译错误:对于 of((Integer element, Downstream<? super Integer> downstream) -> ...) 找不到合适的方法 |
Gatherer.of(Integrator<Void, T, R>) 里的状态类型是 Void,编译器需要类型见证(type witness)才能推断:
1 | .gather(Gatherer.<Integer, Integer>of( |
坑二:ofSequential 的类型参数顺序是 <T, A, R>(元素类型、状态类型、结果类型),不是 <T, R, A>。写错了编译器会给出”ArrayList 无法转换为 Event”这类看起来毫不相关的错误——上面 4.1 节的例子就是修过一次之后的写法。
总结
gather是中间操作:Collector对应终结处的自定义,Gatherer对应管道中间的自定义- JDK 25 内建只有五个:
windowFixed、windowSliding、fold、scan、mapConcurrent windowFixed会给出不足长度的尾窗口(10 个元素按 3 分批 → 最后一批只有 1 个),批量提交时必须处理- 自定义 gatherer 的四个部件里
integrator是必需的;finisher负责把状态里的残留推给下游(会话切分忘了它就会丢最后一段) integrator的返回值控制短路:返回false之后元素不再进入;用ofGreedy声明”永不短路”可以让运行时优化- 有状态且无法跨分区合并的 gatherer 用
ofSequential,否则在并行流上会得到与顺序流不同的结果 mapConcurrent内部用虚拟线程、保序、并发度由你指定,是替代parallelStream做下游调用的正解- 泛型推断是主要的使用摩擦点:短路型 integrator 需要类型见证,
ofSequential的类型参数顺序是<T, A, R>
参考资料
- JEP 485: Stream Gatherers
- JEP 461: Stream Gatherers (Preview)
- JDK 25 API:java.util.stream.Gatherers
- JDK 25 API:java.util.stream.Gatherer
系列索引:Java 系列,语言特性与运行时的长文集






