Skip to content

7. 库 RxJava

RxJava 库基于以下概念:一个类型为 T 的 Observable<T> 元素流由一个或多个订阅者(Subscriber<T>)进行观察。 RxJava 库允许 Observable<T> 流在 T1 线程中运行,其观察者 Subscriber<T> 在 T2 线程中运行,而开发人员无需无需担心管理这些线程的生命周期,也无需处理诸如线程间数据共享和为执行全局任务而进行的线程同步等自然而然的难题。因此,它简化了异步编程。

一个 Observable<T> 流会生成类型为 T 的元素,这些元素在生成时即可被观察。 如果观察者和可观察对象(此处为方便起见,将 Observable<T> 称为“可观察对象”)位于同一线程中,那么可观察对象只能在观察者消耗完元素 i 之后,才可生成元素 (i+1)。这种架构具有实际意义的情况很少。 如果观察者和可观察对象不在同一个线程中,那么可观察对象及其观察者将具有独立的行为:可观察对象按自己的节奏生成元素,观察者按自己的节奏消费元素。这正是该库的价值所在。到目前为止,我们一直只讨论了一个观察者。实际上,一个可观察对象可以拥有任意数量的观察者。

RxJava 库特别适合用于第 2 节中介绍的架构,该架构在此重述如下:

Image

  • 在 [1] 中,一个服务层提供各类服务,其中部分服务需要较长时间才能获取(例如网络请求);
  • 该服务层由图形用户界面 [1](Swing、Android、JavaFx)调用。 如果服务层与调用它的方法 [swing] 在同一线程中运行,则在等待服务结果期间,图形用户界面会处于冻结状态(无响应);
  • 在 [2] 中,通过 RxJava 实现的轻量级适配层,可向图形层提供该服务的异步实现:该服务可在与调用它的图形层方法不同的线程中执行。 在此情况下,图形界面 [3] 仍保持响应性:用户可以继续与其交互,例如在第一个请求并行时触发新的网络请求,最重要的是,可以为用户提供取消耗时过长处理的选项,而如果图形界面处于冻结状态,这是无法实现的;
  • 调用 [4] 是同步的,而调用 [5-6] 则是异步的;

在此架构中,[2] 层提供返回 Observable<T> 类型的服务,图形层 [3] 的方法可以订阅这些服务。 随后,[2] 层的一个服务会逐一返回其结果,而 [3] 层可以对每个结果做出响应,例如更新图形界面的一个或多个组件。

Observable<T> 类拥有数十种方法。这也是该库的一大难点:它功能非常丰富,很难掌握其所有可能性。我们将介绍其中的一部分。其余方法的掌握将随着时间的推移而逐渐掌握。

7.1. 创建可观察对象并订阅

7.1.1. 示例-01:[Observable.from]方法

  

请看以下代码:


package dvp.rxjava.observables;

import rx.Observable;
import rx.functions.Action0;
import rx.functions.Action1;

import java.util.Arrays;

public class Exemple01 {
  public static void main(String[] args) {
    // 整数可观测量
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    obs1.subscribe(new Action1<Integer>() {
      @Override
      public void call(Integer integer) {
        System.out.printf("next : %s%n", integer);
      }
    }, new Action1<Throwable>() {
      @Override
      public void call(Throwable throwable) {
        System.out.println(throwable);
      }
    }, new Action0() {
      @Override
      public void call() {
        System.out.println("completed");
      }
    });
  }
}
  • 第 12 行:基于一个整数列表创建了一个 Observable<Integer> 类型。

Observable<T> 是一个可观察的 T 类型元素流,通常(但并非必须)以异步方式在元素生成时进行观察。其定义如下:

 

如前所述,Observable<T> 类拥有数十个方法。 其中部分方法与第 5 节中探讨的 Stream<T> 类中的方法相似。RxJava 的文档中包含 [2] 的“大理石图”,这些图示说明了这些方法的工作原理:

  • 第 3 行展示了可观测量随时间的变化;
  • 方法 [4] 应用于该观测量发出的元素。它通常会产生一个新的观测量;
  • 第 5 行展示了所得的新观测量;

方法 [Observable.from] 的签名如下:

 

静态方法 [Observable.from] 允许根据一组类型为 T 的元素创建一个 Observable<T>。这是入门可观察对象的一种非常简单的方式。该行:


    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));

将发出三个元素。它不会立即发出这些元素,而是在每次有观察者声明订阅时才完整地发出它们。这被称为冷可观察对象。该可观察对象会为每个新订阅者重新发出其元素。

我们可以将上述语句视为对可观察对象的配置操作。该对象只需配置一次,如果出现 n 个订阅者,则会执行 n 次。

如何订阅?

一种方法是使用 [Observable.subscribe] 方法,此处的定义如下:

 
  • 该方法的第一个参数 [Action1<T> onNext](参见第 6.2 节)是当可观察对象发出新元素 T 时要执行的方法;
  • 该方法的第二个参数 [Action1<Throwable> onError] 是在可观察对象抛出异常时需执行的方法;
  • 该方法的第三个参数 [Action0 onComplete](参见第 6.1 节)是在可观察对象抛出异常时要执行的方法;
  • 该方法返回类型为 [Subscription];

类型 [Subscription] 表示对可观察对象的订阅。其定义如下:

 

该 [1] 接口的意义在于其 [2] 方法,该方法可用于取消订阅。

在本例中,可观察对象订阅的代码如下:


    obs1.subscribe(new Action1<Integer>() {
      @Override
      public void call(Integer integer) {
        System.out.printf("next : %s%n", integer);
      }
    }, new Action1<Throwable>() {
      @Override
      public void call(Throwable throwable) {
        System.out.println(throwable);
      }
    }, new Action0() {
      @Override
      public void call() {
        System.out.println("completed");
      }
});
  • 第 1 行:类型为 [Subscription] 的结果被忽略;
  • 第 1-15 行:这三个参数是匿名类的实例。我们还将使用 lambda 表达式。匿名类的优势在于,可以清晰地看到这些类中唯一方法所期望的数据类型;
  • 第 2-5 行:实现类型为 [Action1<Integer>] 的第一个参数;
  • 第 6-10 行:实现类型为 [Action1<Throwable>] 的第二个参数;
  • 第11-15行:实现类型为[Action0]的第三个参数;

完整代码如下:


package dvp.rxjava.observables;

import rx.Observable;
import rx.functions.Action0;
import rx.functions.Action1;

import java.util.Arrays;

public class Exemple01 {
  public static void main(String[] args) {
    // 整数可观察量
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    // 订阅
    obs1.subscribe(new Action1<Integer>() {
      @Override
      public void call(Integer integer) {
        System.out.printf("next : %s%n", integer);
      }
    }, new Action1<Throwable>() {
      @Override
      public void call(Throwable throwable) {
        System.out.println(throwable);
      }
    }, new Action0() {
      @Override
      public void call() {
        System.out.println("completed");
      }
    });
  }
}

第12行的可观察对象会在第14行调用方法[subscribe]时开始发布其3个元素。从这一刻起:

  • 每当发出一个元素时,第15-18行就会执行。
  • 当 3 个元素全部发送完毕后,第 24-29 行代码将执行;
  • 第19-24行将永远不会执行,因为可观察对象在此处未抛出异常;

默认情况下,可观察对象和观察者运行在同一个线程中。虽然存在一些预定义的可观察对象会在主线程(此处指 main 方法的线程)之外的线程中运行,但大多数情况并非如此。 因此,此处所有操作均在 [main] 方法的线程中进行:

  • 可观察对象发布元素 1;
  • 第 15-18 行代码执行并显示该元素;
  • 可观察对象发布元素 2;
  • 第 15-18 行代码执行并显示该元素;
  • 可观察对象发布元素 3;
  • 第 15-18 行代码执行并显示该元素;
  • 可观察对象发布通知 [completed];
  • 第 24-29 行代码执行;

以下是获得的结果:

1
2
3
4
next : 1
next : 2
next : 3
completed

类 [Exemple02] 继承了 [Exemple01],但这次使用了 lambda 函数作为方法 [Observable.subscribe] 的参数:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;

public class Exemple02 {
  public static void main(String[] args) {
    // 整数可观察量
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    // 订阅
    obs1.subscribe(
      (integer) -> System.out.printf("next : %s%n", integer),
      (th) -> System.out.println(th),
      () -> System.out.println("completed"));
  }
}

7.1.2. 示例-03:Observer类

  

用于订阅可观察对象的 [Observable.subscribe] 方法有多种版本,其中包括以下版本:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;
import rx.Observer;

public class Exemple03 {
    public static void main(String[] args) {
        // 整数可观察量
        Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
        // 订阅
        obs1.subscribe(new Observer<Integer>() {
            @Override
            public void onCompleted() {
                System.out.println("completed");
            }

            @Override
            public void onError(Throwable th) {
                System.out.printf("throwable %s", th);
            }

            @Override
            public void onNext(Integer integer) {
                System.out.printf("next : %s%n", integer);
            }
        });
    };
}

第 13 行,不再向方法 [subscribe] 传递三个参数,而是传递如下 [Observer] 类型:

 

类型 [Observer] 是一个包含三个方法的接口:

  • [onNext(T t)],该方法在可观察对象每次发出元素 t 时被调用;
  • [onError(Throwable th)]:当可观察对象抛出异常 th 时被调用;
  • [onCompleted]:当可观察对象指示其已完成发布时被调用;

该代码的工作原理与前文所述类似。结果如下:

1
2
3
4
next : 1
next : 2
next : 3
completed

7.1.3. 示例-04:方法 [Observable.create]

  

静态方法 Observable.create 的定义如下:

 
  • 方法 [create] 返回类型 Observable<T>;
  • 方法 [create] 的参数是一个类型为 [Observable.OnSubscribe<T>] 的函数,其定义如下:
 

类型 [Observable.OnSubscribe<T>] 是一个函数接口,它本身继承自函数接口 [Action1<Subscriber<? super T>>]。 该接口的 [call] 方法期望接收一个 [Subscriber] 类型(订阅者、监听者),其定义如下:

 

从 [1] 可见,类 [Subscriber<T>] 实现了第 7.1.2 节中介绍的 [Observer<T>] 接口。

最终,方法 [<T> Observable.create]:

  • 期望的参数是一个类型为 [Observable.OnSubscribe<T>] 的实例,该实例具有唯一的方法签名:void call(Subscriber<T> s)。 类型 [Subscriber<T>] 继承自类型 [Observer<T>],因此拥有方法 onNextonErroronCompleted
  • 返回一个 Observable<T> 类型;

方法 [<T> Observable.create] 返回一个已配置的可观察对象。 目前尚未发布任何元素。当 [Subscriber<T> s] 订阅者订阅此可观察对象时,将调用作为 [<T> Observable.create] 方法参数传递的函数中的 [void call(s)] 方法。 该方法的作用是发布类型为 T 的 t 元素,并在每次发布时调用观察者的 [s.onNext(t)] 方法。 当该方法执行完毕后,必须调用观察者的 [s.onCompleted(t)] 方法,并终止 [call] 方法。 如果方法 [call] 遇到异常 th,则应调用观察者的方法 [s.onError(th)],并终止方法 [call];

为说明这一复杂的工作原理,我们将使用以下代码 [Exemple04]:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;

import java.util.Random;

public class Exemple04 {
    public static void main(String[] args) {
        // 实数可观测配置
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                for (int i = 0; i < 3; i++) {
                    // 元素 i 的发布
                    subscriber.onNext(new Random((i + 1)).nextDouble());
                }
                // 发布结束
                subscriber.onCompleted();
            }
        });
        // 订阅,因此也发布
        obs1.subscribe((d) -> System.out.printf("onNext %s%n", d), (th) -> System.out.printf("onError %s%n", th),
                () -> System.out.println("onCompleted"));
    }
}
  • 第 11 行:创建一个发布 Double 类型的可观察对象;
  • 第 11-21 行:[create] 方法的参数通过一个匿名类进行实例化,该类包含第 12-20 行中唯一的 [call] 方法。 第 11 行创建的可观察对象已准备好进行发布,但只有当观察者加入时才会发布;
  • 第 13-21 行:方法 [call] 接收了一个观察者的引用;
  • 第14-17行:向观察者发送3个元素;
  • 第19行:向观察者发送发送结束通知;
  • 第23-24行:订阅第11行的可观察对象。通过三个lambda表达式实现方法[subscribe]的三个参数[onNext, onError, onCompleted]。 此订阅将创建订阅者 [Subscriber<Double>],该订阅者将被传递给第 13 行中的方法 [call]。随后将开始发布元素;
  • 所有操作均在同一线程中进行:可观察对象和观察者;

结果如下:

1
2
3
4
onNext 0.7308781907032909
onNext 0.7311469360199058
onNext 0.731057369148862
onCompleted

[Observable.create]方法允许基于任何事件创建可观察对象。我们在探索部分的第2段中正是使用了该方法,将同步接口转换为异步接口。

7.1.4. 示例-05:[Exemple-04]的重构

  

以下示例展示了静态方法 [Observable.subscribe] 的新版本:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Random;

public class Exemple05 {
    public static void main(String[] args) {
        // 实数可观测量的配置
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                showInfos("Observable.call start");
                for (int i = 0; i < 3; i++) {
                    // 等待
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // 错误
                        subscriber.onError(e);
                    }
                    // 操作
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // 完成
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // 一个订阅者
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
            }

            @Override
            public void onError(Throwable e) {
                showInfos(String.format("Subscriber.onError (%s)", e));
            }

