【发布时间】:2019-04-12 18:37:16
【问题描述】:
我运行了这个:
Mono<Void> mono = Mono.empty();
System.out.println("mono.block: " + mono.block());
它会产生:
mono.block: null
正如预期的那样。换句话说,如果Mono 已经完成,调用block 将立即返回。
另一个例子,类似于现实世界的场景。我有一个源通量,例如:
Flux<Integer> ints = Flux.range(0, 2);
我制作了一个可连接的通量,我将使用它来允许多个订阅者:
ConnectableFlux<Integer> publish = ints.publish();
对于这个例子,假设有一个真实的订阅者:
publish
.doOnComplete(() -> System.out.println("publish completed"))
.subscribe();
和另一个只产生元素计数的订阅者:
Mono<Long> countMono = publish
.doOnComplete(() -> System.out.println("countMono completed"))
.count();
countMono.subscribe();
我连接可连接的通量并打印元素计数:
publish.connect();
System.out.println("block");
long count = countMono.block();
System.out.println("count: " + count);
打印出来:
publish completed
countMono completed
block
换句话说,两个订阅者都成功订阅并完成,但随后countMono.block() 无限期阻塞。
为什么会这样,我该如何做?我的最终目标是计算元素的数量。
【问题讨论】:
标签: java project-reactor