Stream API 常见误用》讲的是内建操作怎么用错;这篇讲内建操作不够用时怎么办。
过去要在流中间做”滑动窗口””按状态切分””限流并发”这类事,答案通常是”收集成 List 再用 for 循环”,或者硬塞一个 Collector——后者只能在终结处用,中间管道依然插不进去。
JDK 24 转正的 JEP 485 补上了这块:Stream.gather(Gatherer) 是中间操作,和 mapfilter 处在同一个位置上。全部代码在 JDK 25 上编译运行,不需要预览开关。

实验环境

1
2
3
4
$ java -version
java 25 2025-09-16 LTS
Java(TM) SE Runtime Environment (build 25+37-LTS-3491)
Java HotSpot(TM) 64-Bit Server VM (build 25+37-LTS-3491, mixed mode, sharing)

一、Gatherer 在管道里的位置

JEP 485 的定位说得很清楚:

Stream::gather(Gatherer) is to intermediate operations what Stream::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
2
3
4
5
6
### Gatherers 工厂类
Gatherer windowFixed[int]
Gatherer windowSliding[int]
Gatherer fold[Supplier, BiFunction]
Gatherer scan[Supplier, BiFunction]
Gatherer mapConcurrent[int, Function]

用同一份输入跑一遍:

1
2
3
4
5
6
var nums = List.of(1, 2, 3, 4, 5, 6, 7, 8, 9);

System.out.println("windowFixed(3): " + nums.stream().gather(Gatherers.windowFixed(3)).toList());
System.out.println("windowSliding(3): " + nums.stream().gather(Gatherers.windowSliding(3)).toList());
System.out.println("fold(求和): " + nums.stream().gather(Gatherers.fold(() -> 0, Integer::sum)).findFirst().orElseThrow());
System.out.println("scan(前缀和): " + nums.stream().gather(Gatherers.scan(() -> 0, Integer::sum)).toList());
1
2
3
4
windowFixed(3):   [[1, 2, 3], [4, 5, 6], [7, 8, 9]]
windowSliding(3): [[1, 2, 3], [2, 3, 4], [3, 4, 5], [4, 5, 6], [5, 6, 7], [6, 7, 8], [7, 8, 9]]
fold(求和): 45
scan(前缀和): [1, 3, 6, 10, 15, 21, 28, 36, 45]

几个容易搞混的点:

  • windowFixedwindowSliding 的区别:前者不重叠、窗口数是 ⌈n/size⌉;后者滑动步长为 1、窗口数是 n - size + 1。凑不满的尾部窗口会怎样?windowFixed(3) 对 10 个元素会给出最后一组 [10]——短窗口照给,不会丢弃。这一点在批量提交场景里要特别小心(见 4.3 节)。
  • foldscan 的区别fold 把整个流塌缩成一个值(下游只收到一个元素),scan 是”前缀和”,每个输入对应一个输出。
  • fold 常常配合 findFirst() 取值,因为它的下游只会有一个元素。

2.1 scan 可以配合 limit 做”无限流前缀”

1
2
System.out.println("scan 配合 limit:  " +
Stream.iterate(1, i -> i + 1).gather(Gatherers.scan(() -> 0, Integer::sum)).limit(5).toList());
1
scan 配合 limit:  [1, 3, 6, 10, 15]

Stream.iterate 是无限的,但整条管道是惰性的:limit(5) 决定了下游只拉 5 个元素,scan 也只处理 5 个。

三、惰性与短路:两个可验证的行为

这两个性质经常被口头描述,但都能用几行代码验证。

3.1 惰性:不消费就不执行

1
2
3
4
5
6
var pipeline = Stream.of("a", "b", "c").gather(Gatherer.<String, List<String>, List<String>>ofSequential(
ArrayList::new,
(state, element, downstream) -> { System.out.println(" 处理元素 " + element); state.add(element); return true; },
(state, downstream) -> downstream.push(state)));
System.out.println(" 管道已构造,尚未看到任何「处理元素」输出");
System.out.println(" 消费后: " + pipeline.toList());
1
2
3
4
5
6
=== 惰性验证:下面构造了管道但不消费 ===
管道已构造,尚未看到任何「处理元素」输出
处理元素 a
处理元素 b
处理元素 c
消费后: [[a, b, c]]

构造管道时一行 处理元素 都没打印,toList() 之后才逐个出现——gather 是中间操作,行为与 map 一致。

3.2 短路:integrator 的返回值决定还能收多少个