            @Override
            public void onNext(Double aDouble) {
                showInfos(String.format("Subscriber.onNext (%s)", aDouble));
            }
        };

        // 订阅
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        showInfos("après souscription");

    }

    private static void showInfos(String message) {
        System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message, Thread.currentThread().getName(),
                new SimpleDateFormat("ss:SSS").format(new Date()));
    }
}
  • 第 56 行:静态方法 [Observable.subscribe] 的新版本接受我们在上一段中介绍的类型 [Subscriber] 作为参数;
  • 第37-52行:订阅者(订阅者、观察者)。 它实现了接口 Observer 及其三个方法 onNextonErroronCompleted
  • 第 61-64 行:接下来我们将关注可观察对象及其观察者所运行的线程;
  • 第 62 行:线程名称;
  • 第 63 行:当前时间,以秒和毫秒为单位。这将使我们能够观察可观察对象在时间轴上发布元素以及观察者对其进行处理的过程;
  • 这段代码的功能与前面的代码相同,只是对后者进行了重构;

所得结果如下:

avant souscription ------Thread[main] ---- Time[31:685]
Observable.call start ------Thread[main] ---- Time[31:691]
Observable.call onNext(80.39999999999999) ------Thread[main] ---- Time[32:194]
Subscriber.onNext (80.39999999999999) ------Thread[main] ---- Time[32:195]
Observable.call onNext(73.2) ------Thread[main] ---- Time[32:595]
Subscriber.onNext (73.2) ------Thread[main] ---- Time[32:595]
Observable.call onNext(106.8) ------Thread[main] ---- Time[32:897]
Subscriber.onNext (106.8) ------Thread[main] ---- Time[32:897]
Observable.call onCompleted ------Thread[main] ---- Time[32:898]
Subscriber.onCompleted ------Thread[main] ---- Time[32:898]
après souscription ------Thread[main] ---- Time[32:899]
  • 结果第1行:在代码第56行之前,尚未发生任何变化。可观察对象仅被配置;
  • 结果第2行:代码第56行触发了第15行中[call]方法的调用。第3行,实数80.39被发送给观察者;
  • 第4行:观察者接收到了发送的数值;
  • 第5-8行:上述过程重复了2次;
  • 第9行:可观察对象发送发送结束通知;
  • 第10行:观察者接收该通知;
  • 第11行:由代码第57行显示;

由此可见,仅第56行的订阅操作就触发了第2至10行结果的显示。初次接触RxJava库时,人们往往会疑惑各项操作是如何相互关联的,尤其是观察者与被观察对象之间的联系。 这里可以看到,第56行(对可观察量的订阅)

  • 触发了可观察对象所有元素的发布;
  • 可观察对象与观察者运行在同一线程中;
  • 正因如此,我们观察到以下序列:发出元素 i,观察元素 i,发出元素 (i+1),观察元素 (i+1),……

我们记得,发布者在发布元素之前会进行等待:


                    // 待处理
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // 错误
                        subscriber.onError(e);
}

其中第3行中的i代表发送序号(0<=i<3)。若观察可观测对象各元素的发送时间:

  • 第2、3行:元素0在订阅开始后约500毫秒被发送;
  • 第3、5行:元素1在元素0之后约400毫秒发出;
  • 第5、7行:元素2在元素1之后约300毫秒被发布;

7.2. 执行线程、观察线程

7.2.1. 示例-06:[main]以外线程中的可观察对象与观察者

  

我们将前面的示例 [Exemple06] 重构如下:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;
import rx.schedulers.Schedulers;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Random;
import java.util.concurrent.CountDownLatch;

public class Exemple06 {
    public static void main(String[] args) {

        // 门卫
        CountDownLatch latch = new CountDownLatch(1);

        // 实数观测量的配置
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                showInfos("Observable.call start");
                for (int i = 0; i < 3; i++) {
                    // 等待
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // 错误
                        subscriber.onError(e);
                    }
                    // 操作
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // 完成
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // 一个订阅者
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // 降低门槛
                latch.countDown();
            }

            @Override
            public void onError(Throwable e) {
                showInfos(String.format("Subscriber.onError (%s)", e));
            }

            @Override
            public void onNext(Double aDouble) {
                showInfos(String.format("Subscriber.onNext (%s)", aDouble));
            }
        };

        // 可观察配置后续
        obs1 = obs1.subscribeOn(Schedulers.computation());
        // 认购
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // 在闸门前等待
        try {
            showInfos("début attente barrière");
            latch.await();
            showInfos("fin attente barrière");
        } catch (InterruptedException e1) {
            System.out.println(e1);
        }
        showInfos("après souscription");

    }

    private static void showInfos(String message) {
        System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message, Thread.currentThread().getName(),
                new SimpleDateFormat("ss:SSS").format(new Date()));
    }
}
  • 第 16 行:创建一个带有 [CountDownLatch] 类型对象的屏障(信号量)。该对象用于在线程之间进行同步。此处将其初始化为值 1,我们称之为屏障(或信号量)的值。一个线程通过以下操作进入屏障等待状态:

latch.await();

如果栅栏值 > 0,线程将被阻塞。线程可以增加或减少栅栏的内部值。第 48 行,栅栏值被减 1。

  • 第63行:可观察对象被配置为在调度器[Schedulers.computation()]提供的线程上执行。该调度器可提供的线程数量与执行机器上的核心数量相同。 示例应用章节展示了其他调度器的使用(参见第2.8节);

代码的工作原理如下:

  • 方法 [main] 在主线程(main)中执行;
  • 第 66 行:触发可观察对象的元素发布。这些元素将在与主线程不同的线程上发布;
  • 第70行:由于门禁值设为1(参见第16行),主线程被阻塞。只有当该值变为0时,主线程才能继续执行。这一变化发生在第48行。当观察者接收到可观察对象已完成发出的通知时,会将其门禁值设为0;

运行结果如下:

avant souscription ------Thread[main] ---- Time[09:268]
Observable.call start ------Thread[RxComputationThreadPool-1] ---- Time[09:278]
début attente barrière ------Thread[main] ---- Time[09:278]
Observable.call onNext(44.4) ------Thread[RxComputationThreadPool-1] ---- Time[09:783]
Subscriber.onNext (44.4) ------Thread[RxComputationThreadPool-1] ---- Time[09:783]
Observable.call onNext(18.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:183]
Subscriber.onNext (18.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:184]
Observable.call onNext(54.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:486]
Subscriber.onNext (54.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:488]
Observable.call onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[10:489]
Subscriber.onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[10:490]
fin attente barrière ------Thread[main] ---- Time[10:491]
après souscription ------Thread[main] ---- Time[10:493]
  • 第 1 行:订阅即将开始;
  • 第2行:这将触发在[RxComputationThreadPool-1]线程上执行[call]方法。现在有两个线程并行执行;
  • 第 3 行:出于不明原因,线程 [RxComputationThreadPool-1] 已释放控制权。 随后线程 [main] 接管了控制权,但被守护程序阻塞(代码第 70 行)。从这一刻起,只有线程 [RxComputationThreadPool-1] 能够运行;
  • 第4-11行:可观察对象与其观察者之间呈现出先前观察到的行为,但此时所有操作均在线程[RxComputationThreadPool-1]中进行;
  • 第12-13行:观察者已放下道闸(代码第48行),线程[RxComputationThreadPool-1]终止。线程[main]接管并显示两条消息;

7.2.2. 示例-07:位于两个不同线程中的可观察对象和观察者

  

我们将前面的示例修改如下:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;
import rx.schedulers.Schedulers;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Random;
import java.util.concurrent.CountDownLatch;

public class Exemple07 {
    public static void main(String[] args) {

        // 道闸看守
        CountDownLatch latch = new CountDownLatch(1);

        // 实数观测量的配置
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                showInfos("Observable.call start");
                for (int i = 0; i < 3; i++) {
                    // 等待
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // 错误
                        subscriber.onError(e);
                    }
                    // 操作
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // 完成
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // 一个订阅者
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // 降低门槛
                latch.countDown();
            }

            @Override
            public void onError(Throwable e) {
                showInfos(String.format("Subscriber.onError (%s)", e));
            }

            @Override
            public void onNext(Double aDouble) {
                showInfos(String.format("Subscriber.onNext (%s)", aDouble));
            }
        };

        // 可观察配置后续
        obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());
        // 认购
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // 等待触及门槛
        try {
            showInfos("début attente barrière");
            latch.await();
            showInfos("fin attente barrière");
        } catch (InterruptedException e1) {
            System.out.println(e1);
        }
        showInfos("après souscription");

    }

    private static void showInfos(String message) {
        System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message, Thread.currentThread().getName(),
                new SimpleDateFormat("ss:SSS").format(new Date()));
    }
}

代码与前一个示例完全相同,仅第 63 行有所不同:


obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());

该行配置了可观察对象(subscribeOn)和观察者(observeOn),使其在调度器[Schedulers.computation()]提供的线程之一上运行。

所得结果如下:

avant souscription ------Thread[main] ---- Time[09:643]
début attente barrière ------Thread[main] ---- Time[09:656]
Observable.call start ------Thread[RxComputationThreadPool-4] ---- Time[09:656]
Observable.call onNext(39.6) ------Thread[RxComputationThreadPool-4] ---- Time[10:162]
Subscriber.onNext (39.6) ------Thread[RxComputationThreadPool-3] ---- Time[10:163]
Observable.call onNext(98.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[10:562]
Subscriber.onNext (98.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[10:564]
Observable.call onNext(46.8) ------Thread[RxComputationThreadPool-4] ---- Time[10:864]
Observable.call onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[10:866]
Subscriber.onNext (46.8) ------Thread[RxComputationThreadPool-3] ---- Time[10:866]
Subscriber.onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[10:868]
fin attente barrière ------Thread[main] ---- Time[10:869]
après souscription ------Thread[main] ---- Time[10:870]

可以注意到以下几点:

  • 可观察对象在线程 [RxComputationThreadPool-4] 中执行(第 3-4 行、第 6 行、第 8-9 行);
  • 观察者在线程 [RxComputationThreadPool-3] 中运行(第 5、7、10-11 行);
  • 它们各自独立运行。因此,在第 8-9 行,可观察对象会先发出 2 条通知(onNext、onCompleted),随后观察者才获取通知 [onNext](第 10 行);

RxJava 库负责将数据(即发布)从可观察对象的线程传递到观察者的线程。开发者无需为此操心。

我们已经了解了如何创建可观察对象(Observable.from、Observable.create)。现在,我们将了解库中预定义的可观察对象 RxJava。

7.3. 预定义的可观察对象

7.3.1. 示例-08:[Observable.range]方法

 

从现在开始,我们将使用专门的类来处理被观察的进程及其观察者。其目的是能够记录它们的名称、执行线程和执行时间,以便随时间推移进行追踪。

类 [Process] 仅是一个可命名的可观察对象。它将实现以下接口 [IProcess]:


package dvp.rxjava.observables.utils;

import rx.Observable;

public interface IProcess<T> {

    // 可观察对象名称
    public String getName();

    // 可观测项
    public Observable<T> getObservable();

}

该接口可由以下 [Process<T>] 类实现:


package dvp.rxjava.observables.utils;

import rx.Observable;
import rx.Scheduler;

public class Process<T> implements IProcess<T>{

    // 可观测对象名称
    protected String name;
    // 被观测过程
    protected Observable<T> observable;

    // 构造函数
    public Process(String name, Observable<T> observable) {
        // 局部初始化
        this.name = name;
        this.observable = observable;
    }

    // 获取器和设置器
    public String getName() {
        return name;
    }

    public Observable<T> getObservable() {
        return observable;
    }

}
  • 第9行:进程名称;
  • 第11行:被观测的观测量;
  • 第14-18行:构造函数;

观察者将由以下类 [Observateur] 描述:


package dvp.rxjava.observables.utils;

import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

import rx.Subscriber;

public class Observateur<T> extends Subscriber<T> {

...
}
  • 第 11 行,类 Observateur<T> 继承了我们在第 7.1.3 节中简要介绍过的类 Subscriber<T>。我们将它用作方法 [Observable.subscribe] 的参数:

// 可观察的执行(观察)
obs1.subscribe(observateur);

上文第2行中使用的[Observable.subscribe]方法定义如下:

 

[Subscriber] 的主要作用是通过 [Observer] 接口的方法,管理其订阅的可观察对象所发出的元素: onNextonErroronCompleted。[Subscriber] 类具有以下方法:

 

在类 [Observateur] 的代码中,我们将使用方法 [1] isUnsubscribed,以判断投保人的投保是否已被取消。 完整的 [Observateur<T>] 类如下:


package dvp.rxjava.observables.utils;

import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

import rx.Subscriber;

public class Observateur<T> extends Subscriber<T> {

    // 一个信号量
    private CountDownLatch latch;
    // 显示方法
    private Consumer<String> showInfos;
    // 观察者的名称
    private String observerName;
    // 被观察进程的名称
    private String processName;

    // 构造函数
    public Observateur() {

    }

    public Observateur(String name, CountDownLatch latch, Consumer<String> showInfos, String observedName) {
        this.observerName = name;
        this.latch = latch;
        this.showInfos = showInfos;
        this.processName = observedName;
    }

    // --------------------------- Observer<T> 接口的实现
    @Override
    public void onCompleted() {
        // 结束广播
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s,%s].onCompleted", observerName, processName));
        }
        // 主线程阻塞结束
        latch.countDown();
    }

    @Override
    public void onError(Throwable e) {
        // 发送错误
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s, %s].onError (%s)", observerName, processName, e));
        }
    }

    @Override
    public void onNext(T value) {
        // 额外一次发布
        if (!isUnsubscribed()) {
            try {
                showInfos.accept(String.format("Subscriber [%s,%s] : onNext (%s)", observerName, processName,
                        new ObjectMapper().writeValueAsString(value)));
            } catch (JsonProcessingException e) {
                showInfos.accept(String.format("Subscriber [%s,%s].onNext (%s)", observerName, processName, e));
            }
        }
    }
}
  • 除了具备 Subscriber 的特性外,观察者 Observateur 还将携带以下信息:
    • 第 14 行:一个互斥锁或信号量,用于阻塞主线程,直到观察者收到可观察对象发出的所有元素。这将在代码第 36 行实现,即当观察者收到可观察对象发出的发送结束通知时;
    • 第 16 行:一个 Consumer<String> 实例,用于在控制台显示一条消息;
    • 第18行:观察者的名称,用于在存在多个观察者时进行区分;
    • 第 20 行:被观察进程的名称;
  • 第 36、46、54 行:接口 [Observer<T>] 中的方法 [onCompleted, onError, onNext],该接口由抽象类 [Subscriber<T>] 实现。该类并未实现这些方法。 因此必须在子类中实现这些方法。在这些方法中执行任何操作之前,需检查观察者是否已从其所观察的可观察对象上取消订阅;
  • 第 59 行:观察者的 [onNext] 方法将接收到的元素的字符串 jSON 写入。这将使我们能够显示各种类型的元素;

