前回はRxJavaの構成についてまとめました(記事はこちら)。今回も引き続きRxJavaの構成をまとめてみます。
| 戻り値の型 | メソッド | 説明 |
|---|---|---|
| Disposable | subscribe() | Flowable/Observableの処理だけ行いSubscriber/Observerは何もしない |
| Disposable | subscribe(Consumer onNext) | データの通知(onNext)を受け取った時の処理のみ引数で定 義してあるように行う |
| Disposable | subscribe( Consumer onNext, Consumer onError) | データの通知(onNext)とエラーの通知(onError)を受け取った時の処理のみ引数で定義してあるように行う |
| Disposable | subscribe( Consumer onNext, Consumer onError, Action onComplete) | データの通知(onNext)とエラーの通知(onError)と完了の通知(onComplete)を受け取った時の処理のみ引数で定義してあるように行う |
| Disposable | subscribe( Consumer onNext, Consumer onError, Action onComplete, Consumer onSubscribe) | データの通知(onNext)とエラーの通知(onError)と完了の通知(onComplete)を受け取った時の処理を引数で定義してあるように行い、さらに購読開始時の処理(onSubscribe)も定義してあるように行う |
※subscribeメソッドでは、デフォルトでリクエストするデータ数にはLong.MAX_VALUEが設定されている。
※エラー時の処理を指定していない場合は、エラーの通知を受けてもスタックトレースを出力するだけでそれ以外は処理しない。
// 購読を開始する。
Disposable disposable = flowable.subscribe(data -> System.out.println("data=" + data));
// 購読を解除する
disposable.dispose();
RxJava2.xではsubscribeメソッドに加え、新たに購読を行うためのメソッドとして、Subscriber/Observerを引数に取り、戻り値も返すsubscribeWithメソッドがある。
引数にSubscriber/Observerを渡すと内部でそのSubscriber/Observerをsubscribeメソッドに渡して実行し、戻り値としてその引数となったSubscriber/Observerを返す。
public final <E extends Subscriber<? super T>> E subscribeWith(E subscriber) {
subscribe(subscriber);
return subscriber;
}
Disposable disposable = flowable.subscribeWith(new ResourceSubscriber(){
// 中略
});
CompositeDisposableは複数のDisposableをまとめることで、CompositeDisposableのdisposableのメソッドを呼ぶことで保持している全てのDisposableのdisposeメソッドを呼ぶことができる。
public static void main(String[] args) throws Exception {
// Disposableをまとめる
CompositeDisposable compositeDisposable = new CompositeDisposable();
compositeDisposable.add(Flowable.range(1, 3)
.doOnCancel(() -> System.out.println("No.1 canceled"))
.observeOn(Schedulers.computation())
.subscribe(data -> {
Thread.sleep(100L);
System.out.println("No.1: " + data);
}));
compositeDisposable.add(Flowable.range(1, 3)
.doOnCancel(() -> System.out.println("No.2 canceled"))
.observeOn(Schedulers.computation())
.subscribe(data -> {
Thread.sleep(100L);
System.out.println("No.2: " + data);
}));
// しばらく待つ
Thread.sleep(150L);
// まとめて購読の解除を行う
compositeDisposable.
| クラス | 説明 |
|---|---|
| Single | データを1件だけ通知するか、もしくはエラーを通知するクラス |
| Maybe | データを1件だけ通知するか、1件も通知せず完了を通知するか、もしくはエラーを通知するクラス |
| Completable | データを1件も通知せず完了を通知するか、もしくはエラーを通知するクラス |
| 生産者 | 消費者 |
|---|---|
| Single | SingleObserver |
| Maybe | MaybeObserver |
| Completable | CompletableObserver |
Singleはデータを1件だけ通知するか、もしくはエラーを通知するクラス
データを通知することは処理が完了したことも意味するので、完了の通知は存在しない。
Singleが持つ通知プロトコルにonNextとonCompleteはなく、onSuccessという1件のデータを通知し完了したことを意味する通知プロトコルが用意されている。
データが1件しか通知されないので、データ数をリクエストする必要がない。
| メソッド | 説明 |
|---|---|
| onSubscribe(Disposable disposable) | 通知の準備ができたら呼ばれるメソッド。引数に購読の解除を行えるDisposableを受け取る |
| onSuccess(T data) | データを受け取り処理を行うメソッド。これ以降のデータの通知はないので、onSuccessメソッドが呼ばれることはSingleの処理が完了したことを意味する |
| onError(Throwable error) | 通知処理を行っている間にエラーが発生したら呼ばれるメソッド。 引数に発生したエラーのオブジェクトが渡される |
public static void main(String[] args) {
// Singleの作成
Single<DayOfWeek> single = Single.create(emitter ->
emitter.onSuccess(LocalDate.now().getDayOfWeek());
});
// 購読する
single.subscribe(new SingleObserver<DayOfWeek>() {
// 購読の準備ができた際の処理を行う
@Override
public void onSubscribe(Disposable disposable) {
// 何もしない
}
// データの通知を受け取った際の処理を行う
@Override
public void onSuccess(DayOfWeek value) {
System.out.println(value);
}
// エラーの通知を受け取った際の処理を行う
@Override
public void onError(Throwable e) {
System.out.println("エラー=" + e);
}
});
}
Maybeはデータをデータを1件だけ通知するか、1件も通知せず完了を通知するか、もしくはエラーを通知するクラス
Maybeではデータを通知することは処理が完了したことも意味するので、あえて再び完了の通知を行わない。
Maybeが完了の通知を行う場合は、データが1件もなく処理が正常に終了した場合となる。
Maybeの完了の通知(onComplete)は、正常に処理が終了した際に、必ず呼ばれるわけではない
| メソッド | 説明 |
|---|---|
| onSubscribe(Disposable disposable) | 通知の準備ができたら呼ばれるメソッド。引数に購読の解除を行えるDisposableを受け取る |
| onSuccess(T data) | データを受け取り処理を行うメソッド。これ以降のデータの通知はないので、onSuccessメソッドが呼ばれることは、Maybeの処理が完了したことを意味し、onCompleteは呼ばれない |
| onComplete | データを通知することなく、Maybeの処理が完了した際に実行されるメソッド |
| onError(Throwable error) | 通知処理を行っている間にエラーが発生したら呼ばれるメソッド。 引数に発生したエラーのオブジェクトが渡される |
public static void main(String[] args) {
// Maybeの作成
Maybe<DayOfWeek> maybe = Maybe.create(emitter -> {
emitter.onSuccess(LocalDate.now().getDayOfWeek());
});
// 購読する
maybe.subscribe(new MaybeObserver<DayOfWeek>() {
// 購読の準備ができた際の処理を行う
@Override
public void onSubscribe(Disposable disposable) {
// 何もしない
}
// データの通知を受け取った際の処理を行う
@Override
public void onSuccess(DayOfWeek value) {
System.out.println(value);
}
// 完了の通知を受け取った際の処理を行う
@Override
public void onComplete() {
System.out.println("完了");
}
// エラーの通知を受け取った際の処理を行う
@Override
public void onError(Throwable e) {
System.out.println("エラー=" + e);
}
});
Completableはデータを通知することなく完了を通知するか、もしくはエラーを通知するクラス
他の生産者となるクラスと異なり、データを通知することはない
Completableは他の生産者となるクラスと役割が異なり、Completable内で何らかの副作用が発生する処理を行う。
処理が終わった場合に完了の通知を行い、エラーが発生したらエラーの通知を行う。
Completableを使う場合は、その処理の購読を呼び出しているスレッドと異なるスレッド上で行わせないと、RxJavaを使わない場合と同じ処理内容になり、Completableを使う意味がなくなる。
| メソッド | 説明 |
|---|---|
| onSubscribe(Disposable disposable) | 通知の準備ができたら呼ばれるメソッド。引数に購読の解除を行えるDisposableを受け取る |
| onComplete | Completableの処理が完了した際に実行されるメソッド |
| onError(Throwable error) | 通知処理を行っている間にエラーが発生したら呼ばれるメソッド。引数に発生したエラーのオブジェクトが渡される |
public static void main(String[] args) throws Exception {
// Completableの作成
Completable completable = Completable.create(emitter ->
// …略 何らかの処理を行う
// 完了を通知する
emitter.onComplete();
});
completable
// Completableを非同期で行う
.subscribeOn(Schedulers.computation())
// 購読する
.subscribe(new CompletableObserver() { // ❸
// 購読の準備ができた際の処理を行う
@Override
public void onSubscribe(Disposable disposable) {
// 何もしない
}
// 完了の通知を受け取った際の処理を行う
@Override
public void onComplete() {
System.out.println("完了");
}
// エラーの通知を受け取った際の処理を行う
@Override
public void onError(Throwable e) {
System.out.println("エラー=" + e);
}
});
// しばらく待つ
Thread.sleep(100L);
}
RxJavaは軽量化されているため基本的には必要最低限の機能しか持たない
次のモジュールはRxJavaの拡張として用意されているもの
Android用としては次のモジュールがある。
written by tamito0201
プログラミングとのご縁結びならプロマリへ。
オンラインプログラミング学習スクールのプロマリは、プログラミングの初学者の皆様を応援しています。プログラミング講師と一緒に面白いアプリを作りませんか。
The programming school "Promari" will help you learn programming. "Promari" is supporting the first scholars of programming. Let's develop an application with our programming instructor.