Java Stream Gatherers 实战:滑动窗口、按 key 去重与会话切分怎么写

2026-09-20 0 933

先讲个具体的活儿。一批用户行为日志,需要把同一用户间隔 30 秒以内的请求归成一个会话,输出每个会话的开始时间、结束时间和访问路径列表。日志按时间排好序了,几百万条,不想全部读进内存。

用 Stream 写,写到一半就卡住了。filtermap 都是无状态的,处理第 8 条记录时需要知道第 7 条记录的时间,这个信息 Stream 里没有地方放。于是退而求其次,用 collect 收集成 List 再循环,或者硬塞一个 AtomicReference 进 lambda 里,代码一下就丑了。

这个缺口最近几年补上了。JDK 22 以预览形式引入了 Gatherer,JDK 24 转正,现在写 stream().gather(...) 不需要再加 --enable-preview。它补的正好是「有状态的中间操作」这块。

为什么 collect 和 mapMulti 都不够用

动手之前先说清楚这东西解决的是什么问题,不然后面容易用错地方。

mapMulti 看起来很像,它可以做一对多的转换,也能提前终止。但它没有跨元素的记忆——处理每个元素时都是干干净净的一双手,拿不到上一个元素留下了什么。

collect 确实能携带状态,问题是它是个终结操作。一旦用了 collect,整条流就被物化成集合了,后面的 limitfindFirst 都失去了短路的意义。上面那个会话切分的需求,如果日志有 800 万条,而我们其实只需要前 100 个会话,collect 的写法会把 800 万条全部读完。

想要的是:保留流的惰性,同时允许中间操作记住点东西。这就是 Gatherer 的位置。

Gatherer 的四个零件

一个 Gatherer 由四部分构成,最常用的其实只有前两个:

  • initializer:给每个流准备一份初始状态,比如 ArrayList::newHashSet::new。不给的话状态就是 null。
  • integrator:核心。每来一个元素调用一次,在里面读状态、改状态、决定要不要往下游推结果。返回 false 表示不再接收元素了。
  • finisher:上游走完之后调用一次,用来处理攒在状态里没来得及输出的尾巴,比如滑动窗口最后一段没凑满的元素。
  • combiner:只在并行流下用得着,负责把两个分片的状态合并起来。不提供就别指望并行。

对应到代码,静态工厂有三个常用的:

Gatherer.ofSequential(integrator)
Gatherer.ofSequential(initializer, integrator)
Gatherer.ofSequential(initializer, integrator, finisher)

Sequential 字样的意味着不参与并行合并。大多数有状态的逻辑本来就该用这一组,后面第五部分会具体说。

从最小的滑动窗口开始

第一个例子做固定大小的窗口:每收满 3 个元素,就打包成一个 List 推出去。

import java.util.ArrayList;
import java.util.List;
import java.util.stream.Gatherer;

public final class Gatherers {

    static <T> Gatherer<T, List<T>, List<T>> windowed(int size) {
        if (size < 1) {
            throw new IllegalArgumentException("size 必须大于 0");
        }
        return Gatherer.ofSequential(
                () -> new ArrayList<T>(size),
                (List<T> buf, T element, Gatherer.Downstream<? super List<T>> down) -> {
                    buf.add(element);
                    if (buf.size() < size) {
                        return true;
                    }
                    List<T> snapshot = List.copyOf(buf);
                    buf.clear();
                    return down.push(snapshot);
                }
        );
    }
}

有几个地方值得停一下。

List.copyOf(buf) 这行不能省。如果直接 down.push(buf) 再把 buf 清空,下游拿到的是同一个对象引用,等它真正消费的时候里面已经空了。这个坑在自己写收集器的时候也踩过,本质是「推出去的对象不能再改」。

方法返回的是 down.push(snapshot),而不是无条件 return true。push 的返回值代表下游还收不收。如果下游已经拒绝(比如前面挂了个 limit(2),第三个窗口就会被拒),这里必须诚实地把 false 传上去,整个流才会停下来。写 return true 不会有编译错误,但会让流多跑很多无用功。

尾巴怎么办?上面这个版本没写 finisher,所以输入 7 个元素、窗口大小为 3 时,最后 1 个元素被静默丢弃。这是有意的设计选择——滑动平均这类场景通常就不要不满的窗口。但如果业务上需要,加上第三个参数就行:

