Java 24 Stream Gatherers 实战:连续去重、区间提取与并发映射的完整写法

2026-09-19 0 288

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,不用每次检查」。框架拿到这个声明之后,循环里就省掉了返回值判断的分支。windowFixedwindowSliding 这些内置实现用的就是 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 组合就够了——新特性是用来解决问题的,不是用来刷存在感的。

Java 24 Stream Gatherers 实战:连续去重、区间提取与并发映射的完整写法
收藏 (0) 打赏

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

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

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

淘吗网 java Java 24 Stream Gatherers 实战:连续去重、区间提取与并发映射的完整写法 https://www.taomawang.com/server/java/2787.html

常见问题

相关文章

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

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