8
10

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

More than 5 years have passed since last update.

RxJava のサンプルコードを Reactor-Core で書き直す

8
Last updated at Posted at 2016-06-22

概要

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 ではこのライブラリを使えない可能性があります)

  1. RxJava と同じことができる
  2. クラス名やメソッド名が RxJava より短いものを採用している
  3. 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 に、下記を参考に追記してください。

build.gradle
dependencies {
  compile "io.projectreactor:reactor-core:2.5.0.M4"
......省略......
}

Observable#range での FizzBuzz

単純なサンプルであれば、 rx.Observable を reactor.core.publisher.Flux に変えるだけで動作します。

元のコード:RxFizzBuzz.java

書き換え後

ReactorFizzBuzz.java
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修正
import reactor.core.publisher.Flux; // rx.Observable;
Observable=>Fluxに変更
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

書き換え後

FluxImplementation.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();

ファイルウォッチャー

結構長くなってしまったので、コードはリンク先をご覧ください。主な変更点を以下に書いていきます。

元のコード:RxFileWatcher.java

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 環境では動作しないかもしれません。


リンク

GitHub repositories

  1. Reactor-Core
  2. RxJava
  3. 今回のコード(GitHub repository)

Qiita 記事

  1. Reactor Core 2.5: もう一つのJava向けReactive Extensions実装
  2. Observable を実装して RxJava に入門し直す
8
10
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
8
10

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?