用了这么多年 Stream,有些操作一直写得不顺手。最典型的是滑动窗口——算一条曲线的 5 点移动平均,或者判断连续三个采样点都超过阈值。Stream 里没有这个操作。常见的绕过写法有两种:一种是先 collect 成 List,再用 for 循环处理;另一种是拿索引开一个 IntStream.range,在里面 list.subList(i, i + 5)。前者退回了命令式,后者每次 subList 都要开个新视图,最后整个链子看起来像在跟 Stream 较劲。
JDK 24 把这个缺口补上了。Stream.gather(Gatherer) 是一个新的中间操作,允许你自己定义「怎么从上游元素产生下游元素」。这篇就讲清楚 Gatherer 的四个部件,然后用五个真实场景把它用一遍,最后说说几个必须提前知道的坑。
文中 API 以 JDK 24 为准,部分重载在不同版本间有过调整,如果你用的是别的版本,以手边文档为准,思路是一样的。
一、Gatherer 到底由什么组成
先看接口的形状,不用记细节,看个大概:
public interface Gatherer<T, A, R> {
default Supplier<A> initializer() { return () -> null; }
Integrator<A, T, R> integrator();
default BinaryOperator<A> combiner() { return null; }
default BiConsumer<A, Downstream<? super R>> finisher() { return null; }
}
三个泛型参数:T 是上游元素类型,A 是内部状态类型,R 是下游元素类型。四个部件:
initializer:造一个状态对象。如果你这个 Gatherer 不需要跨元素保存东西(比如纯粹做格式转换),返回 null 就行,可以不实现。
integrator:核心。每来一个上游元素就被调一次,签名是 (state, element, downstream) -> boolean。你在这里决定要不要往下游推东西,以及还要不要继续接收元素。
combiner:只在并行流里用到,负责把两个状态合并成一个。返回 null 表示不支持并行。
finisher:上游元素全部处理完之后调一次。用来把缓冲区里剩下的东西吐出去——比如最后那个不满一批的尾巴。
integrator 的返回值是最容易被忽略的地方,它代表「我还想不想继续收到元素」。返回 false 就是从上游断开,整个流会尽快终止。这个语义后面单独用一节讲,因为它是最容易写错、而且写错了不会报错的地方。
还有一个东西叫 Downstream,只有 push(R) 一个方法,返回 boolean。它的返回值的意思是「下游还收不收」,如果下游已经满足了(比如 limit(10) 拿够了,或者 findFirst() 找到了),push 会返回 false,这时候你就该把 integrator 也返回 false,让上游停掉。
二、案例一:滑动窗口算指标
先上最简单的,用内建的 windowSliding。
record Sample(Instant at, double value) {}
List<Double> movingAverage = samples.stream()
.gather(Gatherers.windowSliding(5))
.map(window -> window.stream()
.mapToDouble(Sample::value)
.average()
.orElse(Double.NaN))
.toList();
想法很直白:把流按 5 个一组、每次滑动一格切开,每个窗口算个平均。要判断「连续 3 个点超阈值」,换个 map 就行,比用 for 循环维护一个 boolean 标志清楚得多。
几个行为需要说清楚,不然容易写错:
窗口长度不足时会被丢掉。如果原始数据只有 4 个点,windowSliding(5) 什么都不会发出。所以上面那段代码在数据点很少的时候会得到一个空列表。如果你的业务需要「不够也输出」,得自己写 Gatherer,在 finisher 里处理尾巴。
窗口是重叠的。windowSliding(3) 对 [1,2,3,4,5] 会得到 [1,2,3]、[2,3,4]、[3,4,5]。跟它对应的 windowFixed(3) 是不重叠的,对同样输入得到 [1,2,3]、[4,5],最后那个不满 3 个的窗口会被发出来。这两个的行为差异挺拧的,一个丢尾巴一个留尾巴,用之前最好拿几个小数组试一下。
每个窗口是一个不可变的 List。不能往下游推同一个 List 然后自己清空复用——下游可能还持有它。这点如果你写自定义 Gatherer,是必须注意的。
三、案例二:按大小分批,内建的不够用了
有一批日志行要打包上传,要求是「每包的总字符数不超过 8192,且同一行不能被拆开」。这不是按数量分批,是按累计重量分批,windowFixed 帮不上忙。
先定义状态:
final class LinePack {
final List<String> lines = new ArrayList<>();
int chars = 0;
void reset() {
lines.clear();
chars = 0;
}
List<String> snapshot() {
return List.copyOf(lines);
}
}
然后写 Gatherer:
static Gatherer<String, LinePack, List<String>> packBySize(int maxChars) {
return Gatherer.ofSequential(
LinePack::new,
(pack, line, downstream) -> {
// 这一行加进来会超限,而且包里已经有东西了 —— 先封包
if (pack.chars > 0 && pack.chars + line.length() > maxChars) {
List<String> batch = pack.snapshot();
pack.reset();
if (!downstream.push(batch)) {
return false; // 下游不要了,直接停
}
}
pack.lines.add(line);
pack.chars += line.length();
return true;
},
// finisher:把最后一个不满的包发出去
(pack, downstream) -> {
if (!pack.lines.isEmpty()) {
downstream.push(pack.snapshot());
}
}
);
}
用法:
List<List<String>> packs = Files.lines(path)
.gather(packBySize(8192))
.toList();
这里有个细节值得停下来看一下。pack.chars > 0 这个判断不能省。因为如果单行本身就超过了 8192,不加这个判断就会出现死循环——包里永远是空的,每次加进来都「超限」,封一个空包出去,然后还是超限。加上这个判断,语义就变成了「单行超长时独占一个包」,至少能往下走。
另外注意 snapshot() 用的是 List.copyOf。如果直接把 pack.lines 推给下游然后 clear(),下游拿到的是一个马上被清空的引用。这个 bug 在同步的流里可能碰巧不出问题,等到哪天并行或者下游加了缓冲,就莫名其妙地出现空包。所以状态对象里的集合,往下游推之前一律做不可变拷贝,这是条不用犹豫的规矩。
四、案例三:连续去重,而不是全局去重
distinct() 做的是全局去重,靠的是 hashCode 和 equals。但日志场景要的往往是另一种东西:连续重复的行合并成一条。比如一堆重试日志连着刷了 40 遍,你只想看一遍并标注次数。
final class LastSeen<T> {
T value;
boolean present;
}
static <T> Gatherer<T, LastSeen<T>, T> dedupeConsecutive() {
return Gatherer.ofSequential(
LastSeen::new,
(state, element, downstream) -> {
if (state.present && Objects.equals(state.value, element)) {
return true; // 和上一个相同,吞掉
}
state.value = element;
state.present = true;
return downstream.push(element);
}
);
}
用 Objects.equals 而不是 ==,是因为上游元素可能是 null,也可能是不重写 equals 的对象。这个写法对 null 是安全的:第一个元素是 null 时,present 还是 false,会正常发出去。
如果还想带上重复次数,把状态换成 记录 value + count,在「遇到不同的值」时先发出上一个(带次数),最后在 finisher 里补发最后一个。这个改动很小,但能让输出从「去重后的行」变成「行 + 出现次数」,对排查很有用。
五、案例四:保留原元素的累计值
内建的 Gatherers.scan(initial, f) 是累积的,对每个输入元素发出一个累计结果。看着挺适合算账户流水的累计余额,但真拿它算的时候会发现一个尴尬:它只发累计值,不发原始元素。
List<Long> balances = txns.stream()
.gather(Gatherers.scan(() -> 0L, (sum, t) -> sum + t.deltaFen()))
.toList();
// 结果是 [10, -5, 25, ...] 一堆数字,丢掉了对应哪笔交易
要保留关联,得自己写:
record Running(long balanceFen, Txn txn) {}
static Gatherer<Txn, ?, Running> runningBalance(long openingFen) {
return Gatherer.ofSequential(
() -> new long[]{ openingFen },
(state, txn, downstream) -> {
state[0] += txn.deltaFen();
return downstream.push(new Running(state[0], txn));
}
);
}
用了一个单元素 long[] 当可变容器。写起来最省事,但可读性一般。如果这个状态不止一个字段,老老实实写个小的持有类,或者用 record 的包装,别硬塞进数组里——三个月后你自己都不记得 state[1] 是什么。
这段代码还暴露了一件事:scan 和 fold 这两个内建的内建,一个发 n 个值(每个元素一个累计),一个只发 1 个值(最终结果)。但它们都只关心累计值本身。一旦你要「累计值 + 原始元素」,内建就用不上了。这不是设计缺陷,是接口的边界,知道边界在哪就行。
六、案例五:带并发上限的映射
一批 ID 要调外部接口查详情,接口能给 200ms 左右。原来的写法是用 CompletableFuture 加自定义线程池,还得自己保证输出顺序和原始的 ID 顺序对得上,代码量不小。
List<Detail> details = ids.stream()
.gather(Gatherers.mapConcurrent(64, id -> remoteApi.fetch(id)))
.toList();
一行。几个关键行为:
它底层用的就是虚拟线程。所以这是并发的,而且不需要你去配线程池。64 指的是同时在飞的任务数上限,不是线程数。
输出顺序和输入顺序一致。这是它比「自己丢线程池」省事的地方——不用收集完之后再排一次序。
异常会传播。任意一个映射抛异常,整条流就带着这个异常失败,在飞的其他任务会被取消。这通常是你想要的,但也意味着一个坏 ID 会让整批失败。如果业务允许部分失败,就在映射函数里自己 try-catch 返回一个失败占位对象。
然后是最容易搞错的一点:maxConcurrency 不是下游服务的限流器。它的作用是控制「同一时刻有多少个任务在飞」,目的是避免瞬间开出几万个虚拟线程把内存撑爆。但下游服务能承受的是 QPS,不是并发数。如果你的调用平均耗时 200ms,64 并发就意味着 320 QPS——服务只能扛 50 的话,照样打得它报警。
所以真正的限流还得另外做,比如用信号量按时间窗口控制,或者在远程客户端那一层做。这一点我在改造时踩过一次:把并发从 8 调到 64,压测环境一切正常,预发环境下游容量小,一上去就触发熔断。
七、那个 boolean 返回值,不尊重它会静默变慢
写自定义 Gatherer 的时候,downstream.push() 的返回值和 integrator 的返回值,是最容易含糊过去的地方。它们合起来决定了短路能不能生效。
拿一个常见的需求举例:取元素直到某个条件命中(含命中的那个)。
static <T> Gatherer<T, ?, T> takeThrough(Predicate<T> stop) {
return Gatherer.ofSequential(
(state, element, downstream) -> {
// 先推,再判断要不要停
if (!downstream.push(element)) {
return false;
}
return !stop.test(element);
}
);
}
两个返回值的分工很清楚:push 返回 false 说明下游已经不需要了,那就告诉上游也停;stop 命中说明业务条件达成,同样停下来。
现在看如果偷懒会怎样:
// 别这么写
(state, element, downstream) -> {
downstream.push(element);
return true; // 永远 true
}
功能上「看起来」是对的——下游会自己忽略多余的元素。但上游的循环不会停。假设上游是一百万条记录,下游写的是 .limit(10),这版代码会把一百万个元素全部走一遍,只是后面 99.99% 的 push 都白推了。
更糟的是这种写法不会报错、不会警告,测试也全绿,只是在生产环境上多跑了几百毫秒。所以养成习惯:只要调了 push,就检查它的返回值;只要在循环里判断出了终止条件,就把它体现在 integrator 的返回值上。
顺带提一下 takeWhile。标准库里已经有 Stream.takeWhile 了,标准行为是「不含命中的那个元素」。上面这个 takeThrough 是「含」的版本,两者不冲突。想清楚要哪个再写,别把语义搞混了。
八、并行和 combiner:值不值得写
前面所有例子都用了 Gatherer.ofSequential,意思是这些 Gatherer 不支持并行。如果在一个并行流上使用它们,含这个 Gatherer 的阶段会按顺序求值——也就是说你拿不到并行加速,但结果是对的。
那能不能加个 combiner 让它并行起来?技术上可以,但我想提醒一下代价。
对于滑动窗口这种,状态是最近 N 个元素。两个并行任务的窗口状态要合并,意味着你得处理「A 的尾巴和 B 的开头重叠」这种情况。这个合并逻辑的代码量不小,而且还容易写错——错的地方往往只在某种特定的分块边界上才出现,测试很容易漏。
实际的做法通常是:让 Gatherer 保持串行,但把耗时的部分放在它之前或之后。比如「从数据库读一百万个点算滑动平均」这个任务,并行的价值在读数据,不在算平均。正确的切法是:
// 直接用串行流,滑动窗口这一段也不会有并行问题
List<Double> avg = jdbcTemplate.queryForStream(sql, rowMapper)
.gather(Gatherers.windowSliding(5))
.map(Window::average)
.toList();
数据库的 IO 本身在驱动里已经是并发的。而算平均这一步就是加法,八核并行也就快那么一点,写 combiner 的成本远大于收益。
什么情况下值得写 combiner?我的判断标准是:状态能廉价合并(比如就是个计数器、求和),而且单个元素的处理确实很重(比如要做一次复杂解析)。除此之外,保持串行、用 mapConcurrent 单独把 IO 那段并发掉,是更省心的组合。
九、几个实际会撞上的事
Stream 只能消费一次,Gatherer 也一样。这点在 Stream 上已经人人皆知,但加上一个有状态的 Gatherer 之后更容易出洋相——因为状态对象在流被消费时就绑定了,第二次消费会拿到一个已经跑完的状态。别想着先算一次 toList 再复用同一个 Stream 变量。
调试会变难。整条链全是 lambda,异常堆栈里看到的是一层层 Stream.lambda$...,很难看出是哪一步炸的。我的做法是:任何超过五行的 Gatherer 逻辑,都抽成一个命名方法,而不是写在调用点上。堆栈里出现方法名,排查效率完全不一样。
别为了用而用。如果一个需求用 for 循环加两个局部变量,十行写完而且一眼看得懂,就不要硬塞进 Gatherer。Gatherer 的价值在于「能继续链下去」——滑动窗口之后还想 filter、还想 limit、还想分组,这时候它比 for 循环强。如果处理完就结束了,for 循环没什么不好。
版本和依赖。Gatherers 是 JDK 24 转正的,JDK 22 和 23 上是预览特性,需要 --enable-preview。有些团队的构建工具链还是 JDK 21 LTS,这个特性就用不了,此时退回 collect + 手动处理是唯一选择。升级 JDK 之前先确认一下,别写完才发现 CI 跑不过。
十、结尾
Stream API 从 Java 8 出来到现在十年了,中间加过 takeWhile、dropWhile、ofNullable、iterate 的三参数版本,都是一些小补丁。Gatherer 是这十年里对中间操作模型最大的一次扩展——它把「怎么把上游元素变成下游元素」这个能力开放给了使用者。
回头看那些绕开的写法会发现,它们其实一直在暗示这个缺失:先 collect 再 reassign 成新 stream,本质就是想插一个自定义的中间阶段,只是没办法表达。现在有地方放了。
如果手头正好有一个用了 IntStream.range 加 subList 做滑动窗口的方法,可以拿它当第一个试点。改完之后对比一下两种写法,就知道值不值得在项目里推广了。

