如何在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,无需二次处理。
优化点说明
- 对每个分组使用窗口内聚合而非
scan:仅在buffer窗口结束时输出当前分组的最终聚合结果,避免中间冗余数据; - 简化分组Key:用
Pair替代临时构建的Event对象,提升效率; - 用
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
相关产品推荐
相关产品推荐