基于此,让我们研究 Observable 类中的另一个方法——[range] 方法:

 

可观测对象 Observable.range(n,m) 输出 (m) 个整数,取值范围从 n 到 n+m-1。我们通过以下 [Exemple08] 代码来研究它:


package dvp.rxjava.observables.exemples;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import dvp.rxjava.observables.utils.Observateur;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple08 {
    public static void main(String[] args) throws InterruptedException {

        // 观察者数量
        final int nbObservateurs = 2;

        // 信号量
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // 可观察配置
        Observable<Integer> obs1 = Observable.range(15, 3).subscribeOn(Schedulers.computation());
        // 可观察的执行(观察)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            obs1.subscribe(new Observateur<>(String.format("observateur[%d]", i), latch, showInfos,"obs1"));
        }
        // 等待
        showInfos.accept("main : attente fin observation");
        latch.await();
        // 结束
        showInfos.accept("main : fin observation");
    }

    // 显示
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}
  • 第 16 行:我们将使用两个观察者;
  • 第19行:信号量初始化为2,因为每个观察者将运行在不同的线程上。因此主线程需要等待两个观察线程结束;
  • 第22行:将可观察对象配置为在调度器[Schedulers.computation()]的线程上运行。观察者将与可观察对象位于同一线程;
  • 第25-27行:为可观察对象订阅两个观察者。这将触发可观察对象为每个观察者完整执行:整数15、16和17将被发布;
  • 第 30 行:主线程等待观察者完成;

所得结果如下:

main : début observation ------Thread[main] ---- Time[27:875]
main : attente fin observation ------Thread[main] ---- Time[27:893]
Subscriber[observateur[1],obs1] : onNext (15) ------Thread[RxComputationThreadPool-2] ---- Time[28:245]
Subscriber[observateur[0],obs1] : onNext (15) ------Thread[RxComputationThreadPool-1] ---- Time[28:245]
Subscriber[observateur[1],obs1] : onNext (16) ------Thread[RxComputationThreadPool-2] ---- Time[28:247]
Subscriber[observateur[0],obs1] : onNext (16) ------Thread[RxComputationThreadPool-1] ---- Time[28:248]
Subscriber[observateur[1],obs1] : onNext (17) ------Thread[RxComputationThreadPool-2] ---- Time[28:249]
Subscriber[observateur[1],obs1].onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[28:250]
Subscriber[observateur[0],obs1] : onNext (17) ------Thread[RxComputationThreadPool-1] ---- Time[28:251]
Subscriber[observateur[0],obs1].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[28:252]
main : fin observation ------Thread[main] ---- Time[28:252]
  • 第2行:主线程被阻塞,等待两个观察者的执行结束;
  • 第3-4行:可以看到观察者0位于线程[RxComputationThreadPool-1]上,观察者1位于线程[RxComputationThreadPool-2]上;
  • 第3-10行:可以看到两个观察者接收到的元素完全相同;

我们将使用如此定义的 Observateur 类来说明其他类型可观察对象的行为。

7.3.2. 示例-09:Observable.[interval, take, doNext]的方法

  
 

此示例演示了可观察对象 Observable.interval(长整型,TimeUnit 单位)的使用,该对象以固定时间间隔发布长整型数据。 需要注意的是 [1] 点:默认情况下,可观察对象 [Observable.interval] 在调度器 [Schedulers.computation] 的某个线程上运行。

代码如下:


package dvp.rxjava.observables.exemples;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;

import dvp.rxjava.observables.utils.Observateur;
import rx.Observable;

public class Exemple09 {
    public static void main(String[] args) throws InterruptedException {

        // 观察者数量
        final int nbObservateurs = 2;

        // 信号量
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // 可观察配置
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // 可观察的执行(观察)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            obs1.subscribe(new Observateur<>(String.format("observateur [%d]", i), latch, showInfos,
                    "obs1"));
        }
        // 等待
        showInfos.accept("main : attente fin observation");
        latch.await();
        // 结束
        showInfos.accept("main : fin observation");
    }

    // 显示
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}
  • 第 22 行:该可观察对象每 500 毫秒输出一个长整型。序列从数字 0 开始;
  • 第 22 行:该可观察对象会发出无限数量的值。方法 [Observable.take(n)] 创建了一个新的可观察对象,该对象仅保留前 n 个已发出的元素;
 

让我们回顾一下可观察对象的代码:


Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
.doOnNext(l -> showInfos.accept(l.toString()));

第 2 行,每当可观察对象发布新元素时,方法 [Observable.doOnNext] 就会执行。这通常用于记录信息。在此,我们希望记录元素的发布时间,以验证 500 毫秒的间隔是否得到正确执行。 方法 [Observable.doOnNext] 不会修改其所应用的可观察对象。其定义如下:

 

执行结果如下:

main : début observation ------Thread[main] ---- Time[55:892]
main : attente fin observation ------Thread[main] ---- Time[55:911]
0 ------Thread[RxComputationThreadPool-1] ---- Time[56:412]
0 ------Thread[RxComputationThreadPool-2] ---- Time[56:413]
Subscriber[observateur [1],obs1] : onNext (0) ------Thread[RxComputationThreadPool-2] ---- Time[56:723]
Subscriber[observateur [0],obs1] : onNext (0) ------Thread[RxComputationThreadPool-1] ---- Time[56:723]
1 ------Thread[RxComputationThreadPool-1] ---- Time[56:906]
Subscriber[observateur [0],obs1] : onNext (1) ------Thread[RxComputationThreadPool-1] ---- Time[56:908]
1 ------Thread[RxComputationThreadPool-2] ---- Time[56:912]
Subscriber[observateur [1],obs1] : onNext (1) ------Thread[RxComputationThreadPool-2] ---- Time[56:914]
2 ------Thread[RxComputationThreadPool-1] ---- Time[57:405]
Subscriber[observateur [0],obs1] : onNext (2) ------Thread[RxComputationThreadPool-1] ---- Time[57:407]
Subscriber[observateur [0],obs1].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[57:408]
2 ------Thread[RxComputationThreadPool-2] ---- Time[57:412]
Subscriber[observateur [1],obs1] : onNext (2) ------Thread[RxComputationThreadPool-2] ---- Time[57:414]
Subscriber[observateur [1],obs1].onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[57:415]
main : fin observation ------Thread[main] ---- Time[57:416]
  • 第 3、7 和 11 行:可以看到,发送间隔大致接近 500 毫秒;
  • 这两个观察者显然位于不同的线程上,尽管该可观察对象并未配置为在特定调度器下运行。这里展示的是可观察对象 [Observable.interval] 的默认行为;

7.3.3. 示例-10/12:Observable.[error, empty, never] 方法

 

接下来,我们将更简洁地演示 [Observable] 类的各种方法。之前的代码如下:


package dvp.rxjava.observables;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;

import rx.Observable;

public class Exemple09 {
    public static void main(String[] args) throws InterruptedException {

        // 观察者数量
        final int nbObservateurs = 2;

        // 信号量
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // 可观察配置
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // 可观察的执行(观察)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            obs1.subscribe(new Observateur<>(String.format("observateur [%d]", i), latch, showInfos,
                    "obs1"));
        }
        // 等待
        showInfos.accept("main : attente fin observation");
        latch.await();
        // 结束
        showInfos.accept("main : fin observation");
    }

    // 显示
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}

该代码已在前面的示例中使用过。仅第21-22行有所变化。因此,我们将把大部分代码提取到以下[ProcessUtils]类中:


package dvp.rxjava.observables.utils;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import rx.Observable;

public class ProcessUtils {

    @SafeVarargs
    public static void subscribe(int nbObservateurs, IProcess<?>... processes) throws InterruptedException {

        // 信号量
        CountDownLatch latch = new CountDownLatch(nbObservateurs * processes.length);

        // 可观察执行(观察)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            for (IProcess<?> process : processes) {
                Observable<?> obs = process.getObservable();
                obs.subscribe(new Observateur<>(String.format("observateur[%d]", i), latch, showInfos, process.getName()));
            }
        }
        // 等待
        showInfos.accept("main : attente fin observation");
        latch.await();
        // 结束
        showInfos.accept("main : fin observation");
    }

    // 显示
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}
  • 第13行:该方法接受两个参数:
    • nbObservateurs:作为第二个参数传递的进程的观察者数量;
    • processes:待观察的进程(即命名的可观察对象)。借助[IProcess<?>]这种表示法,进程可以发布不同类型的元素;
  • 第 16 行:当所有观察者完成所有观察操作后,信号量应转为绿色。因此,信号量的初始值为观察者数量乘以观察次数;
  • 第20-25行:将每个观察者订阅到需要观察的所有进程;
  • 第23行:从进程中获取可观察对象(参见第7.3.1节);
  • 第23行:为其订阅一个观察者。向该观察者传递4项信息:
    • 其名称;
    • 当其接收所观察可观测对象的发送结束通知时,需递减的信号量;
    • 当它需要将信息记录到控制台时应使用的方法;
    • 它将要观察的进程名称;

定义完这些类后,示例 10 将如下所示:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple10 {
    public static void main(String[] args) throws InterruptedException {
        // 可观察配置
        Observable<?> obs = Observable.error(new RuntimeException("Erreur !!!")).subscribeOn(Schedulers.computation());
        // 可观察的执行(观察)
        ProcessUtils.subscribe(2,new Process<>("process1", obs));
    }
}

第 11 行,静态方法 [Observable.error] 的定义如下:

 

因此,第 8 行配置了一个仅向其订阅者的 [onError] 方法抛出异常的可观察对象。执行结果如下:


main : début observation ------Thread[main] ---- Time[22:618]
main : attente fin observation ------Thread[main] ---- Time[22:636]
Subscriber[observateur[1], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-2] ---- Time[22:638]
Subscriber[observateur[0], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-1] ---- Time[22:638]

第 3 行和第 4 行,两个订阅者的 [onError] 方法都接收到了可观察对象抛出的异常。

此次执行有一个特殊之处:两个观察者的 [onCompleted] 方法均未被调用。因此,屏障未被放下,主线程仍被阻塞在静态方法 [ProcessUtils.subscribe] 的第 3 行:


// 等待
showInfos.accept("main : attente fin observation");
latch.await();
// 结束
showInfos.accept("main : fin observation");

由此发现,当可观察对象发生错误时,订阅者的 [onCompleted] 方法并未被调用。因此,我们将 [Observateur.onError] 方法修改如下:


    @Override
    public void onError(Throwable e) {
        // 发送错误
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber[%s, %s].onError (%s)", observerName, processName, e));
        }
        // 主线程阻塞结束
        latch.countDown();
}

我们添加了第7-8行代码,以在可观察对象发生错误时解除限制。使用此新代码,执行结果如下:


main : début observation ------Thread[main] ---- Time[40:750]
main : attente fin observation ------Thread[main] ---- Time[40:764]
Subscriber[observateur[0], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-1] ---- Time[40:766]
Subscriber[observateur[1], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-2] ---- Time[40:766]
main : fin observation ------Thread[main] ---- Time[40:767]

我们得到了之前未出现的第 5 行。

