如何高效合并已按时间戳排序的Flux数据流并求和?
问题
- 需要合并多个包含时间戳与数值的大型Flux数据流,当多条数据时间戳相同时,对应数值求和
- 所有输入的Flux数据流均已按时间戳升序排列
- 针对小数据流会用
groupBy处理,但面对大量数据时该方法效率不足,希望利用数据流已排序的特性来实现需求,寻求可行的工具或实现方式
预期逻辑伪代码
var flux1 = Flux.just( new Data(ZonedDateTime.parse("2025-01-01T00:00:00"), 1.0), new Data(ZonedDateTime.parse("2025-03-01T00:00:00"), 1.0) ); var flux2 = Flux.just( new Data(ZonedDateTime.parse("2025-02-01T00:00:00"), 2.0), new Data(ZonedDateTime.parse("2025-03-01T00:00:00"), 2.0), new Data(ZonedDateTime.parse("2025-04-01T00:00:00"), 2.0) ); var flux3 = Flux.just( new Data(ZonedDateTime.parse("2025-02-01T00:00:00"), 5.0) ); var input = List.of(flux1, flux2, flux3); var output = Flux.create(sink -> { List<ZonedDateTime> nextEntries = input.stream().map(Flux::next).toList(); do { ZonedDateTime nextTimestamp = nextEntries.stream().map(Data::getTimestamp).min(ZonedDateTime::compareTo).get(); List<Integer> affectedStreams = IntStream.range(0, input.size()).filter(i -> nextTimestamp == nextEntries[i].getTimestamp()).toList(); double nextOutput = affectedStreams.stream().mapToDouble(i -> nextEntries[i].getValue()).sum(); sink.next(new Data(nextTimestamp, nextOutput)); affectedStreams.forEach(i -> nextEntries[i] = input.get(i).next()); } while (!allFluxAreConsumed); });
预期输出
[ Data(ZonedDateTime.parse("2025-01-01T00:00:00"), 1.0), Data(ZonedDateTime.parse("2025-02-01T00:00:00"), 7.0), Data(ZonedDateTime.parse("2025-03-01T00:00:00"), 3.0), Data(ZonedDateTime.parse("2025-04-01T00:00:00"), 2.0) ]
注:原伪代码中预期输出的
2025-05-01应为笔误,已修正为2025-04-01以匹配输入数据
内容的提问来源于stack exchange,提问作者Bernard
相关产品推荐
相关产品推荐