1
2
3
4
5
6
7
8
9
10
final int[] n = {0};
var take3 = Stream.iterate(1, i -> i + 1)
.gather(Gatherer.<Integer, Integer>of(
(Void state, Integer element, Gatherer.Downstream<? super Integer> downstream) -> {
n[0]++;
System.out.println(" tryAccept " + element + " (第 " + n[0] + " 次)");
downstream.push(element);
return n[0] < 3; // false = 后面的元素被拒绝
}))
.limit(10).toList();
1
2
3
4
5
=== 短路验证:拿到 3 个元素就停 ===
tryAccept 1 (第 1 次)
tryAccept 2 (第 2 次)
tryAccept 3 (第 3 次)
结果: [1, 2, 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
2
3
4
5
6
7
8
9
10
Gatherer.of(...)            —— 可并行的版本(提供 combiner 时)
of[Integrator]
of[Integrator, BiConsumer] // + finisher
of[Supplier, Integrator, BinaryOperator, BiConsumer]
of[Supplier, Integrator, BinaryOperator] // ?
Gatherer.ofSequential(...) —— 声明为顺序专用
ofSequential[Integrator]
ofSequential[Integrator, BiConsumer]
ofSequential[Supplier, Integrator]
ofSequential[Supplier, Integrator, BiConsumer]

有状态、且状态在分区之间无法独立合并的,用 ofSequential;能给出合法 combiner 的,用 of 换取并行能力。

4.1 相邻去重:丢弃连续重复

需求:日志事件流里,同一用户连续重复的同一动作只保留第一条。

1
2
3
4
5
6
7
8
9
10
11
12
13
record Event(long ts, String user, String action) {}

var deduped = events.stream().gather(Gatherer.<Event, ArrayList<Event>, Event>ofSequential(
ArrayList::new,
(state, e, downstream) -> {
if (state.isEmpty() || !state.getLast().user().equals(e.user())
|| !state.getLast().action().equals(e.action())) {
state.add(e);
return downstream.push(e);
}
return true; // 与上一条相同:丢弃但不终止
},
(state, downstream) -> {})).toList();
1
2
输入: alice:login, alice:login, alice:view, bob:login, bob:login, bob:logout, alice:login
相邻去重后: [alice:login, alice:view, bob:login, bob:logout, alice:login]

注意状态管理的小细节:状态里只保留”最后一条通过的记录”,所以内存是 O(1),与流的长度无关。如果写成把全部事件缓存进 List 再比对,就退化成了一次全量收集——gatherer 的价值在于让这种”只记一点点状态”的实现方式变得自然。

4.2 会话切分:遇到新用户就断开窗口

需求:把事件流按”连续同一个用户”切分成会话,每次切换就吐出一个完整的会话。

1
2
3
4
5
6
7
8
9
10
11
12
var sessions = events.stream().gather(Gatherer.<Event, ArrayList<Event>, List<Event>>ofSequential(
ArrayList::new,
(state, e, downstream) -> {
if (!state.isEmpty() && !state.getLast().user().equals(e.user())) {
var closed = List.copyOf(state);
state.clear();
if (!downstream.push(closed)) return false; // 下游不要了就停
}
state.add(e);
return true;
},
(state, downstream) -> { if (!state.isEmpty()) downstream.push(List.copyOf(state)); })).toList();
1
会话切分: [alice×3, bob×3, alice×1]

这个例子展示了两个容易漏的点:

  1. finisher 不能忘。流结束时状态里还剩最后一段没有切换点,只能靠 finisher 推出去——不写它,最后一个会话就丢了。
  2. 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
2
source.gather(a).gather(b).gather(c).collect(...);          // 串联三次 gather
source.gather(a.andThen(b).andThen(c)).collect(...); // 组合成一个 gatherer

区别在于: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
2
3
var results = ids.stream()
.gather(Gatherers.mapConcurrent(4, id -> call(id, 200)))
.toList();
1
2
3
4
=== mapConcurrent(4):12 个下游调用,最多 4 个并发 ===
结果数 = 12,峰值并发 = 4,总耗时 = 623ms
顺序保持 = [item-1 -> 200, item-2 -> 200, item-3 -> 200, item-4 -> 200] ...
串行需要 2400ms,mapConcurrent(4) 理论下限 600ms
  • 峰值并发精确停在 4(用 AtomicInteger 在函数内部计数验证)
  • 总耗时 623ms,接近理论下限 600ms
  • 输出顺序与输入一致:这是它比”自己开线程池提交”更好用的地方——并发执行但保序返回

6.1 与 parallelStream 的区别

同一批任务换个写法:

1
2
=== 对照:parallel() 与 mapConcurrent 的并发度 ===
parallel() 峰值并发 = 10(公共 ForkJoinPool),耗时 = 412ms

本机可用核数 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 / DeliveredRelease: 24

写代码时实测遇到的两类坑:

坑一:Integrator 的泛型推断要显式指定。 下面这行编译不过:

1
2
3
// 编译错误:对于 of((Integer element, Downstream<? super Integer> downstream) -> ...) 找不到合适的方法
var take3 = Stream.iterate(1, i -> i + 1).gather(Gatherer.of(
(Integer element, Gatherer.Downstream<? super Integer> downstream) -> { ... }));

Gatherer.of(Integrator<Void, T, R>) 里的状态类型是 Void,编译器需要类型见证(type witness)才能推断:

1
2
.gather(Gatherer.<Integer, Integer>of(
(Void state, Integer element, Gatherer.Downstream<? super Integer> downstream) -> { ... }))

坑二:ofSequential 的类型参数顺序是 <T, A, R>(元素类型、状态类型、结果类型),不是 <T, R, A>。写错了编译器会给出”ArrayList 无法转换为 Event”这类看起来毫不相关的错误——上面 4.1 节的例子就是修过一次之后的写法。

总结

  • gather中间操作Collector 对应终结处的自定义,Gatherer 对应管道中间的自定义
  • JDK 25 内建只有五个:windowFixedwindowSlidingfoldscanmapConcurrent
  • windowFixed 会给出不足长度的尾窗口(10 个元素按 3 分批 → 最后一批只有 1 个),批量提交时必须处理
  • 自定义 gatherer 的四个部件里 integrator 是必需的;finisher 负责把状态里的残留推给下游(会话切分忘了它就会丢最后一段)
  • integrator 的返回值控制短路:返回 false 之后元素不再进入;用 ofGreedy 声明”永不短路”可以让运行时优化
  • 有状态且无法跨分区合并的 gatherer 用 ofSequential,否则在并行流上会得到与顺序流不同的结果
  • mapConcurrent 内部用虚拟线程、保序、并发度由你指定,是替代 parallelStream 做下游调用的正解
  • 泛型推断是主要的使用摩擦点:短路型 integrator 需要类型见证,ofSequential 的类型参数顺序是 <T, A, R>

参考资料

系列索引:Java 系列,语言特性与运行时的长文集