Stream API 从 Java 8 到现在,中间操作就那么几个:filter、map、flatMap、distinct、sorted、peek、limit、skip、takeWhile、dropWhile。用了十年,局限性大家心里都有数。
最典型的痛点是有状态的、需要跨元素判断的操作写不了。比如「把连续重复的元素去重」,distinct 是全流去重,做不到只去相邻的重复;「找出一段连续满足条件的区间」也没有直接对应的操作,只能写 for 循环或者用一些很别扭的 reduce 技巧。
JDK 24 正式落地的 Stream Gatherers(JEP 485)就是为了填这个坑。它不是替换 Collector,而是在中间操作这一层开了一个口子,让你能自己造操作符。
一、先搞清楚 Gatherer 和 Collector 的区别
这两个名字容易被搞混,但它们的定位完全不同。
Collector 是终点操作,stream.collect(...) 之后流就没了,你拿到的是一个最终结果,比如一个 List、一个 Map、一个数字。它没法再往后接 filter、再往后接 map。
Gatherer 是中间操作,用法是 stream.gather(myGatherer),返回的还是个 Stream,后面可以接着接其他操作。它和 map、filter 是一个层级的。
这个区别听着简单,但影响很大。正因为 Gatherer 是中间操作,它才能和短路操作配合。举个直观的例子:
// Collector 版本:必须把整个流跑完才能拿到结果
List<List<Integer>> runs = data.stream().collect(collectRuns());
List<Integer> first = runs.get(0); // 拿到第一个,但前面的都白算了
// Gatherer 版本:findFirst 一触发,后面的数据根本不会被遍历
List<Integer> first = data.stream()
.gather(runsWhere(v -> v > 80))
.findFirst()
.orElse(List.of());
第一个版本不管你要不要,整个流都必须跑完。第二个版本如果第一个区间出现在流的前 10%,后面 90% 的数据压根不会被读取。在处理大文件、网络流或者数据库游标的时候,这个差别不是一点半点。
二、Gatherer 长什么样
核心接口 Gatherer<T, A, R> 三个类型参数:
T,输入元素类型A,内部状态类型R,输出元素类型
它有四个组成部分:
initializer,一个 Supplier,用来为每次流的执行创建一份独立的状态对象integrator,核心逻辑,每来一个元素调一次,决定要不要往下游推数据combiner,只有并行流才需要,用来合并两个分段的状态finisher,流走到尽头时调用一次,用来处理状态里剩下的东西
日常写的 90% 场景只需要前两个,JDK 提供了便捷工厂方法 Gatherer.ofSequential(...),专门用于顺序流。
三、案例一:连续去重
先说个最简单的。业务背景是一批日志按时间排序,同一台设备短时间内会重复上报同一个状态码,需要把相邻的重复项去掉,但不相邻的重复项要保留。
用 distinct 是不行的,它会把所有重复都干掉。之前同事的写法是先转成 List,再手动遍历建一个新的,二十行代码,全是下标和边界判断。
用 Gatherer 写出来是这样:
import java.util.Objects;
import java.util.stream.Gatherer;
import java.util.stream.Stream;
public final class Gatherers {
private Gatherers() {}
private static final class DedupState<T> {
T last;
boolean hasLast;
}
public static <T> Gatherer<T, DedupState<T>, T> dedupConsecutive() {
return Gatherer.ofSequential(
DedupState::new,
(DedupState<T> state, T element, Gatherer.Downstream<? super T> downstream) -> {
if (!state.hasLast || !Objects.equals(state.last, element)) {
state.last = element;
state.hasLast = true;
downstream.push(element);
}
return true;
}
);
}
}
用起来:
List<String> raw = List.of("OK", "OK", "OK", "WARN", "OK", "OK", "ERROR", "ERROR");
List<String> cleaned = raw.stream()
.gather(dedupConsecutive())
.toList();
// 结果:[OK, WARN, OK, ERROR]
几个细节值得单独拎出来讲。
状态类的泛型必须显式声明。
很多人第一次写会尝试用匿名内部类,或者直接传一个 new Object() { T last; }。Java 的泛型推断在这种嵌套场景下会失效,编译报错提示 cannot infer type arguments。老老实实定义一个静态嵌套类,虽然多写几行,但后面维护的时候清楚得多。
状态不能存在 Gatherer 实例里。
这是最容易犯的错。假设你写成了这样:
// 错误示范
public static Gatherer<String, ?, String> dedupBroken() {
T[] holder = new Object[1]; // 状态存在闭包里
return Gatherer.ofSequential(
() -> null,
(s, element, downstream) -> { ... }
);
}
问题在于 Gatherer 实例通常是复用的,被定义成 static final 常量很常见。如果状态存在闭包里,同一个 Gatherer 被两个流同时用,或者被同一个流跑两遍,状态就会串。第一次跑完之后 last 还留着值,第二次跑的开头会莫名其妙把第一个元素吃掉。
initializer 存在的意义就是给每次流的执行提供一个干净的状态对象。规范上说初始化的时机是实现相关的,可能在整个流开始时调一次,也可能在并行流的每个分段开始时各调一次。你唯一应该依赖的是「这个 Supplier 返回的对象是这个流独享的」。
integrator 的返回值是「还能不能继续接收元素」。
返回 true 表示继续,返回 false 表示到此为止,后面的元素不用再送过来了。上面这个去重逻辑永远不会提前结束,所以无条件返回 true。
四、案例二:提取所有「连续达标」的区间
这个场景更实际一点。有一批按时间排序的 CPU 使用率采样,需要把连续超过 80% 的时段都揪出来,作为容量分析的输入。
数据结构就是一个泛型化的版本:把「连续满足某个条件的元素」聚成一个 List,遇到不满足的就切一刀。
import java.util.ArrayList;
import java.util.List;
import java.util.function.Predicate;
import java.util.stream.Gatherer;
public static <T> Gatherer<T, List<T>, List<T>> runsWhere(Predicate<? super T> condition) {
return Gatherer.ofSequential(
ArrayList::new,
(List<T> run, T element, Gatherer.Downstream<? super List<T>> downstream) -> {
if (condition.test(element)) {
run.add(element);
} else if (!run.isEmpty()) {
downstream.push(List.copyOf(run));
run.clear();
}
return true;
},
(List<T> run, Gatherer.Downstream<? super List<T>> downstream) -> {
if (!run.isEmpty()) {
downstream.push(List.copyOf(run));
}
}
);
}
用法:
List<Integer> samples = List.of(12, 45, 88, 91, 95, 40, 30, 85, 87, 20);
List<List<Integer>> peaks = samples.stream()
.gather(runsWhere(v -> v > 80))
.toList();
// 结果:[[88, 91, 95], [85, 87]]
这里用到了三参数版本的 ofSequential,第三个参数就是 finisher。
finisher 的作用是把「流结束了但状态里还有货」的情况处理掉。
如果不写 finisher,上面例子里的 [85, 87] 就会被丢掉——因为这两个元素后面没有出现不满足条件的元素,触发不了 push 的分支。这个 bug 很隐蔽,测试数据如果恰好以不满足条件的元素结尾,你就发现不了。上线之后遇到以异常值结尾的文件,最后一个区间就消失了,排查起来能让人想砸键盘。
finisher 里必须再判断一次 run.isEmpty()。因为如果流本身是空的,或者最后一个元素刚好触发了 push 并清空了状态,finisher 会被调用但状态是空的。少了这个判断就会多推一个空 List 出去。
另外注意 List.copyOf(run)。不能直接 push 那个 ArrayList 本身,因为接下来 run.clear() 会把它清空。你推过去的引用和状态对象是同一个,下游拿到的 List 会在下一步变成空的。这个坑非常经典,跟当年用 subList 之后改原 List 导致视图跟着变是一个道理。
和短路操作配合的效果。
// 只看第一个峰值区间,后面的采样根本不遍历
List<Integer> firstPeak = samples.stream()
.gather(runsWhere(v -> v > 80))
.findFirst()
.orElse(List.of());
这里就体现出 Gatherer 相比 Collector 的价值了。findFirst 拿到第一个结果之后会通知上游停止推数据,gatherer 的 integrator 自然也就不再被调用。如果你的数据源是 10GB 的日志文件,这个特性可以把处理时间从几分钟压到几秒钟。
不过这里有个细节要注意:短路发生之后,不要指望 finisher 还能正常干活。规范里对短路之后 finisher 是否被调用没有给出明确保证,各实现版本可能有差异。所以关键的清理动作、关键的数据推送,都不要放在 finisher 里指望它一定执行。finisher 只用来处理「流正常结束」这个场景。
五、案例三:并发映射,把 IO 等待时间叠起来
这个是我觉得 JDK 24 这批内置 Gatherer 里最有实用价值的一个。
需求很常见:手上有 500 个 URL,需要一个个请求拿回标题。串行跑一遍要几分钟,因为每个请求都要等网络往返。
以前的写法要么手动搞线程池加 CompletableFuture,要么用 parallelStream。parallelStream 在这种 IO 密集场景下效果很差,因为默认的 ForkJoinPool 大小等于 CPU 核心数,8 核机器就只能并发 8 个请求,而且会阻塞公共池影响其他并行流。
Gatherers.mapConcurrent 就是干这个的:
import java.util.stream.Gatherers;
List<String> urls = loadUrls();
List<Article> articles = urls.stream()
.gather(Gatherers.mapConcurrent(64, this::fetchArticle))
.toList();
第一个参数是最大并发数,第二个参数是单个元素的映射函数。它内部用的是虚拟线程,所以设成 64 甚至 256 都不会像平台线程那样把内存吃爆——每个虚拟线程的栈是堆上的一小块,几百 KB 的量级。
用 64 并发跑 500 个 URL,耗时基本等于最慢那个请求的耗时加上几轮调度,从几分钟降到十几秒很正常。
几个必须知道的点:
输出顺序默认和输入顺序一致。
并发执行,但结果按输入的顺序推下去。这意味着如果你处理的是「第 N 个元素对应第 N 条记录」这种场景,不需要额外做序号对齐,下游拿到的顺序是对的。
代价是它内部要缓存已经完成但还不能推的结果。如果第一个 URL 特别慢,后面 63 个都已经完成了,它们会先攒着,等第一个完成之后才一股脑推出去。极端情况下会导致内存占用升高,但一般请求耗时的方差不会那么大,问题不严重。
不要和 parallel stream 混用。
// 别这么写
urls.parallelStream()
.gather(Gatherers.mapConcurrent(64, this::fetchArticle))
.toList();
mapConcurrent 内部已经在开并发了,外面再套一层并行,会变成嵌套并发,实际并发数远超预期,把你调用的一方打挂是很常见的后果。顺序流上用它就对了。
异常处理要提前想清楚。
如果某个元素的映射函数抛了异常,整个流会中断,异常会向外传播。500 个请求里有 1 个失败就让整批结果失效,通常不是我们想要的。所以映射函数内部要把异常吞掉,返回一个表示失败的结果对象。
record FetchResult(String url, String title, String error) {}
List<FetchResult> results = urls.stream()
.gather(Gatherers.mapConcurrent(64, url -> {
try {
return new FetchResult(url, fetchTitle(url), null);
} catch (Exception e) {
return new FetchResult(url, null, e.getMessage());
}
}))
.toList();
这样下游可以自己决定怎么处理失败的元素,而不是被一个异常全部打断。
六、其他几个内置 Gatherer
JDK 24 的 java.util.stream.Gatherers 工具类里预置了几个常用实现,很多场景直接用它就行,不用自己写。
windowFixed(n),把流切成固定大小的块。适合批量入库,比如每 500 条写一次数据库。
List<List<Order>> batches = orders.stream()
.gather(Gatherers.windowFixed(500))
.toList();
最后一块不足 500 个也照样输出,不会丢掉。
windowSliding(n),滑动窗口,窗口之间相互重叠。
// 每 3 个一组,滑一步
List<List<Double>> windows = prices.stream()
.gather(Gatherers.windowSliding(3))
.toList();
做移动平均、趋势判断的时候很好用。写这个自己实现也不难,但内置的省事。
fold(initial, folder),把整个流聚成一个值,然后往下游推一次。
// 流为空时给出默认值,这是它比 reduce 强的地方
String summary = metrics.stream()
.gather(Gatherers.fold(() -> "无数据", (acc, m) -> acc + m + "; "))
.findFirst()
.orElse("");
reduce 在流为空时会返回 Optional.empty(),你得在外面判断。fold 直接给你一个带默认值的结果,链式写起来更顺。
scan(initial, scanner),累积但每步都推。做前缀和、累计计数用。
List<Integer> prefixSum = List.of(1, 2, 3, 4, 5).stream()
.gather(Gatherers.scan(() -> 0, Integer::sum))
.toList();
// 结果:[1, 3, 6, 10, 15]
七、Greedy 和短路,这个细节影响性能
Gatherer.Integrator 有两种构造方式。
一种是从三段式 lambda 直接推出来的(就是前面案例里用的那样),它是非 greedy 的。每调用一次 integrator,框架都要检查返回值决定下一步。
另一种是用 Integrator.ofGreedy(...) 显式构造,意思是「我这个 integrator 永远返回 true,不用每次检查」。框架拿到这个声明之后,循环里就省掉了返回值判断的分支。windowFixed、windowSliding 这些内置实现用的就是 greedy 版本。
return Gatherer.ofSequential(
DedupState::new,
Gatherer.Integrator.ofGreedy((state, element, downstream) -> {
// 这里永远不返回 false,用 ofGreedy 更合适
...
})
);
但千万别滥用。如果你声明成 greedy,实际上又在某些情况下返回了 false,行为是未定义的。可能在某个 JDK 版本上能跑,换个版本就出诡异的问题。判断标准很简单:这段逻辑有没有可能提前结束?可能就不加,不可能再加。
反过来说,如果确实需要短路,非 greedy 版本才正确。比如实现一个「取满足条件的前 N 个元素」的 gatherer:
public static <T> Gatherer<T, int[], T> takeWhileMatches(Predicate<? super T> condition, int limit) {
return Gatherer.ofSequential(
() -> new int[1],
(int[] count, T element, Gatherer.Downstream<? super T> downstream) -> {
if (!condition.test(element)) {
return true; // 不满足条件,跳过但继续
}
if (count[0] >= limit) {
return false; // 已经够了,通知上游别再送了
}
downstream.push(element);
count[0]++;
return true;
}
);
}
返回 false 之后的元素根本不会被送进 integrator,上游的遍历会尽早停下来。如果你用 ofGreedy 包这段逻辑,框架认为你永远不返回 false,短路就失效了,整个流会被完整遍历一遍,白做。
状态用 int[1] 是个偷懒写法,因为基本类型不能当泛型参数。正规做法还是定义一个小的可变状态类,可读性更好。
八、并行流上的 Gatherer
顺序场景用 Gatherer.ofSequential,它明确不支持并行。如果要在 parallelStream 上用,得用四参数的 Gatherer.of(...) 并提供 combiner。
但这里有个根子上的问题:不是所有 gatherer 都能正确并行化。
并行流的处理方式是把数据分段,每段自己跑一遍完整流程,最后用 combiner 把各个分段的状态合并起来。对于 runsWhere 这种依赖「相邻元素之间关系」的操作,分段之后相邻关系就被切断了——A 段末尾的连续高负载和 B 段开头的连续高负载,本来是同一个区间,分段之后会变成两个。combiner 能合并状态,但没法还原本来应该合并的区间。
所以像连续去重、连续区间提取这类顺序敏感的 gatherer,在并行流上跑出来的结果和顺序流不一致,这不是 bug,是语义决定的。真需要并行,只能换算法或者放弃这些操作。
反过来 mapConcurrent 这类在语义上天然并行的操作就完全没问题,而且它本身就是为此设计的。
我的建议很直接:除非你在实现一个明确可并行的 gatherer,否则就用 ofSequential,别去碰 combiner。写错 combiner 导致结果丢失的问题,比顺序流跑得慢要严重得多,而且极难排查。
九、什么时候不该用 Gatherer
这个 API 很容易让人上头,见到什么需求都想用 Gatherer 写一遍。我列几个反例。
能用 map / filter 解决的,别用 Gatherer。
Gatherer 有状态管理的开销,哪怕是个最简单的转换,也比直接 map 慢。数据量大的时候差别能看出来。如果需求是「把每个元素转成另一种形式」,老老实实用 map。
排序、分组这种有专门操作的,也别硬上。
排序用 sorted,分组用 Collectors.groupingBy。这些是高度优化过的,自己用 gatherer 实现只会更慢更难懂。
逻辑超过二十行就该考虑抽成方法。
Gatherer 的 lambda 里写复杂逻辑,可读性会急剧下降。遇到那种需要判断好几个分支、还要处理各种边界情况的逻辑,直接用 for 循环写反而最清楚。别为了「用新特性」而牺牲可维护性,接手你代码的人不一定熟悉这个 API。
十、踩坑清单
- 状态类用一个静态嵌套类,别用匿名内部类,泛型推断会出问题
- 状态对象的创建必须走 initializer,不能存在 Gatherer 实例的字段或闭包里
- push 的集合要先
List.copyOf()拷一份,否则后续 clear 会把推出去的数据也清掉 - 有尾部数据的逻辑必须写 finisher,否则最后一段会被静默丢弃
- finisher 里要先判断状态是否为空,不然会多推一个空结果
- 短路之后不要再 push,也不要依赖 finisher 一定被执行
- 声明 greedy 就必须永远返回 true,做不到就别用
Integrator.ofGreedy mapConcurrent只用在顺序流上,不要和parallelStream叠加mapConcurrent的映射函数要自己吞异常,否则一个失败会中断整批- 顺序敏感的 gatherer 在并行流上结果会变,这不是 bug 是语义问题
- 项目里的 JDK 版本要确认好,这是 JDK 24 才 final 的 API,21 上的预览版签名可能不一样
最后说一句判断标准。Gatherer 真正的价值在于两个场景:一是有状态的、跨元素判断的逻辑,以前只能写循环;二是需要和短路操作配合的时候,能让惰性真正起作用。如果需求不落在里面,用老的 map、filter、Collector 组合就够了——新特性是用来解决问题的,不是用来刷存在感的。