return Gatherer.ofSequential(
        () -> new ArrayList<T>(size),
        (List<T> buf, T element, Gatherer.Downstream<? super List<T>> down) -> {
            buf.add(element);
            if (buf.size() < size) {
                return true;
            }
            List<T> snapshot = List.copyOf(buf);
            buf.clear();
            return down.push(snapshot);
        },
        (List<T> buf, Gatherer.Downstream<? super List<T>> down) -> {
            if (!buf.isEmpty()) {
                down.push(List.copyOf(buf));
            }
        }
);

finisher 的第二个参数是 Downstream,但没有返回值可用——它是个 BiConsumer,push 的结果会被丢掉。如果下游此时已经关闭,这次 push 就无声无息地失败了。绝大多数情况下无所谓,但如果你在 push 前后做了副作用,记得考虑这种可能。

有状态的去重:按业务 key 而不是 equals

Stream.distinct() 只能按 equals 去重。实际业务里更常见的是「同一个订单号只保留第一条」「同一个用户只取最近一次登录」。这时候需要按字段去重,Gatherer 写起来很直接:

static <T, K> Gatherer<T, Set<K>, T> distinctBy(Function<? super T, ? extends K> keyFn) {
    return Gatherer.ofSequential(
            HashSet::new,
            (Set<K> seen, T element, Gatherer.Downstream<? super T> down) -> {
                if (seen.add(keyFn.apply(element))) {
                    return down.push(element);
                }
                return true;
            }
    );
}

用法:

List<Order> unique = orders.stream()
        .gather(distinctBy(Order::orderNo))
        .toList();

这里 seen 是个无限增长的内存结构,数据量大又基数很高的时候要注意。如果场景允许,换成布隆过滤器或者加个容量上限会更稳妥。

另外注意 seen.add(...) 返回 true 时才 push,返回 false 时直接 return true——是「跳过这个元素继续」,不是「停止处理」。这两个 false 的含义完全不同,写反了会导致流在遇到第一个重复项时就整个终止。

完整案例:把日志切成会话

回到开头那个需求。先定义数据结构,用 record 就够了。

import java.time.Duration;
import java.time.Instant;

public record LogEntry(String userId, Instant at, String path) {
}

public record Session(Instant start, Instant end, List<String> paths) {
}

Gatherer 的实现思路是:维护一个缓冲区,每来一条记录就跟缓冲区最后一条比对时间差。超过间隔阈值,就把当前缓冲区封成一个 Session 推出去,然后清空重新开始。

static Gatherer<LogEntry, List<LogEntry>, Session> sessionize(Duration gap) {
    return Gatherer.ofSequential(
            () -> new ArrayList<LogEntry>(),
            (List<LogEntry> buf, LogEntry entry,
             Gatherer.Downstream<? super Session> down) -> {

                if (!buf.isEmpty()) {
                    LogEntry previous = buf.get(buf.size() - 1);
                    Duration delta = Duration.between(previous.at(), entry.at());
                    if (delta.compareTo(gap) > 0) {
                        Session closed = seal(buf);
                        buf.clear();
                        if (!down.push(closed)) {
                            return false;
                        }
                    }
                }
                buf.add(entry);
                return true;
            },
            (List<LogEntry> buf, Gatherer.Downstream<? super Session> down) -> {
                if (!buf.isEmpty()) {
                    down.push(seal(buf));
                }
            }
    );
}

private static Session seal(List<LogEntry> buf) {
    List<String> paths = new ArrayList<>(buf.size());
    for (LogEntry e : buf) {
        paths.add(e.path());
    }
    return new Session(
            buf.get(0).at(),
            buf.get(buf.size() - 1).at(),
            paths
    );
}

主流程:

public static void main(String[] args) {
    List<LogEntry> logs = List.of(
            new LogEntry("u1", Instant.parse("2025-05-06T10:00:00Z"), "/home"),
            new LogEntry("u1", Instant.parse("2025-05-06T10:00:05Z"), "/search"),
            new LogEntry("u1", Instant.parse("2025-05-06T10:00:40Z"), "/detail/7"),
            new LogEntry("u1", Instant.parse("2025-05-06T10:00:45Z"), "/cart"),
            new LogEntry("u1", Instant.parse("2025-05-06T10:02:10Z"), "/pay")
    );

    logs.stream()
            .gather(sessionize(Duration.ofSeconds(30)))
            .forEach(s -> System.out.println(
                    s.start() + " ~ " + s.end() + " " + s.paths()));
}