示例 11 如下:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple11 {
    public static void main(String[] args) throws InterruptedException {
        // 可观察配置
        Observable<?> obs1 = Observable.empty();
        // 可观察的执行(观察)
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

第10行,静态方法 [Observable.empty] 创建了一个不发布任何元素的可观察对象。它仅发布结束通知;

 

执行上述示例代码将得到以下结果:

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[37:073]
Subscriber[observateur[0],process1].onCompleted ------Thread[main] ---- Time[37:086]
Subscriber[observateur[1],process1].onCompleted ------Thread[main] ---- Time[37:086]
main : attente fin observation ------Thread[main] ---- Time[37:087]
main : fin observation ------Thread[main] ---- Time[37:087]
  • 第 2 行和第 3 行:可以看到两个观察者都收到了发射结束通知,但此前并未收到任何元素。

人们可能会质疑这种方法究竟有何用处。我们可以将其类比为一个集合:初始为空,随后向其中累积元素:

1
2
3
4
Observable obs=Observable.empty() ;
for(Observable o : observables){
    obs=obs.mergeWith(o) ;
}

第3行,将初始可观测量obs(第1行)与其他可观测量合并。

示例 12 演示了静态方法 [Observable.never]:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple12 {
    public static void main(String[] args) throws InterruptedException {
        // 可观察配置
        Observable<?> obs1 = Observable.never();
        // 可观察的执行(观察)
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

静态方法 [Observable.never] 创建了一个永远不会发出的可观察对象:

 

运行该示例将得到以下结果:

main : début observation ------Thread[main] ---- Time[27:018]
main : attente fin observation ------Thread[main] ---- Time[27:030]

第 2 行,主线程处于无限等待状态。这是因为没有任何可观察对象发出 [onCompleted] 通知,该通知用于将信号量(道闸)置为绿色(放下道闸)。

7.4. Multi-threading

7.4.1. 示例-13:操作线程与观察线程

我们在第 7.1.3 节中使用静态方法 [Observable.create] 创建了一个可观察对象:

 
  • 方法 [create] 返回类型 Observable<T>;
  • 方法 [create] 的参数是一个类型为 [Observable.OnSubscribe<T>] 的函数,定义如下:
 

类型 [Observable.OnSubscribe<T>] 是一个函数接口,它本身继承自函数接口 [Action1<Subscriber<? super T>>]。 该接口的 [call] 方法期望接收类型 [Subscriber](订阅者、订阅方、观察者)。在本文档的后续内容中,我们将有时将类型 [Observable.OnSubscribe<T>] 称为操作。 我们将创建具有名称的自定义操作。这些操作将是以下 [IProcessAction] 接口的实例:

  

package dvp.rxjava.observables.utils;

import rx.Observable;

public interface IProcessAction<T> extends Observable.OnSubscribe<T> {

    // 该操作有一个名称
    public String getName();
}
  • 第 5 行:接口 [IProcessAction<T>] 具有接口 [Observable.OnSubscribe<T>] 的所有特征;
  • 第 8 行:它还具有一个名为 [getName] 的方法,该方法返回实现该接口的实例的名称;

我们将使用名为 [ProcessAction01] 的以下操作:


package dvp.rxjava.observables.utils;

import java.util.Random;

import rx.Subscriber;
import rx.functions.Func1;

public class ProcessAction01<T> implements IProcessAction<T> {

    // 数据
    private String name;
    private int nbValues;
    private Func1<Integer, T> func1;

    // 构造函数
    public ProcessAction01(String name, int nbValues, Func1<Integer, T> func1) {
        this.name = name;
        this.nbValues = nbValues;
        this.func1 = func1;
    }

    @Override
    public void call(Subscriber<? super T> subscriber) {
        ProcessUtils.showInfos.accept(String.format("Observable (%s) call start", getName()));
        for (int i = 0; i < nbValues; i++) {
            // 等待
            try {
                Thread.sleep(new Random().nextInt(500));
            } catch (InterruptedException e) {
                // 错误
                ProcessUtils.showInfos.accept(String.format("Observable (%s) onError", getName()));
                subscriber.onError(e);
            }
            // 发送元素
            T value = func1.call(i);
            ProcessUtils.showInfos.accept(String.format("Observable (%s,%s) onNext (%s)", getName(), i, value));
            subscriber.onNext(value);
        }
        // 完成
        ProcessUtils.showInfos.accept(String.format("Observable (%s) onCompleted", getName()));
        subscriber.onCompleted();
    }

    @Override
    public String getName() {
        return name;
    }

}
  • 第 8 行:类 [ProcessAction01<T>] 实现了接口 [IProcessAction<T>],因此也实现了接口 [Observable.OnSubscribe<T>];
  • 第 11 行:操作名称;
  • 第 12 行:要发出的值个数;
  • 第13行:一个类型为[Func1<Integer, T>]的实例,它根据一个整数生成类型T,该类型将由可观察对象发布(第35行和第37行);
  • 第16-20行:将操作名称、待发出的值个数以及发射函数传递给构造函数;
  • 第 23-42 行:进程的代码;
  • 第 23 行:方法 [call] 接收与该进程关联的可观察对象的订阅者作为参数;
  • 第28行:流程在经过随机时长的等待后发布其元素;
  • 第 32 行:触发错误;
  • 第 37 行:一次正常发布;
  • 第 41 行:发送发送结束通知;
  • 第 25-38 行:该操作在经过随机等待时间(第 30 行)后,发布 nbValues 实数;
  • 第35行:要发出的值由作为构造函数参数(第16行)传递的函数[func1]提供;

我们重构类 [Process](参见第 7.3.1 节),使其也能通过命名操作进行构造。为此,我们为其添加了以下构造函数:


public Process(IProcessAction<T> na, Scheduler schedulerObserved, Scheduler schedulerObserver) {
        // 进程名称=操作名称
        name = na.getName();
        // 操作 --> 可观察对象
        observable = Observable.create(na);
        // 被观察进程的执行线程
        if (schedulerObserved != null) {
            observable = observable.subscribeOn(schedulerObserved);
        }
        // 观察员的观察线程
        if (schedulerObserver != null) {
            observable = observable.observeOn(schedulerObserver);
        }
    }
  • 第 1 行,构造函数接受 3 个参数:
    1. 用于构建可观察对象的命名操作(第5行);
    2. 被观察进程的调度器(可能是 null);
    3. 观察者的调度器(例如 null);
  • 第 5 行:根据作为参数传递的操作创建可观察对象;

以下代码 [Exemple13] 观察了不同的可观察对象:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple13 {
    public static void main(String[] args) throws InterruptedException {
        // 进程 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 1, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // 进程 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // 进程 3
        Process<Integer> process3 = new Process<>(new ProcessAction01<Integer>("process3", 3, i -> i * 2), null,
                Schedulers.computation());
        // 进程 4
        Process<Boolean> process4 = new Process<>(new ProcessAction01<Boolean>("process4", 4, i -> i % 2 == 0), null, null);
        // 订阅
        ProcessUtils.subscribe(1, process1);
        ProcessUtils.subscribe(1, process2);
        ProcessUtils.subscribe(1, process3);
        ProcessUtils.subscribe(1, process4);
    }
}
  • 第13-15行:进程process1在计算线程上生成1个实数,该值将在另一个计算线程上被观察;
  • 第 17-18 行:进程 process2 在一个计算线程上生成 2 个字符串,且未指定观察者的线程。结果表明,默认情况下观察操作在与进程执行相同的线程上进行;
  • 第20-21行:进程process3在未指定线程上生成3个整数,这些整数将在计算线程上被观察。结果表明,该进程默认在主线程上执行;
  • 第23行:进程process4在未指定线程上生成4个布尔值,这些值将在未指定线程上被观察。结果表明,该进程的执行及其观察默认都在主线程上进行;

该代码的执行结果如下:

main : début observation ------Thread[main] ---- Time[18:642]
main : attente fin observation ------Thread[main] ---- Time[18:660]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[18:660]
Observable (process1,0) onNext (68.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[19:093]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[19:094]
Subscriber[observateur[0],process1] : onNext (68.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[19:396]
Subscriber[observateur[0],process1].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[19:397]
main : fin observation ------Thread[main] ---- Time[19:397]
main : début observation ------Thread[main] ---- Time[19:398]
main : attente fin observation ------Thread[main] ---- Time[19:399]
Observable (process2) call start ------Thread[RxComputationThreadPool-5] ---- Time[19:399]
Observable (process2,0) onNext (valeur-0) ------Thread[RxComputationThreadPool-5] ---- Time[19:630]
Subscriber[observateur[0],process2] : onNext ("valeur-0") ------Thread[RxComputationThreadPool-5] ---- Time[19:631]
Observable (process2,1) onNext (valeur-1) ------Thread[RxComputationThreadPool-5] ---- Time[20:094]
Subscriber[observateur[0],process2] : onNext ("valeur-1") ------Thread[RxComputationThreadPool-5] ---- Time[20:095]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[20:096]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[20:096]
main : fin observation ------Thread[main] ---- Time[20:097]
main : début observation ------Thread[main] ---- Time[20:097]
Observable (process3) call start ------Thread[main] ---- Time[20:098]
Observable (process3,0) onNext (0) ------Thread[main] ---- Time[20:188]
Subscriber[observateur[0],process3] : onNext (0) ------Thread[RxComputationThreadPool-6] ---- Time[20:213]
Observable (process3,1) onNext (2) ------Thread[main] ---- Time[20:336]
Subscriber[observateur[0],process3] : onNext (2) ------Thread[RxComputationThreadPool-6] ---- Time[20:338]
Observable (process3,2) onNext (4) ------Thread[main] ---- Time[20:676]
Observable (process3) onCompleted ------Thread[main] ---- Time[20:677]
main : attente fin observation ------Thread[main] ---- Time[20:677]
Subscriber[observateur[0],process3] : onNext (4) ------Thread[RxComputationThreadPool-6] ---- Time[20:678]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[20:679]
main : fin observation ------Thread[main] ---- Time[20:679]
main : début observation ------Thread[main] ---- Time[20:680]
Observable (process4) call start ------Thread[main] ---- Time[20:680]
Observable (process4,0) onNext (true) ------Thread[main] ---- Time[21:065]
Subscriber[observateur[0],process4] : onNext (true) ------Thread[main] ---- Time[21:067]
Observable (process4,1) onNext (false) ------Thread[main] ---- Time[21:187]
Subscriber[observateur[0],process4] : onNext (false) ------Thread[main] ---- Time[21:188]
Observable (process4,2) onNext (true) ------Thread[main] ---- Time[21:624]
Subscriber[observateur[0],process4] : onNext (true) ------Thread[main] ---- Time[21:625]
Observable (process4,3) onNext (false) ------Thread[main] ---- Time[21:765]
Subscriber[observateur[0],process4] : onNext (false) ------Thread[main] ---- Time[21:766]
Observable (process4) onCompleted ------Thread[main] ---- Time[21:767]
Subscriber[observateur[0],process4].onCompleted ------Thread[main] ---- Time[21:767]
main : attente fin observation ------Thread[main] ---- Time[21:767]
main : fin observation ------Thread[main] ---- Time[21:768]
  • 进程 process1 在计算线程 [RxComputationThreadPool-4] 上生成 1 个实数(第 4 行),该值在计算线程 [RxComputationThreadPool-3] 上被观察到(第 6 行);
  • 进程 process2 在计算线程 [RxComputationThreadPool-5] 上生成 2 个字符串(第 12、14 行),这些字符串在同一线程上被观察到(第 13、15 行);
  • 进程 process3 在主线程上生成 3 个整数(第 21、23、25 行),这些整数在计算线程 [RxComputationThreadPool-6] 上被观察到(第 22、24、28 行);
  • 进程 process4 在主线程上生成 4 个布尔值(第 34、36、38、40 行),这些值在同一主线程上被观察到(第 33、35、37、39 行);

请读者关注以下内容:

  • 被观察进程及其线程的生命周期;
  • 其观察者的生命周期及其线程;

Rx 库的许多优势正源于这种多线程特性,开发者无需自行管理。

7.5. 多个可观察对象的组合

7.5.1. 示例-14:使用 [Observable.merge] 合并两个可观察对象

现在,我们将介绍 [Observable] 类的静态方法,这些方法可用于将多个可观察对象组合成一个结果可观察对象。

此类示例的首个示例如下:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.ProcessAction01;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple14 {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // 流程2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // 合并
        Process<?> process12 = new Process<>("process12",
                Observable.merge(process1.getObservable(), process2.getObservable()));
        // 订阅
        ProcessUtils.subscribe(1, process12);
    }
}
  • 第 15-17 行:一个名为 [process1] 的进程将在计算线程上输出 3 个实数。该进程也将在计算线程上被观察;
  • 第19-20行:名为[process2]的进程将在计算线程上输出2个字符串。观察线程未被指定。此前我们已看到,在此情况下,观察线程即为计算线程;
  • 第23行:两个进程被合并,即创建了一个其元素同时来自两个进程的可观察对象。为此,我们使用静态方法 [Observable.merge]:
 

与上图所示不同,在合并过程中,流 1 的元素可能会插入到流 2 的元素之间。执行结果显示了这一点:

main : début observation ------Thread[main] ---- Time[56:053]
main : attente fin observation ------Thread[main] ---- Time[56:073]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[56:073]
Observable (process2) call start ------Thread[RxComputationThreadPool-5] ---- Time[56:074]
Observable (process1,0) onNext (64.8) ------Thread[RxComputationThreadPool-4] ---- Time[56:263]
Observable (process2,0) onNext (valeur-0) ------Thread[RxComputationThreadPool-5] ---- Time[56:403]
Observable (process2,1) onNext (valeur-1) ------Thread[RxComputationThreadPool-5] ---- Time[56:515]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[56:516]
Subscriber[observateur[0],process12] : onNext (64.8) ------Thread[RxComputationThreadPool-3] ---- Time[56:552]
Subscriber[observateur[0],process12] : onNext ("valeur-0") ------Thread[RxComputationThreadPool-3] ---- Time[56:553]
Subscriber[observateur[0],process12] : onNext ("valeur-1") ------Thread[RxComputationThreadPool-3] ---- Time[56:553]
Observable (process1,1) onNext (56.4) ------Thread[RxComputationThreadPool-4] ---- Time[56:716]
Subscriber[observateur[0],process12] : onNext (56.4) ------Thread[RxComputationThreadPool-3] ---- Time[56:718]
Observable (process1,2) onNext (22.8) ------Thread[RxComputationThreadPool-4] ---- Time[57:082]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[57:083]
Subscriber[observateur[0],process12] : onNext (22.8) ------Thread[RxComputationThreadPool-3] ---- Time[57:084]
Subscriber[observateur[0],process12].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[57:085]
main : fin observation ------Thread[main] ---- Time[57:085]
  • 第 3 行:进程 [process1] 在计算线程 [RxComputationThreadPool-4] 上运行;
  • 第 4 行:进程 [process2] 在计算线程 [RxComputationThreadPool-5] 上运行;
  • 第 9 行:进程 [process12] 被观察到在计算线程 [RxComputationThreadPool-3] 上运行。我不清楚导致这一选择的规则;
  • 第 9-11 行:可以看到观察者同时观察了两个进程 [process1](第 5 行)和 [process2](第 6、7 行)的元素,而这两个进程均未结束(存在混淆);
  • 当两个进程 process1 process2 结束时,进程 [process12] 也会结束(第 17 行);

7.5.2. 示例-15:使用 [Observable.concat] 连接两个可观察对象

现在我们来看以下代码:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.ProcessAction01;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple15 {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // 流程2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, Schedulers.computation());
        // 拼接
        Process<?> process12 = new Process<>("process12",
                Observable.concat(process1.getObservable(), process2.getObservable()));
        // 订阅
        ProcessUtils.subscribe(1, process12);
    }
}
  • 第 15-17 行:一个名为 [process1] 的进程将在一个计算线程上输出 3 个实数。它也将在一个计算线程上被观察;
  • 第19-20行:名为[process2]的进程将在一个未指定线程(此处为默认主线程)上发布2个字符串。该进程将在计算线程上被观察;
  • 第 23 行:将这两个进程进行拼接,即创建一个其元素来自这两个进程的可观察对象。输出的值不会被混合。 进程 [process12] 将首先发布进程 [process1] 的所有值,然后发布进程 [process2] 的值。为此,我们使用静态方法 [Observable.concat]:
 

执行结果如下:

main : début observation ------Thread[main] ---- Time[30:162]
main : attente fin observation ------Thread[main] ---- Time[30:189]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[30:190]
Observable (process1,0) onNext (79.2) ------Thread[RxComputationThreadPool-4] ---- Time[30:681]
Observable (process1,1) onNext (98.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[30:792]
Subscriber[observateur[0],process12] : onNext (79.2) ------Thread[RxComputationThreadPool-3] ---- Time[30:975]
Subscriber[observateur[0],process12] : onNext (98.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[30:976]
Observable (process1,2) onNext (84.0) ------Thread[RxComputationThreadPool-4] ---- Time[31:084]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[31:085]
Subscriber[observateur[0],process12] : onNext (84.0) ------Thread[RxComputationThreadPool-3] ---- Time[31:086]
Observable (process2) call start ------Thread[RxComputationThreadPool-3] ---- Time[31:087]
Observable (process2,0) onNext (valeur-0) ------Thread[RxComputationThreadPool-3] ---- Time[31:556]
Subscriber[observateur[0],process12] : onNext ("valeur-0") ------Thread[RxComputationThreadPool-5] ---- Time[31:557]
Observable (process2,1) onNext (valeur-1) ------Thread[RxComputationThreadPool-3] ---- Time[31:608]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[31:609]
Subscriber[observateur[0],process12] : onNext ("valeur-1") ------Thread[RxComputationThreadPool-5] ---- Time[31:609]
Subscriber[observateur[0],process12].onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[31:610]
main : fin observation ------Thread[main] ---- Time[31:611]
  • 第3-10行:进程[process1]正在运行,且进程[process12]输出由[process1]生成的值;
  • 第 9 行:进程 [process1] 已结束;
  • 第11-17行:进程[process2]正在运行,进程[process12]输出由[process2]输出的值;

进程 process2 存在一个异常情况:我们并未指定其执行线程。因此,按理说它默认应为主线程。但实际情况并非如此。 实际的执行线程是计算线程 [RxComputationThreadPool-3](第 11 行)。因此,当未指定执行线程或观察线程时,无法预判系统将选择哪个线程。

7.5.3. 示例-16:使用 [Observable.zip] 组合两个可观测量

现在我们来看以下代码:


package dvp.rxjava.observables.exemples;

import java.util.Arrays;
import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.functions.FuncN;
import rx.schedulers.Schedulers;

public class Exemple16 {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, null);
        // 两个过程的组合函数
        FuncN<String> funcn = new FuncN<String>() {
            @Override
            public String call(Object... args) {
                if (args.length == 2) {
                    return String.format("double=%s, string=%s", args[0], args[1]);
                } else {
                    throw new RuntimeException("la fonction attend 2 paramètres exactement");
                }
            }
        };
        // 两个流程的压缩包
        Process<String> process12 = new Process<>("process12",
                Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));
        // 订阅
        ProcessUtils.subscribe(1, process12);
    }
}
  • 第16-18行:一个名为[process1]的进程将在一个计算线程上输出3个实数。它也将在一个计算线程上被观察;
  • 第20-21行:一个名为[process2]的进程将在一个非强制线程上输出2个字符串。观察线程同样为非强制线程;
  • 第 23-32 行:使用匿名类实例化类型 [FuncN<String>]。FuncN 是一个函数式接口:
 

方法 [FuncN.call] 接收一个对象数组,并返回类型 R。函数 [funcn] 将用于按此顺序组合流程 process1 process2。 在方法 [FuncN.call] 中:

  • args[0] 将是一个 Double
  • args[1] 将是 String

在此,[funcn.call] 的结果将是第 27 行中的字符串。构建此结果无需了解方法 call 的参数类型。

这两个过程的组合方式如下:


// 两个流程的压缩包
Process<String> process12 = new Process<>("process12",
Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));

方法 [Observable.zip] 的工作原理如下:

 

可以看出:

  • zip 的第一个参数是 Iterable<Observable>。在本例中,我们有一个由两个可观测量组成的有效参数,类型为 List<Observable>
  • zip 的第二个参数是类型 FuncN。在本例中,实际参数是 [funcn];

执行结果如下:

main : début observation ------Thread[main] ---- Time[55:636]
Observable (process2) call start ------Thread[main] ---- Time[55:666]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[55:666]
Observable (process1,0) onNext (69.6) ------Thread[RxComputationThreadPool-4] ---- Time[55:902]
Observable (process2,0) onNext (valeur-0) ------Thread[main] ---- Time[56:076]
Observable (process1,1) onNext (82.8) ------Thread[RxComputationThreadPool-4] ---- Time[56:271]
Subscriber[observateur[0],process12] : onNext ("double=69.6, string=valeur-0") ------Thread[main] ---- Time[56:352]
Observable (process1,2) onNext (14.399999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[56:641]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[56:642]
Observable (process2,1) onNext (valeur-1) ------Thread[main] ---- Time[56:778]
Subscriber[observateur[0],process12] : onNext ("double=82.8, string=valeur-1") ------Thread[main] ---- Time[56:779]
Observable (process2) onCompleted ------Thread[main] ---- Time[56:779]
Subscriber[observateur[0],process12].onCompleted ------Thread[main] ---- Time[56:780]
main : attente fin observation ------Thread[main] ---- Time[56:781]
main : fin observation ------Thread[main] ---- Time[56:781]
  • 第 7、11 行:进程 process12 发出两个元素;
  • 第 8 行:由进程 process1 发出的额外元素在进程 process2 中没有对应对象,因此未被结果进程 process12 发出;

可以看出,未被指定执行线程或观察线程的进程 process2 使用了主线程来处理这两项任务。

7.5.4. 示例-17:使用 [Observable.combineLatest] 组合两个可观察对象

现在我们来看以下代码:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple17 {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null,
                Schedulers.computation());
        // 两个流程的组合
        Process<Double> process12 = new Process<>("process12",
                Observable.combineLatest(process1.getObservable(), process2.getObservable(), (d1, d2) -> d1 + d2));
        // 订阅
        ProcessUtils.subscribe(1, process12);
    }
}
  • 第 14-16 行:一个名为 [process1] 的进程将在一个计算线程上发布 3 个实数。它也将在一个计算线程上被观察;
  • 第18-20行:一个名为[process2]的进程将在一个未指定线程上输出2个实数。这些值将在一个计算线程上被观测;
  • 第23行:这两个可观测量通过以下静态方法[Observable.combineLatest]进行组合:
 

可观测量 [combineLatest] 的工作原理如下:当其中一个可观测量发布 E1 元素时,该元素将由 [combineFunction] 与另一个可观测量发布的最新元素进行组合。

执行此代码将得到以下结果:

main : début observation ------Thread[main] ---- Time[01:768]
Observable (process2) call start ------Thread[main] ---- Time[01:791]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[01:791]
Observable (process1,0) onNext (54.0) ------Thread[RxComputationThreadPool-4] ---- Time[01:991]
Observable (process2,0) onNext (56.0) ------Thread[main] ---- Time[02:245]
Observable (process1,1) onNext (51.6) ------Thread[RxComputationThreadPool-4] ---- Time[02:358]
Subscriber[observateur[0],process12] : onNext (110.0) ------Thread[RxComputationThreadPool-5] ---- Time[02:521]
Subscriber[observateur[0],process12] : onNext (107.6) ------Thread[RxComputationThreadPool-5] ---- Time[02:522]
Observable (process2,1) onNext (261.8) ------Thread[main] ---- Time[02:595]
Observable (process2) onCompleted ------Thread[main] ---- Time[02:596]
main : attente fin observation ------Thread[main] ---- Time[02:596]
Subscriber[observateur[0],process12] : onNext (313.40000000000003) ------Thread[RxComputationThreadPool-5] ---- Time[02:597]
Observable (process1,2) onNext (80.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[02:790]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[02:791]
Subscriber[observateur[0],process12] : onNext (342.2) ------Thread[RxComputationThreadPool-3] ---- Time[02:792]
Subscriber[observateur[0],process12].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[02:792]
main : fin observation ------Thread[main] ---- Time[02:793]
  • 第5行:process2 (56) 的输出与process1 (54,第4行) 输出的最后一个元素组合,并产生第7行的结果;
  • 第6行:process1(51.6)的输出与process2(56,第5行)输出的最后一个元素组合,产生第8行的结果;
  • 第9行:process2(261.8)的输出与process1(51.6,第6行)输出的最后一个元素组合,产生第12行的结果;
  • 第13行:process1(80.39)的传输与process2(261.8,第9行)发出的最后一个元素组合,产生第15行的结果;

这里是可观测对象 [zip] 的一种变体,此次组合的元素不一定是流中位置相同的元素。 值得注意的是,此前未指定执行线程的进程 process2 此次是在主线程上执行的(第 2 行)。

7.5.5. 示例-18:使用 [Observable.amb] 组合两个可观察对象

现在我们来看以下代码:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple18 {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null, null);
        // 两个流程的组合
        Process<Double> process12 = new Process<>("process12",
                Observable.amb(process1.getObservable(), process2.getObservable()));
        // 订阅
        ProcessUtils.subscribe(1, process12);
    }
}
  • 第 14-16 行:一个名为 [process1] 的进程将在计算线程上输出 3 个实数。该进程也将在计算线程上被观察;
  • 第18-20行:一个名为[process2]的进程将在一个未约束的线程上输出2个实数。这些值将在一个未约束的线程上被观测;
  • 第22行:这两个可观测量通过以下静态方法[Observable.amb]进行组合:
 

