【问题标题】:Reactive programming (Reactor) : Why main thread is stuck?反应式编程(Reactor):为什么主线程卡住了?
【发布时间】:2020-09-28 06:43:06
【问题描述】:

我正在学习使用 project-reactor 进行反应式编程。

我有以下测试用例:

@Test
public void createAFlux_just() {
    Flux<String> fruitFlux = Flux.just("apple", "orange");
    fruitFlux.subscribe(f -> {
        try {
            Thread.sleep(5000);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        System.out.println(f);
    });
    System.out.println("hello main thread");
}

通过执行测试,似乎主线程卡住了 5 秒。

我希望订阅的消费者应该在自己的线程中异步运行,也就是说,订阅调用应该立即返回到主线程,因此hello main thread 应该立即打印。

【问题讨论】:

    标签: java reactive-programming project-reactor


    【解决方案1】:

    主线程卡住了,因为订阅发生在main 线程上。如果您希望它异步运行,您需要在main 以外的线程上进行订阅。你可以这样做:

     @Test
    public void createAFlux_just() {
        Flux<String> fruitFlux = Flux.just("apple", "orange");
        fruitFlux.subscribeOn(Schedulers.parallel()).subscribe(f -> {
            try {
                Thread.sleep(5000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            System.out.println(f);
        });
        System.out.println("hello main thread");
    }
    

    注意:我使用了parallel 线程池。你可以使用任何你喜欢的游泳池。 Reactor 的管道默认在调用线程上执行(不像CompletableFuture&lt;T&gt; 默认在ForkJoin 池中运行)。

    【讨论】:

      【解决方案2】:

      如果您有一个异步的 observable (Flux),就会出现这种情况。您通过 just 方法选择使用具有两个现成可用值的 Flux。因为它们立即可用,所以它们被立即传递给订阅对象。

      【讨论】:

        【解决方案3】:

        来自 spring.io documentation

        线程模型 Reactor 操作符通常是并发不可知论者:它们不会强加特定的线程模型,而只是在调用其 onNext 方法的 Thread 上运行。

        调度程序抽象 在 Reactor 中,调度程序是一种抽象,它使用户可以控制线程。调度器可以产生 Worker,它们在概念上是线程,但不一定由线程支持(我们稍后会看到一个例子)。 Scheduler 还包含时钟的概念,而 Worker 则纯粹用于调度任务。

        所以你应该通过subscribeOn 方法订阅不同的线程,Thread.sleep(5000) 将休眠调度程序的线程。您可以在文档中看到更多类似的示例。

        Flux.just("hello")
            .doOnNext(v -> System.out.println("just " + Thread.currentThread().getName()))
            .publishOn(Scheduler.boundedElastic())
            .doOnNext(v -> System.out.println("publish " + Thread.currentThread().getName()))
            .delayElements(Duration.ofMillis(500))
            .subscribeOn(Schedulers.elastic())
            .subscribe(v -> System.out.println(v + " delayed " + Thread.currentThread().getName()));
        

        【讨论】:

          猜你喜欢
          • 2020-11-09
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2016-11-01
          • 2013-03-23
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多