概要
Java の Reactive Extensions 実装といえば RxJava が有名です。ほかにも Reactor-Core というライブラリがあることを先日 「Reactor Core 2.5: もう一つのJava向けReactive Extensions実装 」を読んで知りました。
この記事では先日私が投稿した記事「Observable を実装して RxJava に入門し直す」 の、RxJava のごく簡単なサンプルを、この Reactor-Core で書き直してみます。
Reactor-Core とは
Reactive Extensions の Java 実装のひとつです。名前は「炉心」という意味のようです。Pivotal がこのライブラリのスポンサーになっています。
Sponsored by Pivotal
v2.5.0.M4 の README.md より
詳細は「Reactor Core 2.5: もう一つのJava向けReactive Extensions実装 」で紹介されているので、ぜひご覧ください。
なお、ライセンスは Apache License Version 2.0 です。
何がよいのか?
正直、前述の紹介記事でも述べられているように、 RxJava から乗り換えることを考えると、決定的なところはないという印象です。新規で Reactive な部分を実装したい時には候補となりうるレベルではあります。(追記:Java8 のクラスライブラリ、具体的には java.util.function のクラスを利用できない Android ではこのライブラリを使えない可能性があります)
- RxJava と同じことができる
- クラス名やメソッド名が RxJava より短いものを採用している
- Flux のコードが Observable よりは小さい……それでも5,000行を超える巨大なクラスではありますが
RxJava から書き換える
簡単に使う程度であれば、下記の対応だけ覚えておけばよさそうです。
クラス
| RxJava | Reactor-Core |
|---|---|
| Observable | Flux |
メソッド
| RxJava | Reactor-Core |
|---|---|
| onNext(T t) | next(T t) |
| onCompleted() | complete() |
| onError(Throwable e) | fail(Throwable e) |
map や filter や subscribe といったオペレータはほぼそのまま使えます。
確認に使った動作環境
Java SE 8で Lambda 式を使います。ご了承ください。
| Java | SE 1.8.0_91 |
|---|---|
| Eclipse | Mars 4.5.2 |
| Reactor-Core | 2.5.0.M4 |
依存の追加
build.gradle の dependencies に、下記を参考に追記してください。
dependencies {
compile "io.projectreactor:reactor-core:2.5.0.M4"
......省略......
}
Observable#range での FizzBuzz
単純なサンプルであれば、 rx.Observable を reactor.core.publisher.Flux に変えるだけで動作します。
書き換え後
import reactor.core.publisher.Flux;
public class ReactorFizzBuzz {
public static final void main(final String[] args) {
Flux.range(1, 100)
.map(i -> {
if (i % 15 == 0) {
return "FizzBuzz";
}
if (i % 3 == 0) {
return "Fizz";
}
if (i % 5 == 0) {
return "Buzz";
}
return Integer.toString(i);
})
.subscribe(
(i) -> System.out.print(i + ", "),
(e) -> e.printStackTrace(),
System.out::println
);
}
}
1, 2, Fizz, 4, Buzz, Fizz, 7, 8, Fizz, Buzz, 11, Fizz, 13, 14, FizzBuzz, 16, 17, Fizz, 19, Buzz, Fizz, 22, 23, Fizz, Buzz, 26, Fizz, 28, 29, FizzBuzz, 31, 32, Fizz, 34, Buzz, Fizz, 37, 38, Fizz, Buzz, 41, Fizz, 43, 44, FizzBuzz, 46, 47, Fizz, 49, Buzz, Fizz, 52, 53, Fizz, Buzz, 56, Fizz, 58, 59, FizzBuzz, 61, 62, Fizz, 64, Buzz, Fizz, 67, 68, Fizz, Buzz, 71, Fizz, 73, 74, FizzBuzz, 76, 77, Fizz, 79, Buzz, Fizz, 82, 83, Fizz, Buzz, 86, Fizz, 88, 89, FizzBuzz, 91, 92, Fizz, 94, Buzz, Fizz, 97, 98, Fizz, Buzz,
修正箇所
下記2つのみです。
import reactor.core.publisher.Flux; // rx.Observable;
Flux.range(1, 100) // Observable.range(1, 100)
diff
1c1
+ package jp.toastkid.verification.reactorcore;
---
- package jp.toastkid.verification.rxjava;
3c3
+ import reactor.core.publisher.Flux;
---
- import rx.Observable;
9c9
+ public class ReactorFizzBuzz {
---
- public class RxFizzBuzz {
13c13
+ Flux.range(1, 100)
---
- Observable.range(1, 100)
Observable オブジェクトを実装した FizzBuzz
Observable オブジェクトを独自に実装するケースも、ほぼそのままと言えるレベルで書き換えられます。
元のコード:ObservablelImplementation.java
書き換え後
import reactor.core.publisher.Flux;
public class FluxImplementation {
public static final void main(final String[] args) {
final Flux<String> flux= makeFlux()
.map(i -> {
if (i % 15 == 0) {
return "FizzBuzz";
}
if (i % 3 == 0) {
return "Fizz";
}
if (i % 5 == 0) {
return "Buzz";
}
return Integer.toString(i);
})
.map(i -> i + ", ");
System.out.println("ぬるぽぬるぽぬるぽ");
flux.subscribe(System.out::println);
}
private static Flux<Integer> makeFlux() {
return Flux.create((sub) -> {
for (int i = 1; i <= 100; i++) {
sub.next(i);
}
sub.complete();
});
}
}
先ほど同様にObservableをFluxに変更しました。ほかに気を付けることはメソッド名が微妙に異なることでしょうか。前述の対応表の通り、onNext は next に、 onCompleted は complete に、それぞれ書き換えています。
diff
+ package jp.toastkid.verification.reactorcore;
---
- package jp.toastkid.verification.rxjava;
3c3
+ import reactor.core.publisher.Flux;
---
- import rx.Observable;
6,7c6
+ * FizzBuzz powered by Reactor-Core.
+ *
---
- * FizzBuzz powered by RxJava.
10,15c9
+ public class FluxImplementation {
- public class ObservablelImplementation {
17c11
+ final Flux<String> observable = makeFlux()
---
- final Observable<String> observable = makeObservable()
34,40c28,29
-
+ /**
+ * Make and return simple ranged Flux.
+ * @return Flux object
+ */
+ private static Flux<Integer> makeFlux() {
+ return Flux.create((sub) -> {
---
- private static Observable<Integer> makeObservable() {
- return Observable.create((sub) -> {
42c31
+ sub.next(i);
---
- sub.onNext(i);
44c33
+ sub.complete();
---
- sub.onCompleted();
ファイルウォッチャー
結構長くなってしまったので、コードはリンク先をご覧ください。主な変更点を以下に書いていきます。
import 修正
import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;
Schedulers の使い方
static factory メソッドが一部 RxJava とは異なっているようです。
makeFileWatcher()
.subscribeOn(Schedulers.newElastic("watcher"))
onError を fail に修正
このケースでも Reactor-Core は短い名前を採用しているようです。
sub.fail(e);
diff
1c1
+ package jp.toastkid.verification.reactorcore;
---
- package jp.toastkid.verification.rxjava;
13,14c13,14
+ import reactor.core.publisher.Flux;
+ import reactor.core.scheduler.Schedulers;
---
- import rx.Observable;
- import rx.schedulers.Schedulers;
22c25
+ public class ReactorFileWatcher {
---
- public class RxFileWatcher {
42c45
+ .subscribeOn(Schedulers.newElastic("watcher"))
---
- .subscribeOn(Schedulers.newThread())
81,83c84,86
+ private static Flux<Path> makeFileWatcher() {
+ return Flux.create((sub) -> {
+ Runtime.getRuntime().addShutdownHook(new Thread(() -> sub.complete()));
---
- private static Observable<Path> makeFileWatcher() {
- return Observable.create((sub) -> {
- Runtime.getRuntime().addShutdownHook(new Thread(() -> sub.onCompleted()));
98c101
+ sub.fail(e);
---
- sub.onError(e);
107c110
+ sub.fail(e);
---
- sub.onError(e);
109c112
+ sub.next(entry.getKey());
---
- sub.onNext(entry.getKey());
112c115
+ System.out.printf("Flux sleeping %dms\n", BACKUP_INTERVAL);
---
- System.out.printf("Observable sleeping %dms\n", BACKUP_INTERVAL);
115c118
+ sub.fail(e);
---
- sub.onError(e);
まとめ
Reactor-Core は基本的なメソッド名や API 設計が RxJava に近く、使用感も共通するところが多いので、違和感なく乗り換えが可能です。これから Reactive Extensions を使い始めるのであれば、Reactor-Core を採用することでクラス名やメソッド名が短いことによる書きやすさの恩恵を受けるのもよいかもしれません。
追記
その後 Flux のソースを眺めてみたところ、RxJava では独自に実装し直していた Function や Predicate を、 Reactor-Core では Java8 で導入された java.util.function package のクラスを用いて実装しているようでした。なので Java8 のライブラリが利用できない Android 環境では動作しないかもしれません。