【问题标题】:Subscribe / Iterate a List of Mono订阅/迭代 Mono 列表
【发布时间】:2020-07-23 16:15:31
【问题描述】:

我有一个返回 Mono<Result> 的规则列表,我需要执行它们并返回另一个包含每个函数结果的 Mono 列表。

Mono<List< Rule>> 到 Mono<List< RuleResult>>

我有这个,它有效,但阻止执行对我来说似乎不正确:

List<Rule> RuleSet
...
Mono<List<RuleResult>> result= Mono.just(RuleSet.stream().map(rule -> rule.assess(object).block()).collect(Collectors.toList()));

如何将其转换为更少阻塞的版本?

我尝试了以下方法:

//Create a List of Monos
List<Mono<RuleResult>> ruleStream=statefulRuleSet.stream().map(rule -> rule.assess(assessmentObject)).collect(Collectors.toList());

//Cannot convert type
Flux.zip(ruleStream,...)).collectList();

//Not sure how to do this
Flux.fromIterable(ruleStream...).collectList();

也许我在考虑一个错误的解决方案,有人有任何指示吗?

【问题讨论】:

  • “并返回另一个包含每个函数结果的 Mono 列表” - 从代码中看起来你的意思是“并返回一个包含每个函数结果列表的 Mono”。

标签: java reactive-programming project-reactor


【解决方案1】:

例如:

interface Rule extends Function<Object, Mono<Object>> { }

public static void main(String[] args) {

  Rule rule1 = (o) -> Mono.just(Integer.valueOf(o.hashCode()));
  Rule rule2 = (o) -> Mono.just(o.toString());
  List<Rule> rules = List.of(rule1, rule2);

  Object object = new Object();

  Mono<List<Object>> result = Flux.fromIterable(rules)
    .flatMapSequential(rule -> rule.apply(object)).collectList();

  result.subscribe(System.out::println);
}

使用flatMapSequential 允许您同时等待最多maxConcurrency 个结果。 maxConcurrency 值可以指定为flatMapSequential 的附加参数。在 reactor-core 3.3.8 中,其默认值为 256。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-04-22
    • 2016-06-17
    • 1970-01-01
    • 2020-05-18
    • 2023-03-04
    • 2022-10-01
    • 2020-11-28
    相关资源
    最近更新 更多