如上图所示,可观察量 [Observable.amb(Observable o1, Observable o2)] 会发布第一个可观察量的元素。这得到了所展示示例结果的证实:

main : début observation ------Thread[main] ---- Time[21:594]
Observable (process2) call start ------Thread[main] ---- Time[21:612]
Observable (process1) call start ------Thread[RxComputationThreadPool-3] ---- Time[21:612]
Observable (process2,0) onNext (155.39999999999998) ------Thread[main] ---- Time[21:817]
Observable (process1) onError ------Thread[RxComputationThreadPool-3] ---- Time[21:820]
Observable (process1,0) onNext (90.0) ------Thread[RxComputationThreadPool-3] ---- Time[21:820]
Observable (process1,1) onNext (104.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[21:877]
Subscriber[observateur[0],process12] : onNext (155.39999999999998) ------Thread[main] ---- Time[22:105]
Observable (process1,2) onNext (44.4) ------Thread[RxComputationThreadPool-3] ---- Time[22:122]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[22:123]
Observable (process2,1) onNext (201.6) ------Thread[main] ---- Time[22:581]
Subscriber[observateur[0],process12] : onNext (201.6) ------Thread[main] ---- Time[22:583]
Observable (process2) onCompleted ------Thread[main] ---- Time[22:583]
Subscriber[observateur[0],process12].onCompleted ------Thread[main] ---- Time[22:584]
main : attente fin observation ------Thread[main] ---- Time[22:585]
main : fin observation ------Thread[main] ---- Time[22:586]
  • 第 4 行,首先发布的是进程 process2
  • 第 8、12 行:进程 process12 会发布进程 process2 发布的所有元素(第 4、11 行);

7.6. 可观察对象的处理链

7.6.1. 示例-19:使用 [Observable.map] 转换可观察对象

在之前的示例中,我们探讨了将两个可观察量组合成第三个可观察量的各种方式。 现在,我们将介绍 [Observable] 类的静态方法,这些方法支持对可观察量进行转换、过滤和聚合操作。这里将出现与第 5 节中研究的 [Stream] 类方法类似的方法。

我们的第一个示例如下:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple19 {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<String> process2 = new Process<>("process2",
                process1.getObservable().map(d -> String.format("valeur-%s", d)));
        // 订阅
        ProcessUtils.subscribe(1, process2);
    }
}
  • 第14-16行:一个名为process1的进程将在计算线程上输出3个实数。该进程同样将在计算线程上被观察;
  • 第17-18行:由process1输出的数值将在process2进程中转换为字符串;
  • 第 20 行:观察 process2

第18行的[Observable.map]方法与第5.5节中研究的[Stream.map]方法类似:

 

该示例的结果如下:

