【问题标题】:How to execute list of Mono sequentialy如何按顺序执行 Mono 列表
【发布时间】:2022-01-06 13:50:16
【问题描述】:

我有一个我想处理的单声道列表,但它们必须按顺序执行,并且只有在前一个单声道完成后才应该执行下一个。

private Mono<List<Result>> processGoals(List<> goals,Data data) {
    
    List<Mono<Result>> plans = goals
                              .stream()
                              .map(plan -> processGoal(plan, data))
                              .collect(Collectors.toList());

}

我尝试使用

return Flux.concat(plans).subscribeOn(Schedulers.single()).collectList();

但这会在前一个单声道完成之前执行下一个单声道。

【问题讨论】:

    标签: java spring reactive-programming spring-webflux project-reactor


    【解决方案1】:
    .flatMapSequential(goal -> processGoal(goal, data), 1)
    

    最后一个参数并发1很重要。我之前尝试过没有并发参数,但没有奏效。感谢@Michael McFadyen

    【讨论】:

    • 为什么不concatMap
    • 我相信 concatMap 也可能是另一种解决方案。
    【解决方案2】:

    Flux#concatMap 是这种情况下的最佳选择。

    它将按顺序合并每个映射的发布者并一次触发一个,而无需显式定义concurrency 参数。

    这是一个完整的例子:

    Flux.fromIterable(goals))
            .concatMap(goal -> processGoal(goal, data))
            .collectList();
    

    【讨论】:

      【解决方案3】:

      而不是混合两种范式:Streams 和 Reactive Streams。您可以尝试完全反应。

      尝试以下操作:

      Mono<List<Result>> res = Flux.just(goals.toArray(Goal[]::new))
                                  .flatMapSequential(goal -> processGoal(goal, data))
                                  .subscribeOn(Schedulers.single())
                                  .collectList();
      

      【讨论】:

      • 您可能希望使用flatMapSequential 来满足“它们必须按顺序执行”的要求。您还可以使用concurrency 参数来避免切换Schedulers 即。 .flatMapSequential(goal -&gt; processGoal(goal, data), 1)
      • 你说得对,我错过了。
      • 我确实尝试过 flatmapsequential 但这里的关键是需要设置为 1 的并发参数
      • 我试图弄清上下文,以便我可以解释它是单声道列表,并且需要一次一个顺序地执行。
      • 为什么不concatMap
      猜你喜欢
      • 2019-11-04
      • 2013-12-04
      • 2021-09-04
      • 2018-09-03
      • 2011-07-29
      • 2021-03-16
      • 2011-03-02
      • 2017-09-06
      • 2021-07-14
      相关资源
      最近更新 更多