【问题标题】:How to eagerly merge two Flux?如何急切地合并两个 Flux?
【发布时间】:2019-03-03 10:09:03
【问题描述】:
Flux<Long> flux1 = Flux
        .<Long>create(fluxSink -> {
            for (long i = 0; i < 20; i++) {
                fluxSink.next(i);
            }
        })
        .filter(aLong -> aLong % 2 == 0)
        .doOnNext(aLong -> System.out.println("flux 1 : " + aLong));

Flux<Long> flux2 = Flux
        .<Long>create(fluxSink -> {
            for (long i = 0; i < 20; i++) {
                fluxSink.next(i);
            }
        })
        .filter(aLong -> aLong % 2 == 1)
        .doOnNext(aLong -> System.out.println("flux 2 : " + aLong));

Flux.merge(flux1, flux2)
        .doOnNext(System.out::println)
        .then()
        .block();

像上面的代码一样创建两个Flux&lt;Long&gt;。

flux1 创建偶数流 (0,2,4,6,8 ...) Flux2 创建奇数流 (1,3,5,7,9 ...)

我预计当合并这 2 个flux1 和flux2 时会像这样工作

0,1,2,3,4 ... 或 0,2,1,3,4.. 取决于计算能力

但总是花费flux1和flux2 (flux1 start)0,2,4,6,8, ... 16,18,(flux1 end)(flux2 start)1,3,5,7 ... 17,19

如何订阅多个flux eager事件?

【问题讨论】:

  • 你想做什么?能不能解释的更清楚一点?

标签: project-reactor reactive-streams


【解决方案1】:

两个流在同一个线程上运行。当您订阅 flux1 时开始推送数据,直到完成。只有这样,flux2 才能继续该线程。 merge 运算符按照它们到达的顺序发出值。它不会在第一个和第二个流之间切换。

如果您希望流同时运行,您需要在不同的线程上运行它们,例如通过使用publishOn 运算符。

Flux<Long> flux1 = Flux
    .<Long>create(fluxSink -> {
        for (long i = 0; i < 20; i++) {
            fluxSink.next(i);
        }
    })
    .publishOn(Schedulers.newSingle("thread-x")
    .filter(aLong -> aLong % 2 == 0)
    .doOnNext(aLong -> System.out.println("flux 1 : " + aLong));

【讨论】:

    猜你喜欢
    • 2011-06-24
    • 1970-01-01
    • 2015-10-22
    • 2011-12-16
    • 1970-01-01
    • 2013-07-28
    • 1970-01-01
    • 1970-01-01
    • 2021-07-31
    相关资源
    最近更新 更多