main : début observation ------Thread[main] ---- Time[55:328]
main : attente fin observation ------Thread[main] ---- Time[55:346]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[55:347]
Observable (process1,0) onNext (21.599999999999998) ------Thread[RxComputationThreadPool-4] ---- Time[55:354]
Observable (process1,1) onNext (97.2) ------Thread[RxComputationThreadPool-4] ---- Time[55:512]
Subscriber[observateur[0],process2] : onNext ("valeur-21.599999999999998") ------Thread[RxComputationThreadPool-3] ---- Time[55:615]
Subscriber[observateur[0],process2] : onNext ("valeur-97.2") ------Thread[RxComputationThreadPool-3] ---- Time[55:616]
Observable (process1,2) onNext (98.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[55:803]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[55:804]
Subscriber[observateur[0],process2] : onNext ("valeur-98.39999999999999") ------Thread[RxComputationThreadPool-3] ---- Time[55:804]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[55:805]
main : fin observation ------Thread[main] ---- Time[55:805]
  • 第4、5和8行:process1的输出结果。这些是实数;
  • 第6、7、10行:process2的观测值。这些是字符串;

7.6.2. 示例-20:使用 [Observable.filter] 过滤观测值

示例如下:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple20 {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().filter(i -> i % 2 == 0));
        // 订阅
        ProcessUtils.subscribe(1, process2);
    }
}
  • 第11-12行:一个名为process1的进程将在计算线程上输出0到2的整数。该进程也将在计算线程上被观察;
  • 第14行:process1输出的数字将被过滤,仅将偶数保留在process2中
  • 第20行:观察process2

第18行的方法[Observable.filter]与第5.4节中研究的方法[Stream.filter]类似:

 

该示例的结果如下:

main : début observation ------Thread[main] ---- Time[30:319]
main : attente fin observation ------Thread[main] ---- Time[30:335]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[30:336]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[30:388]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[30:625]
Subscriber[observateur[0],process2] : onNext (0) ------Thread[RxComputationThreadPool-3] ---- Time[30:703]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[30:704]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[30:705]
Subscriber[observateur[0],process2] : onNext (2) ------Thread[RxComputationThreadPool-3] ---- Time[30:706]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[30:707]
main : fin observation ------Thread[main] ---- Time[30:707]
  • 第4、5和7行:process1的观测数据;
  • 第6、9行:观测到的process2的发射;这些是process1中为偶数的元素;

7.6.3. 示例-21:使用 [Observable.flatMap] 转换观测值

示例如下:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple21 {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
            int value = i * 10;
            return Observable.just(value, value + 1, value + 2);
        }));
        // 订阅
        ProcessUtils.subscribe(1, process2);
    }
}
  • 第12-13行:一个名为process1的进程将在一个计算线程上发出0到2的整数。它也将在一个计算线程上被观察;
  • 第15-18行:process1发出的每个数字n都会被转换为一个可观察对象,该对象会发出3个数字(10*n, 10*n+1, 10*n+2)。 如果第 15 行使用的是方法 [map],那么 process2 发出的将是一个 Observable<Integer> 类型,而不是 Integer 类型。 所使用的 [flatMap] 方法可将 (flatten) 将该 Observable<Integer> 类型的元素序列扁平化为 Integer 类型的元素序列,该序列由每个 Observable<Integer> 中的各个元素组成;
  • 第20行:观察到process2

第15行的[Observable.flatMap]方法与第5.6.12节中讨论的[Stream.flatMap]方法类似:

 

该示例的结果如下:

main : début observation ------Thread[main] ---- Time[31:466]
main : attente fin observation ------Thread[main] ---- Time[31:486]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[31:486]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[31:777]
Subscriber[observateur[0],process2] : onNext (0) ------Thread[RxComputationThreadPool-3] ---- Time[32:082]
Subscriber[observateur[0],process2] : onNext (1) ------Thread[RxComputationThreadPool-3] ---- Time[32:085]
Subscriber[observateur[0],process2] : onNext (2) ------Thread[RxComputationThreadPool-3] ---- Time[32:087]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[32:192]
Subscriber[observateur[0],process2] : onNext (10) ------Thread[RxComputationThreadPool-3] ---- Time[32:194]
Subscriber[observateur[0],process2] : onNext (11) ------Thread[RxComputationThreadPool-3] ---- Time[32:196]
Subscriber[observateur[0],process2] : onNext (12) ------Thread[RxComputationThreadPool-3] ---- Time[32:197]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[32:686]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[32:687]
Subscriber[observateur[0],process2] : onNext (20) ------Thread[RxComputationThreadPool-3] ---- Time[32:688]
Subscriber[observateur[0],process2] : onNext (21) ------Thread[RxComputationThreadPool-3] ---- Time[32:690]
Subscriber[observateur[0],process2] : onNext (22) ------Thread[RxComputationThreadPool-3] ---- Time[32:692]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[32:693]
main : fin observation ------Thread[main] ---- Time[32:693]
  • 第5-7行:process1第4行发送后,process2的三个发送;
  • 第9-11行:process2的三个输出,源于process1的第8行输出;
  • 第14-16行:process2的三次发送,紧随process1的第12行发送之后;

以下代码演示了如何基于 process1 和 [Exemple21b] 创建类型 Observable<Integer[]>


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21b {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Integer[]> process2 = new Process<>("process2", process1.getObservable().map(i -> {
            int value = i * 10;
            return new Integer[] { value, value + 1, value + 2 };
        }));
        // 订阅
        ProcessUtils.subscribe(1, process2);
    }
}
  • 第 14 行:使用 [Observable.map] 方法;
  • 第 16 行:该方法返回类型 Integer[]

结果如下:

main : début observation ------Thread[main] ---- Time[58:089]
main : attente fin observation ------Thread[main] ---- Time[58:107]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[58:108]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[58:503]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[58:762]
Subscriber[observateur[0],process2] : onNext ([0,1,2]) ------Thread[RxComputationThreadPool-3] ---- Time[58:792]
Subscriber[observateur[0],process2] : onNext ([10,11,12]) ------Thread[RxComputationThreadPool-3] ---- Time[58:795]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[58:851]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[58:852]
Subscriber[observateur[0],process2] : onNext ([20,21,22]) ------Thread[RxComputationThreadPool-3] ---- Time[58:853]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[58:854]
main : fin observation ------Thread[main] ---- Time[58:854]
  • 第 6、7、10 行:可见 map 的结果;

所有这些可观测量的转换都可以串联起来,因为每次转换都会产生一个新的可观测量。下例 [Exemple21c] 展示了这一点:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple21c {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
            int value = i * 10;
            return Observable.just(value, value + 1, value + 2);
        }).filter(i -> i % 2 == 0));
        // 订阅
        ProcessUtils.subscribe(1, process2);
    }
}
  • 第15-18行:flatMap之后紧跟一个filter

执行结果如下:

main : début observation ------Thread[main] ---- Time[37:993]
main : attente fin observation ------Thread[main] ---- Time[38:016]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[38:017]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[38:124]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[38:366]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[38:380]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[38:381]
Subscriber[observateur[0],process2] : onNext (0) ------Thread[RxComputationThreadPool-3] ---- Time[38:436]
Subscriber[observateur[0],process2] : onNext (2) ------Thread[RxComputationThreadPool-3] ---- Time[38:439]
Subscriber[observateur[0],process2] : onNext (10) ------Thread[RxComputationThreadPool-3] ---- Time[38:441]
Subscriber[observateur[0],process2] : onNext (12) ------Thread[RxComputationThreadPool-3] ---- Time[38:443]
Subscriber[observateur[0],process2] : onNext (20) ------Thread[RxComputationThreadPool-3] ---- Time[38:445]
Subscriber[observateur[0],process2] : onNext (22) ------Thread[RxComputationThreadPool-3] ---- Time[38:446]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[38:447]
main : fin observation ------Thread[main] ---- Time[38:447]
  • 第8-13行:process2仅输出flatMap中的偶数项;

与 [flatMap] 类似的方法是 [flatMapIterable],如下例所示(以 [Exemple21d] 为例):


package dvp.rxjava.observables.exemples;

import java.util.Arrays;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21d {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMapIterable(i -> {
            int value = i * 10;
            return Arrays.asList(value, value + 1, value + 2);
        }).filter(i -> i % 2 == 0));
        // 订阅
        ProcessUtils.subscribe(1, process2);
    }
}

第 16 行,这里使用的是 [flatMapIterable] 方法,而非 [flatMap] 方法。 在这种情况下,转换函数应生成类型为 Iterable<T>(第 18 行)的结果,而不是类型为 Observable<T> 的结果。

结果与之前相同。

让我们回到方法 [flatMap] 的定义:

 

如上所示,一个蓝色元素 [3] 插入到了两个绿色元素 [1-2] 之间。 这意味着在对 Observable<T> 进行扁平化操作时,方法 [flatMap] 遵循了这些不同内部可观量的发射顺序。下例 [Exemple21e] 展示了这一点:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21e {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // 流程 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().flatMap(i -> process2.getObservable()));
        // 订阅
        ProcessUtils.subscribe(1, process3);
    }
}
  • 第 11-12 行:进程 process1 生成整数 [0,1];
  • 第14-15行:进程process2输出整数[10,11,12];
  • 第 17-18 行:process1 发出的每个元素,都与进程 process2 的可观测量相关联。这意味着:
    • process1 的元素 [0] 将关联一个输出 [10,11,12] 的可观测量;
    • 元素1亦同;

最终,将发出6个数字[10, 11, 12, 10, 11, 12]。我们想查看它们的输出顺序。

执行结果如下:

main : début observation ------Thread[main] ---- Time[22:540]
main : attente fin observation ------Thread[main] ---- Time[22:566]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[22:566]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[22:949]
Observable (process2) call start ------Thread[RxComputationThreadPool-6] ---- Time[22:951]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[23:159]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[23:160]
Observable (process2) call start ------Thread[RxComputationThreadPool-8] ---- Time[23:160]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-6] ---- Time[23:286]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-8] ---- Time[23:513]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[23:597]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[23:599]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-6] ---- Time[23:645]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-5] ---- Time[23:647]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-6] ---- Time[23:789]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[23:790]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-5] ---- Time[23:791]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-8] ---- Time[23:976]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-7] ---- Time[23:978]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-8] ---- Time[24:184]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-8] ---- Time[24:184]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-7] ---- Time[24:186]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-7] ---- Time[24:187]
main : fin observation ------Thread[main] ---- Time[24:187]

可见,进程 process3 的发送顺序为:[10, 10, 11, 12, 11, 12](第 11、12、14、17、19、22 行)。 因此,process2 进程输出的元素确实发生了混淆。 可以通过使用方法 [concatMap] 代替方法 [flatMap] 来避免这种情况。以下代码 [Exemple21ef] 展示了这一点:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21ef {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // 流程 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().concatMap(i -> process2.getObservable()));
        // 订阅
        ProcessUtils.subscribe(1, process3);
    }
}

在第 18 行,我们将 [flatMap] 替换为 [concatMap]。执行结果如下:

main : début observation ------Thread[main] ---- Time[45:507]
main : attente fin observation ------Thread[main] ---- Time[45:530]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[45:530]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[45:775]
Observable (process2) call start ------Thread[RxComputationThreadPool-6] ---- Time[45:778]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-6] ---- Time[45:846]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-6] ---- Time[45:890]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[45:947]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[45:948]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-6] ---- Time[46:096]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[46:097]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[46:144]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-5] ---- Time[46:147]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-5] ---- Time[46:148]
Observable (process2) call start ------Thread[RxComputationThreadPool-8] ---- Time[46:149]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-8] ---- Time[46:364]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-7] ---- Time[46:366]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-8] ---- Time[46:529]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-7] ---- Time[46:531]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-8] ---- Time[46:558]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-8] ---- Time[46:559]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-7] ---- Time[46:560]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-7] ---- Time[46:562]
main : fin observation ------Thread[main] ---- Time[46:562]

可以看出,进程 process3 的输出顺序为:[10, 11, 12, 10, 11, 12](第 12-14 行、第 17 行、第 19 行、第 22 行)。 进程 process2 输出的元素未被混入。

[map]方法的另一种变体是[switchMap]方法:

 

在上文中,从可观测对象 [1] 衍生出另外 3 个包含 2 个元素的可观测对象 [2],随后这些对象被扁平化,如 [flatMap] 和 [3] 所示。 可以注意到,结果包含5个元素而非6个。这是因为在第二个可观测对象发出其第2个元素[6]之前,第三个可观测对象已先发出了其第一个元素[5],导致第二个可观测对象被舍弃。 因此,在结果可观测量 [3] 中找不到元素 [6]。

为说明 [switchMap],我们将使用以下示例 [Exemple21eg]:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21eg {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // 流程 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().switchMap(i -> process2.getObservable()));
        // 订阅
        ProcessUtils.subscribe(1, process3);
    }
}

运行该示例将得到以下结果:

main : début observation ------Thread[main] ---- Time[02:388]
main : attente fin observation ------Thread[main] ---- Time[02:419]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[02:419]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[02:641]
Observable (process2) call start ------Thread[RxComputationThreadPool-6] ---- Time[02:643]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-6] ---- Time[02:802]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-6] ---- Time[02:888]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-6] ---- Time[02:957]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[02:958]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[03:005]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[03:007]
Observable (process2) call start ------Thread[RxComputationThreadPool-8] ---- Time[03:007]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-8] ---- Time[03:106]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[03:106]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[03:108]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-8] ---- Time[03:236]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-7] ---- Time[03:238]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-8] ---- Time[03:716]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-8] ---- Time[03:717]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-7] ---- Time[03:718]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-7] ---- Time[03:718]
main : fin observation ------Thread[main] ---- Time[03:719]
  • process1 发出 2 个元素,从而产生 2 个各含 3 个元素的可观测对象 process2
  • 第14行:观察者接收第6行首个可观测量process2发出的第0个元素;
  • 第15行:观察者接收了第2个可观量process2(第13行)发出的第0号元素。 尚不清楚为何观察者此前未收到由第一个可观测对象 process2 在第 7 行和第 8 行发出的第 1 和第 2 个数据项。无论如何,第一个可观测对象 process2 已被放弃;
  • 最终,观测者仅看到4个数据项(第14、15、17、20行),而非实际发送的6个;

