【问题标题】:How can I perform flatMap using multiple threads in Reactor?如何在 Reactor 中使用多个线程执行 flatMap?
【发布时间】:2018-11-21 10:04:03
【问题描述】:

我尝试在Flux range 上运行flatMap,然后是subscribeOn,似乎所有操作都在同一个线程上运行。这正常吗?

Flux.range(0, 1000000).log().flatMap{ it + 1 }.subscribeOn(Schedulers.parallel()).subscribe()

【问题讨论】:

    标签: multithreading kotlin project-reactor


    【解决方案1】:

    您可以按如下方式创建ParallelFlux

    Flux.range(0, 100000).parallel(2).runOn(Schedulers.parallel()).log().map{ it + 1 }.subscribe()
                          ^^^^^^^^^^^  ^^^^^^use runOn ^^^^^^^^^^^
    

    【讨论】:

    • 所以 flatMap 不会将您的流分成并行处理的子流?这就是该方法的记录方式,并且认为只有 subscribeOn 是必要的。
    • flatmap 只有在块内创建新的 Observable(Flux, Mono) 时才会并行化。在您使用it +1 的示例中,运算符flatMap 将无法编译,您应该使用map。我刚刚更新了我的答案。见stackoverflow.com/questions/43269275/…
    【解决方案2】:

    使用.subscribeOn(Schedulers.X) 管理订阅

    使用.publishOn(Schedulers.X) 管理发布

    使用 ParallelFlux 时使用 .parallel(N).runOn(Schedulers.X)

    【讨论】:

      猜你喜欢
      • 2016-02-18
      • 1970-01-01
      • 2015-08-03
      • 1970-01-01
      • 1970-01-01
      • 2019-06-30
      • 2019-12-18
      • 2018-07-07
      • 1970-01-01
      相关资源
      最近更新 更多