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

如何高效合并已按时间戳排序的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 19:50:15