You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在RxJava查询中合理组合GroupBy、Buffer与Scan/Reduce

RxJava 3 优化Event聚合对接API的实现方案

问题背景

我的应用有无限Event数据流,Event结构如下:

@Builder
@EqualsAndHashCode
public class Event {
  String Name;
  String Type;
  String Value;
}

需要完成以下对接目标API的操作:

  • 按Name和Type对所有传入Event分组;
  • 将分组后的事件聚合为ApiEvent结构(OrdinalValues和OrdinalCounts为原数据的无损汇总):
@Builder
@ToString
public class ApiEvent {
    String Name;
    String Type;
    Collection<String> OrdinalValues;
    Collection<Integer> OrdinalCounts;
}
  • 最多收集1000条聚合后的事件,或等待最长60秒;
  • 调用API传入ApiEvent集合。

原实现代码(基于RxJava 3):

PublishProcessor<Event> publishProcessor = PublishProcessor.create();
TestSubscriber<List<ApiEvent>> testSubscriber = new TestSubscriber<>();
TestScheduler testScheduler = new TestScheduler();

publishProcessor
        .groupBy(event -> Event.builder().Name(event.Name).Type(event.Type).build())
        .flatMap(x -> x.scan(new Hashtable<>(), this::reduceEvents)
                .map(a -> (ApiEvent.builder()
                                .Name(x.getKey().Name)
                                .Type(x.getKey().Type)
                                .OrdinalValues(a.keySet())
                                .OrdinalCounts(a.values())
                                .build()
                        )
                )
        )
        .buffer(60, TimeUnit.SECONDS, testScheduler, 1000)
        .filter(l -> !l.isEmpty())
        .flatMap((List<ApiEvent> x) -> {

            val grp = x.stream()
                    .collect(groupingBy(post -> Pair.of(post.Name, post.Type)));

            val reducedData = new ArrayList<ApiEvent>(grp.size());

            grp.keySet().forEach(key -> grp.get(key).stream()
                    .max(Comparator.comparingInt(o -> o.OrdinalCounts.size()))
                    .map(reducedData::add));

            return Flowable.fromArray(reducedData);
        })
        .subscribe(testSubscriber);

publishProcessor.onNext(Event.builder().Name("Name1").Type("Type1").Value("Value1").build());
publishProcessor.onNext(Event.builder().Name("Name1").Type("Type1").Value("Value1").build());
publishProcessor.onNext(Event.builder().Name("Name2").Type("Type2").Value("Value2").build());
publishProcessor.onNext(Event.builder().Name("Name2").Type("Type2").Value("Value2").build());

testScheduler.advanceTimeBy(60, TimeUnit.SECONDS);

testSubscriber.assertValueCount(1);

testSubscriber.values().forEach(a -> {
    System.out.println(a);
    // Result: [
    //    ApiEvent(Name=Name2, Type=Type2, OrdinalValues=[Value2], OrdinalCounts=[2]),
    //    ApiEvent(Name=Name1, Type=Type1, OrdinalValues=[Value1], OrdinalCounts=[2])
    // ]
});

对应的reduce函数:

private Hashtable<String, Integer> reduceEvents(Hashtable<String, Integer> hashtable,
                                                    Event event) {
    if (hashtable.containsKey(event.Value)) {
        hashtable.replace(event.Value, hashtable.get(event.Value) + 1);
    } else {
        hashtable.put(event.Value, 1);
    }
    return hashtable;
}

原实现能得到预期结果,但因使用scan导致每个分组会输出多次中间聚合结果,后续需要二次处理去重,存在性能浪费,寻求更优实现方案。


优化方案

核心思路是让每个groupBy后的分组只输出最终聚合结果,而非每次事件到来都输出中间状态,这样后续buffer收集的就是已经聚合完成的ApiEvent,无需二次处理。

优化点说明

  1. 对每个分组使用窗口内聚合而非scan:仅在buffer窗口结束时输出当前分组的最终聚合结果,避免中间冗余数据;
  2. 简化分组Key:用Pair替代临时构建的Event对象,提升效率;
  3. 用merge方法简化计数逻辑,代码更简洁。

优化后代码(窗口内聚合,不跨窗口累积)

这种方案适合窗口结束后无需保留分组计数的场景,是最贴合原需求的高效实现:

PublishProcessor<Event> publishProcessor = PublishProcessor.create();
TestSubscriber<List<ApiEvent>> testSubscriber = new TestSubscriber<>();
TestScheduler testScheduler = new TestScheduler();

publishProcessor
        // 按Name和Type分组,用Pair作为分组Key更高效
        .groupBy(event -> Pair.of(event.Name, event.Type))
        // 对每个分组,在buffer窗口周期内做聚合,窗口结束时输出一次结果
        .flatMap(group -> 
            group.buffer(60, TimeUnit.SECONDS, testScheduler)
                .filter(buffer -> !buffer.isEmpty())
                .map(events -> {
                    // 聚合当前窗口内的同组事件
                    Hashtable<String, Integer> valueCountMap = new Hashtable<>();
                    for (Event event : events) {
                        valueCountMap.merge(event.Value, 1, Integer::sum);
                    }
                    // 转换为ApiEvent
                    Pair<String, String> key = group.getKey();
                    return ApiEvent.builder()
                            .Name(key.getLeft())
                            .Type(key.getRight())
                            .OrdinalValues(valueCountMap.keySet())
                            .OrdinalCounts(valueCountMap.values())
                            .build();
                })
        )
        // 收集聚合后的ApiEvent,达到1000条或60秒触发一次API调用
        .buffer(60, TimeUnit.SECONDS, testScheduler, 1000)
        .filter(list -> !list.isEmpty())
        .subscribe(testSubscriber);

// 测试事件发送
publishProcessor.onNext(Event.builder().Name("Name1").Type("Type1").Value("Value1").build());
publishProcessor.onNext(Event.builder().Name("Name1").Type("Type1").Value("Value1").build());
publishProcessor.onNext(Event.builder().Name("Name2").Type("Type2").Value("Value2").build());
publishProcessor.onNext(Event.builder().Name("Name2").Type("Type2").Value("Value2").build());

// 推进时间触发窗口
testScheduler.advanceTimeBy(60, TimeUnit.SECONDS);

// 验证结果
testSubscriber.assertValueCount(1);
testSubscriber.values().forEach(System.out::println);

可选方案(跨窗口累积计数)

如果需要跨窗口保留分组的计数状态(即窗口结束后不重置计数),可以结合scan和window操作符,只在窗口结束时取分组的最新聚合状态:

publishProcessor
        .groupBy(event -> Pair.of(event.Name, event.Type))
        .flatMap(group -> 
            group.scan(new Hashtable<>(), this::reduceEvents)
                // 只保留每个分组的最新聚合状态
                .last()
                // 每个窗口周期触发一次输出
                .repeatWhen(completed -> completed.delay(60, TimeUnit.SECONDS, testScheduler))
        )
        // 收集到1000条后触发API调用
        .buffer(1000)
        .subscribe(testSubscriber);

内容的提问来源于stack exchange,提问作者user1574775

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 18:07:13