【问题标题】:Mono.count block indefinitelyMono.count 无限期阻塞
【发布时间】: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


    【解决方案1】:

    您可以使用autoConnectrefCount 来实现此功能,而不是手动调用connect()

    例如:

            Flux<Integer> ints = Flux.range(0, 2);
            Flux<Integer> publish = ints.publish()
                    .autoConnect(2);  // new 
            publish
                    .doOnComplete(() -> System.out.println("publish completed"))
                    .subscribe();
            Mono<Long> countMono = publish
                    .doOnComplete(() -> System.out.println("countMono completed"))
                    .count();
            // countMono.subscribe();
            long count = countMono.block();
            System.out.println("count: " + count);
    

    为什么你的例子不起作用?

    这是我认为在您的示例中发生的情况...但这是基于我有限的知识,我不能 100% 确定它是正确的。

    1. .publish() 将上游源变成热流
    2. 然后您订阅了两次(但这些还没有启动流程,因为可连接的通量尚未连接到上游)
    3. .connect() 订阅上游,并启动流
    4. 上游,以及在connect() 之前注册的两个订阅完成(因为这一切都发生在主线程中)
    5. 此时,ConnectableFlux 不再连接到上游,因为上游已完成(反应器文档对在上游源完成后新订阅到达时 ConnectableFlux 发生的情况的详细信息,所以这就是我我不是 100% 确定的。)
    6. block() 创建一个新的订阅。
    7. 但由于 ConnectableFlux 不再连接,因此没有数据流动
    8. 如果您再次调用connect()(从另一个线程,因为主线程被阻塞),数据将再次流动,block() 将完成。但是,这将是一个序列(不是在步骤 4 中完成的原始序列)

    为什么我的示例有效?

    只创建了两个订阅(而不是您的示例中的 3 个),一个来自 .subscribe() 调用,另一个来自 .block()。 ConnectableFlux 在 2 次订阅后自动连接,因此 block() 订阅完成。两个订阅共享相同的上游序列。

    【讨论】:

    • 谢谢,这是有道理的。如果Mono 在链条下游,即不是ConnectableFlux 的直接订阅者,我该怎么做?或者换一种说法——有没有办法以阻塞的方式从Mono 获取最新结果?
    • 我不太了解您想要什么,但是...要获取以前订阅发出的值,您可以使用 .cache (以便将来的订阅者在一段时间内看到相同的值时间)。或者您可以将发出的值保存在流之外的某个位置(来自doOnNext),然后从保存的位置查找值。
    • 上面我有Mono&lt;Long&gt; countMono = publish.count()。假设我有Flux&lt;Integer&gt; publish2 = publish.flatMap(e -&gt; whatever),然后是Mono&lt;Long&gt; countMono = publish2.count()。 IE。在publishcountMono 之间有一个中间流(或多个)。以上将使publish2自动订阅,但它可能会在countMono.block()被调用之前完成,因此会遇到同样的问题,对吧? cache 会无限期地缓存,对吗?对于长期运行的作业来说可能是个问题。
    猜你喜欢
    • 2018-03-13
    • 1970-01-01
    • 2012-05-14
    • 1970-01-01
    • 2011-02-25
    • 1970-01-01
    • 2018-03-27
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多