7.6.4. 示例-22:[Observable]类的其他方法

类 [Observable] 继承了类 [Stream] 的许多方法,其工作原理类似。以下列举其中部分方法。我们仅给出代码及其结果。

[Exemple22a - take=limit]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22a {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).take(3));
        // 订阅
        ProcessUtils.subscribe(1, process);
    }
}

结果

1
2
3
4
5
6
7
main : début observation ------Thread[main] ---- Time[25:071]
Subscriber[observateur[0],process] : onNext (1) ------Thread[main] ---- Time[25:399]
Subscriber[observateur[0],process] : onNext (2) ------Thread[main] ---- Time[25:402]
Subscriber[observateur[0],process] : onNext (3) ------Thread[main] ---- Time[25:404]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[25:404]
main : attente fin observation ------Thread[main] ---- Time[25:406]
main : fin observation ------Thread[main] ---- Time[25:406]

[Exemple22b - takeLast]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22b {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).takeLast(2));
        // 订阅
        ProcessUtils.subscribe(1, process);
    }
}

结果

1
2
3
4
5
6
main : début observation ------Thread[main] ---- Time[19:440]
Subscriber[observateur[0],process] : onNext (9) ------Thread[main] ---- Time[19:726]
Subscriber[observateur[0],process] : onNext (10) ------Thread[main] ---- Time[19:728]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[19:728]
main : attente fin observation ------Thread[main] ---- Time[19:729]
main : fin observation ------Thread[main] ---- Time[19:730]

[Exemple22c - skip]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22c {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).skip(5).take(2));
        // 订阅
        ProcessUtils.subscribe(1, process);
    }
}

结果

1
2
3
4
5
6
main : début observation ------Thread[main] ---- Time[16:685]
Subscriber[observateur[0],process] : onNext (6) ------Thread[main] ---- Time[17:002]
Subscriber[observateur[0],process] : onNext (7) ------Thread[main] ---- Time[17:004]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[17:005]
main : attente fin observation ------Thread[main] ---- Time[17:006]
main : fin observation ------Thread[main] ---- Time[17:006]

[Exemple22d - reduce]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22d {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).reduce(0, (i, a) -> i + a));
        // 订阅
        ProcessUtils.subscribe(1, process);
    }
}
  • 第 10 行:计算可观测量的元素之和。结果是一个会发出该和值的可观测量;

结果

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[52:412]
Subscriber[observateur[0],process] : onNext (55) ------Thread[main] ---- Time[52:640]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[52:640]
main : attente fin observation ------Thread[main] ---- Time[52:642]
main : fin observation ------Thread[main] ---- Time[52:642]

[Exemple22e - all]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22e {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        Process<Boolean> process = new Process<>("process", Observable.range(1, 10).all(i -> i > 10));
        // 订阅
        ProcessUtils.subscribe(1, process);
    }
}
  • 第 10 行:返回一个 Observable<Boolean>,该可观察对象会发布元素 true;如果方法 [all] 的谓词对所有元素均为真,则返回 false;否则返回 Observable<Boolean>;

结果

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[59:866]
Subscriber[observateur[0],process] : onNext (false) ------Thread[main] ---- Time[00:069]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[00:070]
main : attente fin observation ------Thread[main] ---- Time[00:071]
main : fin observation ------Thread[main] ---- Time[00:071]

[Exemple22f - count]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22f {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).count());
        // 订阅
        ProcessUtils.subscribe(1, process);
    }
}
  • 第 10 行:[Observable.count] 创建了一个包含 1 个元素的可观察对象,该元素是所有被观察元素的总和;

结果

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[16:409]
Subscriber[observateur[0],process] : onNext (10) ------Thread[main] ---- Time[16:634]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[16:634]
main : attente fin observation ------Thread[main] ---- Time[16:635]
main : fin observation ------Thread[main] ---- Time[16:635]

[Exemple22g - distinct]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22g {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        Process<Integer> process = new Process<>("process", Observable.just(1, 2, 1, 3).distinct());
        // 订阅
        ProcessUtils.subscribe(1, process);
    }
}

结果

1
2
3
4
5
6
7
main : début observation ------Thread[main] ---- Time[05:373]
Subscriber[observateur[0],process] : onNext (1) ------Thread[main] ---- Time[05:594]
Subscriber[observateur[0],process] : onNext (2) ------Thread[main] ---- Time[05:595]
Subscriber[observateur[0],process] : onNext (3) ------Thread[main] ---- Time[05:596]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[05:597]
main : attente fin observation ------Thread[main] ---- Time[05:597]
main : fin observation ------Thread[main] ---- Time[05:597]

[Exemple22h - groupBy, asObservable]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.observables.GroupedObservable;

public class Exemple22h {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        Observable<GroupedObservable<Boolean, Integer>> obs = Observable.range(1, 10).groupBy(i -> i % 2 == 0);
        Process<Integer> process = new Process<>("process", obs.concatMap(g -> g.asObservable()));
        // 订阅
        ProcessUtils.subscribe(1, process);
    }
}
  • 第 11 行:方法 [groupBy] 将输出的 10 个元素分为两组,即偶数和奇数。 结果类型为 Observable<GroupedObservable<Boolean, Integer>>,即一个可观察对象,其元素类型为 GroupedObservable<Boolean, Integer>,其中 Boolean 是组键的类型 (此处为 falsetrue),同时也是作为参数传递给 [groupBy] 方法的 lambda 表达式返回值的类型,而 Integer 则是该组元素的类型;
  • 第 12 行:类型 GroupedObservable 有一个方法 [asObservable],可用于基于该类型创建可观察对象。 因此我们将得到两个类型 Observable<Integer>,一个用于偶数,另一个用于奇数。对于这两个可观察对象,方法 [concatMap] 将创建一个单一的可观察对象;

结果

main : début observation ------Thread[main] ---- Time[23:809]
Subscriber[observateur[0],process] : onNext (1) ------Thread[main] ---- Time[24:034]
Subscriber[observateur[0],process] : onNext (3) ------Thread[main] ---- Time[24:036]
Subscriber[observateur[0],process] : onNext (5) ------Thread[main] ---- Time[24:037]
Subscriber[observateur[0],process] : onNext (7) ------Thread[main] ---- Time[24:038]
Subscriber[observateur[0],process] : onNext (9) ------Thread[main] ---- Time[24:039]
Subscriber[observateur[0],process] : onNext (2) ------Thread[main] ---- Time[24:041]
Subscriber[observateur[0],process] : onNext (4) ------Thread[main] ---- Time[24:043]
Subscriber[observateur[0],process] : onNext (6) ------Thread[main] ---- Time[24:044]
Subscriber[observateur[0],process] : onNext (8) ------Thread[main] ---- Time[24:045]
Subscriber[observateur[0],process] : onNext (10) ------Thread[main] ---- Time[24:046]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[24:047]
main : attente fin observation ------Thread[main] ---- Time[24:047]
main : fin observation ------Thread[main] ---- Time[24:048]

[Exemple22i - timestamp]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;
import rx.schedulers.Timestamped;

public class Exemple22i {
    public static void main(String[] args) throws InterruptedException {
        // 流程 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // 流程 2
        Process<Timestamped<Integer>> process2 = new Process<>("process2", process1.getObservable().timestamp());
        // 订阅
        ProcessUtils.subscribe(1, process2);
    }
}
  • 第 15 行,方法 [timestamp] 为处理后的每个可观测量元素关联一个时间;

结果

main : début observation ------Thread[main] ---- Time[59:362]
main : attente fin observation ------Thread[main] ---- Time[59:377]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[59:378]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[59:553]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[59:692]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462975259555,"value":0}) ------Thread[RxComputationThreadPool-3] ---- Time[59:789]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462975259789,"value":1}) ------Thread[RxComputationThreadPool-3] ---- Time[59:791]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[00:025]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[00:027]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462975260026,"value":2}) ------Thread[RxComputationThreadPool-3] ---- Time[00:031]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[00:033]
main : fin observation ------Thread[main] ---- Time[00:034]

在此示例中,很难判断信息 timestamp 代表什么:

  • 第4-5行:可见process1的第1个元素是在第0个元素之后139毫秒发出的;
  • 第6、7行:可见process2的第1个元素是在第0个元素之后234毫秒被观测到的;
  • 第5、8行:可见process1的第2个元素是在第1个元素发出33毫秒后发出的;
  • 第7行和第10行:可见process2的第2个元素是在第1个元素之后37毫秒被观测到的;

这些时间差是由于观测线程与可观测量的执行线程不同所致。如果将第12-13行替换为以下内容(示例22j):


// 流程 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
null);
  • 第2-3行:未指定观察线程。我们知道,在此情况下,可观测量将在其执行位置被观察;

这将产生以下结果:

main : début observation ------Thread[main] ---- Time[43:834]
main : attente fin observation ------Thread[main] ---- Time[43:845]
Observable (process1) call start ------Thread[RxComputationThreadPool-1] ---- Time[43:846]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-1] ---- Time[44:291]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462976384293,"value":0}) ------Thread[RxComputationThreadPool-1] ---- Time[44:552]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-1] ---- Time[44:878]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462976384879,"value":1}) ------Thread[RxComputationThreadPool-1] ---- Time[44:884]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-1] ---- Time[45:274]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462976385275,"value":2}) ------Thread[RxComputationThreadPool-1] ---- Time[45:280]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[45:281]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[45:283]
main : fin observation ------Thread[main] ---- Time[45:284]
  • 第4行和第6行:进程process1在发出第0个元素587毫秒后发出第1个元素;
  • 第5行和第7行:观察者观测到这两个元素的时间间隔为586毫秒;
  • 第 6 行和第 8 行:进程 process1 在发出第 1 个元素 396 毫秒后发出第 2 个元素;
  • 第7行和第9行:观察者观测到这两个元素的时间间隔为396毫秒;

在此,timestamp的数值是连贯的:它们确实代表了元素的发送时间。

7.7. 调度器

7.7.1. 示例-23:调度器 [Schedulers.computation]

现在我们来探讨执行调度器。观察将基于执行线程进行。

调度器这一主题稍显晦涩。StackOverflow 网站上的这篇文章介绍了各种调度器:

 

我们将通过实例来演示这些不同调度器的用法。第一个示例演示了 [Schedulers.computation] 调度器:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple23 {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        @SuppressWarnings("unchecked")
        Process<Double> processes[] = new Process[10];
        for (int i = 0; i < processes.length; i++) {
            processes[i] = new Process<>(
                    new ProcessAction01<Double>(String.format("process%s", i), 1, value -> new Random().nextInt(100) * 1.2),
                    Schedulers.computation(), null);
        }
        // 订阅
        ProcessUtils.subscribe(1, processes);
    }
}
  • 第14-19行:创建一个包含10个进程的数组,这些进程在计算线程上运行;
  • 第17行:每个进程生成一个随机实数;
  • 第21行:订阅所有这些进程;

结果如下:

main : début observation ------Thread[main] ---- Time[01:034]
Observable (process0) call start ------Thread[RxComputationThreadPool-1] ---- Time[01:042]
Observable (process2) call start ------Thread[RxComputationThreadPool-3] ---- Time[01:042]
Observable (process1) call start ------Thread[RxComputationThreadPool-2] ---- Time[01:042]
Observable (process5) call start ------Thread[RxComputationThreadPool-6] ---- Time[01:043]
Observable (process7) call start ------Thread[RxComputationThreadPool-8] ---- Time[01:043]
Observable (process4) call start ------Thread[RxComputationThreadPool-5] ---- Time[01:042]
Observable (process3) call start ------Thread[RxComputationThreadPool-4] ---- Time[01:042]
main : attente fin observation ------Thread[main] ---- Time[01:043]
Observable (process6) call start ------Thread[RxComputationThreadPool-7] ---- Time[01:043]
Observable (process3,0) onNext (70.8) ------Thread[RxComputationThreadPool-4] ---- Time[01:115]
Observable (process1,0) onNext (13.2) ------Thread[RxComputationThreadPool-2] ---- Time[01:153]
Observable (process0,0) onNext (63.599999999999994) ------Thread[RxComputationThreadPool-1] ---- Time[01:215]
Subscriber[observateur[0],process0] : onNext (63.599999999999994) ------Thread[RxComputationThreadPool-1] ---- Time[01:326]
Subscriber[observateur[0],process3] : onNext (70.8) ------Thread[RxComputationThreadPool-4] ---- Time[01:326]
Subscriber[observateur[0],process1] : onNext (13.2) ------Thread[RxComputationThreadPool-2] ---- Time[01:326]
Observable (process3) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[01:326]
Observable (process0) onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[01:326]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[01:327]
Subscriber[observateur[0],process0].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[01:327]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[01:327]
Subscriber[observateur[0],process1].onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[01:327]
Observable (process8) call start ------Thread[RxComputationThreadPool-1] ---- Time[01:329]
Observable (process9) call start ------Thread[RxComputationThreadPool-2] ---- Time[01:329]
...
main : fin observation ------Thread[main] ---- Time[01:610]
  • 第2-10行:前8个进程在8个不同的线程上启动(所用机器有8个核心)。可以注意到它们几乎在同一时间启动;
  • 第17-19行:3个进程结束,从而释放了3个线程;
  • 第23-24行:最后两个进程利用这2个被释放的线程启动;

因此,我们可以总结出:调度器 [Schedulers.computation] 提供了一个包含 n 个线程的线程池,其中 n 代表机器的核心数。这些线程在这些核心上并行执行。

7.7.2. 示例-24:调度器 [Schedulers.io]