输出是三段会话:

2025-05-06T10:00:00Z ~ 2025-05-06T10:00:05Z [/home, /search]
2025-05-06T10:00:40Z ~ 2025-05-06T10:00:45Z [/detail/7, /cart]
2025-05-06T10:02:10Z ~ 2025-05-06T10:02:10Z [/pay]

10:00:05 到 10:00:40 之间隔了 35 秒,超过阈值,所以第一段在这里断开;最后一条记录孤零零的,靠 finisher 收尾推出来。

两个使用前提必须说清楚:一是输入必须已经按时间排序,乱序的数据会让结果完全错误,而且不会报错;二是这个 Gatherer 本身不按 userId 分组,多用户混在一起时需要先 filter 或者分组后再处理,否则会把不同用户的时间线串起来。

第一个前提如果满足不了,就在 gather 前面加一个排序:

logs.stream()
        .sorted(Comparator.comparing(LogEntry::at))
        .gather(sessionize(Duration.ofSeconds(30)))
        .forEach(System.out::println);

并行没那么简单

Sequential 的工厂方法不提供 combiner。在并行流上使用这样的 Gatherer,流实现会退化成顺序执行——不会报错,但也不会更快。如果你写了 .parallel() 却发现 CPU 上不去,多半是这个原因。

真想并行,需要用 Gatherer.of(initializer, integrator, combiner),并且自己保证两个分片的状态能合并。拿会话切分来说,麻烦在于:分片 A 的最后一个元素和分片 B 的第一个元素可能属于同一个会话,combiner 必须能把它们接上。这个逻辑写对并不轻松,绝大多数有状态场景的性价比都不高。

所以实践里的建议是:并行处理放在 gather 之前或者之后,中间的 gather 保持顺序。比如先把大文件切成多个块并行解析,再在顺序流里做会话切分。

顺便提一下,Gatherer 本身还有个 andThen,可以像 Function 那样把两个 Gatherer 串起来,写复杂管道时能拆得更清楚。

六个容易写错的地方

一、push 之后继续改状态。推出去的对象必须当成不可变的。要么构造新对象,要么像上面那样 List.copyOf 之后再清空缓冲。

二、忽略 push 的返回值。下游拒收时返回 false 是契约。忽略它不会编译报错,但会让短路的流白跑,数据量大时能差出一个数量级。

三、把「跳过元素」写成了「停止流」。integrator 返回 true 是继续读下一个元素,返回 false 是收工。去重、过滤这类场景要的是前者。

四、状态对象放成静态字段。initializer 的职责就是给每次流执行准备独立状态。写成 static final List 之后,第二条流会看到第一条流留下的残留数据,而且并行下直接出并发问题。

五、把所有收尾逻辑押在 finisher 上。如果上游被 findFirstanyMatch 这类短路操作提前终止,finisher 未必会执行。写数据库、发消息这类副作用别只靠它,该用 try/finally 还是得用。

六、忘了 List.copyOf 会做一次浅拷贝。元素本身还是共享的。如果元素是可变对象,下游改了它,缓冲区里那份也跟着变。对不可变 record 无所谓,对实体类要留个心眼。

该不该用

判断标准可以简化成一句:这个中间操作需不需要记住上一个元素?需要,就考虑 Gatherer;不需要,mapfiltermapMulti 还是更简单更省事。

Gatherer 的力量来自状态,麻烦也来自状态。上面这几段代码都不长,但它们把「流式处理」的边界往外推了一截——以前只能靠物化集合才能表达的逻辑,现在可以保持惰性、可以短路、可以组合。

windowedsessionize 这两个方法收进项目的工具类,后面遇到类似需求就不用再改写成 List 循环了。

Java Stream Gatherers 实战:滑动窗口、按 key 去重与会话切分怎么写
收藏 (0) 打赏

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

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

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

淘吗网 java Java Stream Gatherers 实战:滑动窗口、按 key 去重与会话切分怎么写 https://www.taomawang.com/server/java/2791.html

常见问题

相关文章

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

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