【问题标题】:How to perform a moving window calculation on a Flux and output the result as a new Flux如何对 Flux 执行移动窗口计算并将结果输出为新的 Flux
【发布时间】:2019-10-19 06:03:42
【问题描述】:

我想对 Flux 执行移动窗口计算并生成包含计算值的 Flux,但我不知道如何完成此操作。

作为一个简化的例子,假设我有一个整数通量,我想用这个通量中每 3 个连续整数的总和生成一个新通量。举例说明:

第一个通量包含从 1 到 8 的整数:{1, 2, 3, 4, 5, 6, 7, 8}

结果通量应包含总和:{1+2+3, 2+3+4, 3+4+5, 4+5+6, 5+6+7, 6+7+8}

我可以很容易地产生第一个通量并推导出包含连续 3 个值的通量,如下所示:

Flux<Integer> f1 = Flux.range(1,8);
Flux<Flux<Integer>> f2 = f1.window(3,1);

我也可以 subscribe() 到 f2 并计算总和,但我不知道如何同时将这些总和作为新的 Flux 发布。

是我遗漏了一些简单的事情,还是这种事情真的很难做到?

【问题讨论】:

    标签: java project-reactor


    【解决方案1】:

    您可以在内部通量上使用.reduce(Integer::sum) 对窗口中的元素进行求和,在外部通量上使用.flatMap 将这些和合并回单个流中。

    请注意,由于 .window 是用 maxSize &lt; skip 调用的,因此尾随窗口中的项目将小于最大大小。

    Flux<Integer> sums = Flux.range(1, 8)                    // Flux<Integer>
            .window(3, 1)                                    // Flux<Flux<Integer>>
            .flatMap(window -> window.reduce(Integer::sum)); // Flux<Integer>
    
    StepVerifier.create(sums)
            .expectNext(6)  // 1+2+3
            .expectNext(9)  // 2+3+4
            .expectNext(12) // 3+4+5
            .expectNext(15) // 4+5+6
            .expectNext(18) // 5+6+7
            .expectNext(21) // 6+7+8
            .expectNext(15) // 7+8
            .expectNext(8)  // 8
            .verifyComplete();
    

    【讨论】:

    • 谢谢。至少我认为它很简单是对的。我将为其他人添加在我的实际应用程序中我将使用.flatMapSequential(...) 因为我需要保留原始元素的顺序。
    猜你喜欢
    • 2021-07-30
    • 1970-01-01
    • 1970-01-01
    • 2023-03-10
    • 2017-12-27
    • 1970-01-01
    • 2014-11-16
    • 1970-01-01
    • 2021-04-15
    相关资源
    最近更新 更多