【问题标题】:RxJava: subscribeOn and observeOn not working as expectedRxJava:subscribeOn 和 observeOn 没有按预期工作
【发布时间】:2017-01-22 09:24:36
【问题描述】:

也许我真的很了解subscribeOnobserveOn 的内部工作原理,但我最近遇到了一些非常奇怪的事情。我的印象是,subscribeOn 确定调度程序最初开始处理的位置(尤其是当我们有很多 maps 改变数据流时)然后observeOn 可以在之间的任何地方使用那些maps 在适当的时候更改调度程序(首先进行网络,然后计算,最后更改 UI 线程)。

但是,我注意到当不直接将这些调用链接到我的 Observable 或 Single 时,它​​不会起作用。这是一个最小的工作示例 JUnit 测试:

import org.junit.Test;
import rx.Single;
import rx.schedulers.Schedulers;

public class SubscribeOnTest {

  @Test public void not_working_as_expected() throws Exception {
    Single<Integer> single = Single.<Integer>create(singleSubscriber -> {
      System.out.println("Doing some computation on thread " + Thread.currentThread().getName());
      int i = 1;
      singleSubscriber.onSuccess(i);
    });
    single.subscribeOn(Schedulers.computation()).observeOn(Schedulers.io());

    single.subscribe(integer -> {
      System.out.println("Observing on thread " + Thread.currentThread().getName());
    });
    System.out.println("Doing test on thread " + Thread.currentThread().getName());
    Thread.sleep(1000);
  }

  @Test public void working_as_expected() throws Exception {
    Single<Integer> single = Single.<Integer>create(singleSubscriber -> {
      System.out.println("Doing some computation on thread " + Thread.currentThread().getName());
      int i = 1;
      singleSubscriber.onSuccess(i);
    }).subscribeOn(Schedulers.computation()).observeOn(Schedulers.io());

    single.subscribe(integer -> {
      System.out.println("Observing on thread " + Thread.currentThread().getName());
    });
    System.out.println("Doing test on thread " + Thread.currentThread().getName());
    Thread.sleep(1000);
  }
}

测试not_working_as_expected() 给了我以下输出

Doing some computation on thread main
Observing on thread main
Doing test on thread main

working_as_expected() 给了我

Doing some computation on thread RxComputationScheduler-1
Doing test on thread main
Observing on thread RxIoScheduler-2

唯一的区别是在第一个测试中,在创建单曲之后有一个分号,然后才应用调度程序,并且在工作示例中,方法调用直接链接到单曲的创建。但这不应该无关紧要吗?

【问题讨论】:

  • 这是一个很常见的错误。每个运算符都返回一个新对象,您应该进一步链接该对象。你只需应用subscribeOn+observeOn,忽略返回的Single并订阅原始未更改的源。

标签: java rx-java reactivex


【解决方案1】:

操作员执行的所有“修改”都是不可变的,这意味着它们返回一个新的流,该流以与前一个不同的方式接收通知。由于您只是调用了subscribeOnobserveOn 运算符并且没有存储它们的结果,因此稍后进行的订阅在未更改的流上。

附注:我不太明白您对subscribeOn 行为的定义。如果您的意思是地图操作员会以某种方式受到它的影响,那么这是不正确的。 subscribeOn 定义了一个调度程序,在该调度程序上调用 OnSubscribe 函数。在您的情况下,您传递给 create() 方法的函数。另一方面,observeOn 定义了调度程序,每个连续的流(应用运算符返回的流)在其上处理来自上游的排放。

【讨论】:

    【解决方案2】:

    .subscribeOn(*) - 返回Observable 的新实例,但在第一次测试中你只是忽略它然后订阅原始Observable,显然默认情况下订阅默认主线程。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-08-13
      • 2013-07-07
      • 2012-12-05
      • 2017-05-22
      • 2019-07-28
      • 2014-05-24
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多