Stream Gatherers 实战:滑动窗口、按大小分批与并发映射的自定义中间操作

2026-09-24 0 486

用了这么多年 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] 是什么。

这段代码还暴露了一件事:scanfold 这两个内建的内建,一个发 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 出来到现在十年了,中间加过 takeWhiledropWhileofNullableiterate 的三参数版本,都是一些小补丁。Gatherer 是这十年里对中间操作模型最大的一次扩展——它把「怎么把上游元素变成下游元素」这个能力开放给了使用者。

回头看那些绕开的写法会发现,它们其实一直在暗示这个缺失:先 collect 再 reassign 成新 stream,本质就是想插一个自定义的中间阶段,只是没办法表达。现在有地方放了。

如果手头正好有一个用了 IntStream.rangesubList 做滑动窗口的方法,可以拿它当第一个试点。改完之后对比一下两种写法,就知道值不值得在项目里推广了。

Stream Gatherers 实战:滑动窗口、按大小分批与并发映射的自定义中间操作
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请发送邮件1506151422@qq.com联系,将会及时下架删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 java Stream Gatherers 实战:滑动窗口、按大小分批与并发映射的自定义中间操作 https://www.taomawang.com/server/java/2808.html

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务