我们使用调度器 [Schedulers.io] 执行上述代码:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple24 {
    public static void main(String[] args) throws InterruptedException {
        // 进程
        @SuppressWarnings("unchecked")
        Process<Double> processes[] = new Process[10];
        for (int i = 0; i < processes.length; i++) {
            processes[i] = new Process<>(
                    new ProcessAction01<Double>(String.format("process%s", i), 1, value -> new Random().nextInt(100) * 1.2),
                    Schedulers.io(), null);
        }
        // 订阅
        ProcessUtils.subscribe(1, processes);
    }
}
  • 第 18 行:进程使用调度器 [Schedulers.io] 的线程运行;

这产生了以下结果:

main : début observation ------Thread[main] ---- Time[03:451]
Observable (process0) call start ------Thread[RxCachedThreadScheduler-1] ---- Time[03:459]
Observable (process1) call start ------Thread[RxCachedThreadScheduler-2] ---- Time[03:459]
Observable (process2) call start ------Thread[RxCachedThreadScheduler-3] ---- Time[03:460]
Observable (process3) call start ------Thread[RxCachedThreadScheduler-4] ---- Time[03:460]
Observable (process4) call start ------Thread[RxCachedThreadScheduler-5] ---- Time[03:464]
Observable (process5) call start ------Thread[RxCachedThreadScheduler-6] ---- Time[03:464]
Observable (process6) call start ------Thread[RxCachedThreadScheduler-7] ---- Time[03:465]
Observable (process8) call start ------Thread[RxCachedThreadScheduler-9] ---- Time[03:465]
Observable (process9) call start ------Thread[RxCachedThreadScheduler-10] ---- Time[03:465]
main : attente fin observation ------Thread[main] ---- Time[03:465]
Observable (process7) call start ------Thread[RxCachedThreadScheduler-8] ---- Time[03:465]
Observable (process7,0) onNext (54.0) ------Thread[RxCachedThreadScheduler-8] ---- Time[03:473]
Observable (process8,0) onNext (116.39999999999999) ------Thread[RxCachedThreadScheduler-9] ---- Time[03:500]
Observable (process6,0) onNext (105.6) ------Thread[RxCachedThreadScheduler-7] ---- Time[03:506]
Observable (process0,0) onNext (96.0) ------Thread[RxCachedThreadScheduler-1] ---- Time[03:509]
Observable (process5,0) onNext (25.2) ------Thread[RxCachedThreadScheduler-6] ---- Time[03:583]
Observable (process3,0) onNext (97.2) ------Thread[RxCachedThreadScheduler-4] ---- Time[03:684]
Subscriber[observateur[0],process7] : onNext (54.0) ------Thread[RxCachedThreadScheduler-8] ---- Time[03:685]
Subscriber[observateur[0],process6] : onNext (105.6) ------Thread[RxCachedThreadScheduler-7] ---- Time[03:685]
Subscriber[observateur[0],process0] : onNext (96.0) ------Thread[RxCachedThreadScheduler-1] ---- Time[03:685]
Subscriber[observateur[0],process8] : onNext (116.39999999999999) ------Thread[RxCachedThreadScheduler-9] ---- Time[03:685]
Observable (process0) onCompleted ------Thread[RxCachedThreadScheduler-1] ---- Time[03:686]
Observable (process6) onCompleted ------Thread[RxCachedThreadScheduler-7] ---- Time[03:686]
Observable (process7) onCompleted ------Thread[RxCachedThreadScheduler-8] ---- Time[03:685]
...
main : fin observation ------Thread[main] ---- Time[03:933]
  • 第2-10行:这10个进程各自在不同的线程上启动。与前一种情况不同,所有进程均已成功启动。值得注意的是,这些进程的启动耗时为6毫秒,而之前仅需1毫秒;
  • 第13-18行:可观察对象依次发出信号,而非像之前那样近乎并行地发出;

调度器 [Schedulers.io] 和 [Schedulers.computation] 之间有何区别?答案可参见 URL [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases]:

 

7.7.3. 示例-25:调度程序 [Schedulers.newThread]

我们使用调度程序 [Schedulers.newThread] 执行上述代码:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple25 {
    public static void main(String[] args) throws InterruptedException {
        // 流程
        @SuppressWarnings("unchecked")
        Process<Double> processes[] = new Process[10];
        for (int i = 0; i < processes.length; i++) {
            processes[i] = new Process<>(
                    new ProcessAction01<Double>(String.format("process%s", i), 1, value -> new Random().nextInt(100) * 1.2),
                    Schedulers.newThread(), null);
        }
        // 订阅
        ProcessUtils.subscribe(1, processes);
    }
}

所得结果与调度程序 [Schedulers.io] 相同:

main : début observation ------Thread[main] ---- Time[17:058]
Observable (process0) call start ------Thread[RxNewThreadScheduler-1] ---- Time[17:065]
Observable (process1) call start ------Thread[RxNewThreadScheduler-2] ---- Time[17:065]
Observable (process2) call start ------Thread[RxNewThreadScheduler-3] ---- Time[17:066]
Observable (process3) call start ------Thread[RxNewThreadScheduler-4] ---- Time[17:066]
Observable (process4) call start ------Thread[RxNewThreadScheduler-5] ---- Time[17:068]
Observable (process5) call start ------Thread[RxNewThreadScheduler-6] ---- Time[17:069]
Observable (process6) call start ------Thread[RxNewThreadScheduler-7] ---- Time[17:069]
Observable (process8) call start ------Thread[RxNewThreadScheduler-9] ---- Time[17:069]
Observable (process7) call start ------Thread[RxNewThreadScheduler-8] ---- Time[17:069]
Observable (process9) call start ------Thread[RxNewThreadScheduler-10] ---- Time[17:069]
main : attente fin observation ------Thread[main] ---- Time[17:069]
Observable (process6,0) onNext (25.2) ------Thread[RxNewThreadScheduler-7] ---- Time[17:120]
Observable (process3,0) onNext (39.6) ------Thread[RxNewThreadScheduler-4] ---- Time[17:193]
Observable (process5,0) onNext (21.599999999999998) ------Thread[RxNewThreadScheduler-6] ---- Time[17:212]
Observable (process0,0) onNext (19.2) ------Thread[RxNewThreadScheduler-1] ---- Time[17:273]
Observable (process8,0) onNext (81.6) ------Thread[RxNewThreadScheduler-9] ---- Time[17:308]
Subscriber[observateur[0],process3] : onNext (39.6) ------Thread[RxNewThreadScheduler-4] ---- Time[17:331]
Subscriber[observateur[0],process0] : onNext (19.2) ------Thread[RxNewThreadScheduler-1] ---- Time[17:331]
Subscriber[observateur[0],process6] : onNext (25.2) ------Thread[RxNewThreadScheduler-7] ---- Time[17:331]
Subscriber[observateur[0],process8] : onNext (81.6) ------Thread[RxNewThreadScheduler-9] ---- Time[17:331]
Subscriber[observateur[0],process5] : onNext (21.599999999999998) ------Thread[RxNewThreadScheduler-6] ---- Time[17:331]
Observable (process8) onCompleted ------Thread[RxNewThreadScheduler-9] ---- Time[17:333]
Observable (process5) onCompleted ------Thread[RxNewThreadScheduler-6] ---- Time[17:333]
Observable (process6) onCompleted ------Thread[RxNewThreadScheduler-7] ---- Time[17:332]
Observable (process0) onCompleted ------Thread[RxNewThreadScheduler-1] ---- Time[17:332]
Observable (process3) onCompleted ------Thread[RxNewThreadScheduler-4] ---- Time[17:332]
...
main : fin observation ------Thread[main] ---- Time[17:571]

在 URL 和 [http://stackoverflow.com/questions/33415881/retrofit-with-rxjava-schedulers-newthread-vs-schedulers-io] 中,我们解释了调度程序 [Schedulers.io] 提供了一个线程池,而调度程序 [Schedulers.newThread] 则不提供。 线程池会自动创建 n 个线程,并将它们分配给需要它们的进程。 当这些进程结束时,其线程不会被删除,而是返回线程池,随后可被其他进程重复利用。这比不断创建/删除线程更为高效。因此,我们可以认为使用调度器 [Schedulers.io] 更为合适。

7.7.4. 示例-26:调度器 [Schedulers.immediate, Schedulers.trampoline]

让我们回顾一下关于这两个调度器的说明:

 

该说明虽然简单易懂,但当试图加以说明时,却发现自己并未真正理解。正是《[Learning Reactive Programming With Java 8]》一书让我得以创建一个示例,该示例借鉴了书中一个例子并对其进行了简化。具体如下:


package dvp.rxjava.observables.exemples;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.function.Consumer;

import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Scheduler;
import rx.Scheduler.Worker;
import rx.functions.Action0;
import rx.schedulers.Schedulers;

public class Exemple26 {
    public static void main(String[] args) throws InterruptedException {

        // 一个调度器
        Scheduler scheduler = Schedulers.immediate();
        // 该调度程序的一个工作者
        Worker worker = scheduler.createWorker();
        // 要在该工作进程上执行的 Action0 类型
        Action0 action02 = new Action0() {
            @Override
            public void call() {
                // Action02日志
                ProcessUtils.showInfos.accept("action02");
            }
        };

        // 要在该工作者上执行的 Action0 类型
        Action0 action01 = new Action0() {
            @Override
            public void call() {
                // 在同一工作进程上安排一项新操作
                worker.schedule(action02);
                // Action01日志
                ProcessUtils.showInfos.accept("action01");
            }
        };
        // action01 已安排在该工作节点上
        worker.schedule(action01);
    }

    // 显示
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));

}
  • 第17行:一个调度器。这里使用的是[Schedulers.immediate],后续也可使用[Schedulers.trampoline];
  • 第19行:可以在调度程序的工作线程上执行Action0类型的操作(第21、20行)。 方法 [Scheduler.createWorker] 用于创建一个 worker。方法 [Worker.schedule(Action0)] 用于让 worker 执行类型为 Action0 的操作;
  • 第21-27行:第一个名为[action02]的操作,将由第19行的工作者执行(第40行);
  • 第30-38行:第二个名为[action01]的操作。其特殊之处在于,它会在与自身相同的 worker 上执行操作 action02(第34行)。 这正是 [Schedulers.immediate] 与 [Schedulers.trampoline] 之间的区别:
    • 如果调度器是 [Schedulers.immediate],那么在第 34 行,action02 操作将立即执行(因此得名该调度器),而正在运行的 action01 操作将被中断。 此时将显示第25行的消息。当操作action02完成后,操作action01将恢复执行,并显示第36行的消息;
    • 如果调度器是 [Schedulers.trampoline],则第 34 行,操作 action02 将被挂起。它仅在当前任务 action01 完成后才会执行。 此时将显示第36行的消息。当操作action01完成后,操作action02将开始执行,并显示第25行的消息;

执行上述代码后,结果如下:

action02 ------Thread[main] ---- Time[38:480]
action01 ------Thread[main] ---- Time[38:485]

若在第17行使用调度程序[Schedulers.trampoline],则会得到相反的结果:

action01 ------Thread[main] ---- Time[42:972]
action02 ------Thread[main] ---- Time[42:976]

话虽如此,很难将其与可观察对象建立联系。我尚未找到令人信服的示例,能够说明在上述两个线程中的任一处执行可观察对象有何意义。不过这里有一个示例,但我认为它完全不自然:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.Scheduler.Worker;
import rx.functions.Action1;
import rx.schedulers.Schedulers;

public class Exemple27 {
    public static void main(String[] args) throws InterruptedException {

        // 工作进程
        Worker worker = Schedulers.immediate().createWorker();
        // Worker worker = Schedulers.trampoline().createWorker();
        // 工作线程上的可观察对象 1
        worker.schedule(() -> Observable.range(1, 2).subscribe(new Action1<Integer>() {

            @Override
            public void call(Integer i) {
                ProcessUtils.showInfos.accept(String.valueOf(i));
                // 同一 worker 上的可观察对象 2
                worker.schedule(() -> Observable.range(100, 2).subscribe(new Action1<Integer>() {
                    @Override
                    public void call(Integer i) {
                        ProcessUtils.showInfos.accept(String.valueOf(i));
                    }
                }));
            }
        }));
    }
}
  • 第13-14行:基于两个调度器[Schedulers.immediate]和[Schedulers.trampoline]之一创建一个工作者;
  • 第16行:在该工作者上调度了第一个可观察对象 obs1,用于发布数值 [1,2]
  • 第22行:每当观察到该可观察量obs1的某个元素时,就会在同一工作器上触发第二个可观察量obs2的观察,以生成数值[100,101];

使用调度器 [Schedulers.immediate],可获得以下结果:

1
2
3
4
5
6
1 ------Thread[main] ---- Time[44:604]
100 ------Thread[main] ---- Time[44:610]
101 ------Thread[main] ---- Time[44:610]
2 ------Thread[main] ---- Time[44:612]
100 ------Thread[main] ---- Time[44:612]
101 ------Thread[main] ---- Time[44:612]

而使用调度器 [Schedulers.trampoline] 时,得到以下结果:

1
2
3
4
5
6
1 ------Thread[main] ---- Time[14:107]
2 ------Thread[main] ---- Time[14:114]
100 ------Thread[main] ---- Time[14:115]
101 ------Thread[main] ---- Time[14:115]
100 ------Thread[main] ---- Time[14:115]
101 ------Thread[main] ---- Time[14:116]

7.8. Conclusion

仍有许多工作需要完成。若要深入了解 RxJava 库,建议读者参考本文开头提供的资料继续学习。尽管如此,我们已经掌握了在 Swing 和 Android 环境中使用 RxJava 的基础知识。接下来我们将演示这一点。