RxJava:本質、コンポーネント、リアクティブプログラミング

著者: IT Sectr 公開日: 2026-05-03 読了時間: 10 分

RxJavaは、関数型変換演算子を備えたObservableパターンを通じて非同期データストリームを実装する、JVM向けのリアクティブプログラミングライブラリです。ReactiveXの概念をJavaとKotlinに移植し、ネットワークリクエスト、データベース、UIイベント、バックグラウンドタスクを扱うための統一APIを提供します。ReactiveX、2025によると、このライブラリはGitHub上の12万以上のプロジェクトで使用され、Kotlin Flow登場までのAndroidリアクティブプログラミングの標準です。RxJavaはAsyncTask、Loader、コールバックを単一のデータ処理チェーンに置き換えます。

重要ポイント

  • RxJavaは、Observable、Flowable、Single、Completable、Maybe型を備えたJava/Kotlin向けReactiveX実装です
  • Observableは、低速コンシューマーで購読する際にFlowableを介してバックプレッシャー管理を行うデータストリームを表します
  • 演算子map、flatMap、switchMap、zip、combineLatestは非同期ストリームをブロッキングなしで変換・結合します
  • Scheduler — Schedulers.io()、computation()、mainThread()は、どのスレッドで処理と購読を実行するかを管理します
  • RxAndroidはリアクティブチェーンからUIを更新するためのAndroidSchedulers.mainThread()を追加します

RxJavaとは?

RxJavaは、Java仮想マシン向けのReactiveX(Reactive Extensions)ライブラリの実装です。RxJavaの最初のバージョンは、サーバーサイドアプリケーションでの非同期呼び出しを管理するためにNetflixによって2013年にリリースされました。作成当時、Javaでの主な代替手段はFutureとコールバックでした — どちらのアプローチもコールバック地獄と複雑なスレッド管理につながりました。RxJavaは、関数型演算子のチェーンを備えたObservableを通じて、非同期操作の合成を導入しました。

RxJavaのアーキテクチャは、Reactive Streams仕様に基づいています — ノンブロッキングバックプレッシャーによる非同期ストリーム処理の標準です。この仕様は、Publisher、Subscriber、Subscription、Processorの4つのインターフェースを定義しています。RxJava 2+は、RxJava 1とは異なりバックプレッシャーの契約に準拠して、Flowable型を通じてReactive Streamsを完全に実装しています。RxJava 2のObservableはバックプレッシャーをサポートしていません — 少数のイベントまたはUIイベントのストリームを対象としています。

JetBrains、2025の調査によると、RxJavaはAndroid開発のトップ3ライブラリの1つです。主なユースケースとしては、Retrofitを介したネットワークリクエスト処理(CallAdapterを介してRxJavaと統合)、Roomとの連携(リアクティブクエリはFlowableまたはMaybeを返す)、RxBindingを介したアニメーションとUIイベント、テキスト入力時のデバウンス検索があります。これらすべてのシナリオは、共通のチェーンパターンを共有しています:ソース(Observable)→ 変換(演算子)→ 購読(subscribe)。

RxJavaのバージョン履歴

RxJava 1(2013年)はObservableと演算子の基礎を築きましたが、バックプレッシャーの問題に悩まされていました — 高速ストリームではデータがメモリに蓄積され、OutOfMemoryErrorを引き起こしました。RxJava 2(2016年)は、Observable(バックプレッシャーなし)とFlowable(バックプレッシャーあり)を分離してアーキテクチャを修正しました。RxJava 3(2020年)はJava 8 Stream APIのサポート、追加の演算子、改善された購読パフォーマンスを追加しました。現在、RxJava 3が新規プロジェクトの推奨バージョンです。

RxJavaのリアクティブストリームの種類

RxJavaは、それぞれ特定のシナリオ向けに設計された5つの主要なリアクティブソース型を提供します。ObservableとFlowableは複数の値を発行し、Singleは1つの値またはエラーを発行し、Completableはデータなしで完了のみを発行し、Maybeは1つ、ゼロ、またはエラーを発行します。正しい型を選択することでコード量が減り、チェーンが自己文書化されます。

イベント数バックプレッシャーシナリオ
Observable0..N、その後完了なしUIイベント、短いストリーム
Flowable0..N、その後完了ありネットワーク応答、DBストリーム
Singleちょうど1つまたはエラーなしHTTPリクエスト、1レコードの読み取り
Completable0(完了のみ)なしDB書き込み、イベント送信
Maybe0、1、またはエラーなしキャッシュ:値ありまたはなし

Flowableは、大規模なデータストリームを扱うための最も柔軟な型です。バックプレッシャーをサポートするReactive Streams Publisherを実装しており、コンシューマーはSubscription.request(n)を介して特定の数の要素を要求できます。これにより、プロデューサーとコンシューマーの速度が一致しない場合のバッファオーバーフローを防ぎます。バックプレッシャーが重要でない場合は、Observableを使用してください — requestメカニズムがないためオーバーヘッドが少なくなります。

SingleはHTTPリクエストに最適な選択です。RxJava CallAdapterを使用したRetrofit 2は、各リクエストに対してSingle<ResponseBody>を返します。Singleはちょうど1回のonSuccessまたはonError呼び出しを保証し、HTTPリクエストのセマンティクス(1つの応答または1つのエラー)と一致します。Completableはデータを返さない書き込み操作(insert、update、delete)に使用されます。Maybeはキャッシュの確認に便利です — 値を返す場合と返さない場合があります。

kotlin
// HTTPリクエストにSingleを使用する例
interface ApiService {
    @GET("users/{id}")
    fun getUser(@Path("id") userId: Int): Single<User>
}

// メインスレッドでの処理を伴う購読
apiService.getUser(42)
    .subscribeOn(Schedulers.io())
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe({ user ->
        textView.text = user.name
    }, { error ->
        Log.e("API", "Error: ${error.message}")
    })
    .addTo(compositeDisposable)

変換およびストリーム管理演算子

演算子はRxJavaにおける高階関数で、1つのリアクティブソースを受け取り、データストリームを変換して別のソースを返します。RxJava 3には400以上の演算子が含まれており、変換、フィルタリング、結合、エラー処理、時間管理のカテゴリに分かれています。各演算子は遅延評価され、チェーンは宣言時に構築され、購読時に実行されます。

変換演算子

mapは、関数を介して各値を変換する基本演算子です。flatMapは各要素に対してObservableを返す関数を受け取り、結果を単一のストリームに平坦化します。switchMapはflatMapに似ていますが、新しい要素が到着すると前のObservableから購読を解除します。concatMapは要素の順序を保持します — flatMapとは異なり、ネストされた各Observableに順次購読します。

kotlin
// 変換とフィルタリングを伴うJSON解析
apiService.getUsers()
    .flatMap { users ->
        Observable.fromIterable(users)
    }
    .filter { user ->
        user.age >= 18
    }
    .map { user ->
        UserDto(user.name, user.age)
    }
    .toList()
    .subscribeOn(Schedulers.computation())
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe({ adapter.submitList(it) },
               { Log.e("エラー", it.message) })

ストリームの結合はRxJavaが特に優れている分野です。zipは複数のObservablesからの要素をインデックスごとにペアで結合します:最初は最初と、2番目は2番目と。combineLatestはいずれかのストリームが変更されたときに新しい値を発行し、すべてのストリームの最新値を結合します。mergeは複数のObservablesを1つに結合し、イベントの到着順を保持します。concatは各Observableに順次購読し、次のObservableに移る前にそのすべてのイベントを渡します。

時間管理には、debounce(発行前にストリームの一時停止を待機)、throttleFirst(最初のイベントを発行し、ウィンドウ内の残りを無視)、timeout(間隔内にイベントが届かない場合のエラー)が含まれます。テキスト入力時のデバウンス検索は最も一般的なシナリオです:searchObservable.debounce(300, MILLISECONDS).distinctUntilChanged()は高速入力時の不要なリクエストを防ぎます。

カテゴリ演算子動作
変換map / flatMap / switchMap値またはストリームを変換
フィルタリングfilter / distinct / take条件に基づいて値を選択
結合zip / combineLatest / merge2つ以上のストリームを結合
エラーonErrorResumeNext / retry障害から回復
ユーティリティdelay / timeout / debounceストリームの時間管理

Schedulerとマルチスレッド

SchedulerはRxJavaにおけるスレッドプールの抽象化です。ライブラリは5つの組み込みSchedulerを提供します:I/O操作用のSchedulers.io()(ネットワーク、ファイル)、CPU集中型タスク用のSchedulers.computation()、毎回新しいスレッド用のSchedulers.newThread()、シングルスレッド実行用のSchedulers.single()、現在のスレッドで即時実行するためのSchedulers.trampoline()。

subscribeOnとobserveOn

subscribeOnは、ソースObservableを実行するSchedulerを決定します。チェーン内に複数のsubscribeOnがある場合、ソースに最も近いものが優先されます。observeOnはダウンストリームを指定されたSchedulerに切り替えます — observeOnを使用するたびに、後続の演算子のスレッドが変更されます。典型的なAndroidパターン:ネットワーク操作にはsubscribeOn(Schedulers.io())、UI更新にはobserveOn(AndroidSchedulers.mainThread())です。

java
// コンテキスト切り替えを伴うマルチスレッド処理
Observable.fromCallable(() -> database.getItems())
    .subscribeOn(Schedulers.io())            // io上のDB
    .map(items -> processItems(items))     // io上の変換
    .observeOn(Schedulers.computation())    // computationに切り替え
    .map(processed -> compressImages(processed))
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(result -> ui.showResult(result))

AndroidSchedulers.mainThread()は、Androidメインスレッドでコードを実行するRxAndroidライブラリのSchedulerです。リアクティブチェーンでのUI更新には必須です。ライブラリは内部的にHandlerを使用し、高負荷下でもUIスレッドでの実行を保証します。バックグラウンド操作には、Schedulers.io()が無制限のスレッドプールをサポートし、あらゆるブロッキング操作に適しています。Schedulers.computation()はCPUコア数に等しい固定プールを使用します。

AndroidでのRxJava:実践的な応用

RxJavaはAndroidで3つの主要なシナリオに使用されます:Roomへのリアクティブクエリ、Retrofitとの統合、RxBindingを介したリアクティブUIバインディング。各シナリオには固有の型セットがあります:Roomは監視可能なクエリにはFlowableを返し、RetrofitはHTTPリクエストにはSingleを返し、RxBindingはUIイベントにはObservableを返します。

Room + RxJava

RoomはGoogleのデータ永続化ライブラリです。Room 2.1以降、データベースはFlowableとObservableのリアクティブ戻り値型をサポートしています。テーブル内のレコードが変更されると、Roomは自動的に新しい値をストリームに送信します。開発者はViewModelでFlowableを購読し、変更のたびに手動クエリなしで最新のデータを取得します。

kotlin
// リアクティブクエリを使用したRoom DAO
@Dao
interface UserDao {
    @Query("SELECT * FROM users WHERE id = :id")
    fun getUserById(@Param("id") userId: Int): Flowable<User>

    @Insert
    fun insertUser(user: User): Completable
}

// ViewModel — Room + Networkの構成
class UserViewModel(private val dao: UserDao) : ViewModel() {
    val users: Flowable<List<User>> = dao.getAllUsers()
        .subscribeOn(Schedulers.io())
}

MVVM + RxJavaパターンは、ViewModelがViewへの参照を持たないという考えに基づいています。ViewModelはリアクティブソース(Flowable、Transformationsを介したLiveData)を公開し、ActivityまたはFragmentがそれらを購読します。これによりテスト容易性が向上します:ViewModelはRxJavaPlugins.setComputationSchedulerを介してSchedulerを置き換えることで、UIなしでテストできます。ViewModelのCompositeDisposableは購読のライフサイクルを管理し、onCleared()ですべての購読が解除されます。

RxJava vs Kotlin Flow

Kotlin FlowはKotlinにおけるコールドストリームのネイティブ実装で、コルーチンに組み込まれ、Kotlin 1.3で導入されました。FlowはRxJavaと同じ問題を解決しますが、根本的な違いがあります:コルーチンの組み込みサポート(suspend関数)、coroutine cancellationによるキャンセル、バックプレッシャーの問題がないこと(Flowはバッファリングの代わりにsuspendを使用)。FlowはKotlin標準ライブラリの一部であり、追加の依存関係を必要としません。

RxJavaは、Javaプロジェクト、Java 7-8をサポートするプロジェクト、既存のRxJavaコードベースにとって依然として好ましい選択肢です。RxJavaのエコシステムははるかに豊富です:Flowの約50に対して400以上の演算子、組み込みCallAdapterを介したRetrofitとの統合、Flowableによるバックプレッシャーサポート、Android向けのRxBinding、RxPermissions、RxLocation。Kotlin Flowは急速に追いついていますが、複雑なストリーム結合シナリオにおけるRxJavaの柔軟性は依然として高いです。

特性RxJavaKotlin Flow
言語Java / KotlinKotlinのみ
キャンセルDisposable / CompositeDisposableCoroutine cancellation
バックプレッシャーFlowable(BUFFER、DROP、LATEST戦略)conflate / bufferによる
演算子400+〜50(拡張可能)
Room統合Flowable、ObservableFlow、StateFlow
ViewModelCompositeDisposableviewModelScope + Flow

よくある質問

RxJavaのObservableとFlowableの違いは何ですか?

Observableはバックプレッシャーをサポートしていません — プロデューサーがコンシューマーより速い場合、イベントがメモリに蓄積されます。FlowableはSubscription.request()を介したバックプレッシャーでReactive Streamsを実装し、速度が一致しない場合のバッファオーバーフローを防ぎます。

Observableの代わりにSingleを使用するのはいつですか?

Singleは、ちょうど1つの値またはエラーを返す操作(HTTPリクエスト、DBからの1レコードの読み取り、結果の計算)に使用されます。Singleは意味的にFutureに対応し、未使用のonCompleteを削除することでコードを削減します。

RxJavaで購読をキャンセルするには?

Disposableのdispose()メソッドが購読をキャンセルします。グループ管理にはCompositeDisposableが使用され、すべてのDisposableを収集し、clear()の呼び出しで同時に破棄します。典型的な場所はViewModelのonCleared()またはActivityのonPause()です。

flatMapとswitchMapの違いは何ですか?

flatMapはすべてのネストされたObservableに購読し、それらのイベントを任意の順序でマージします。switchMapは新しい要素が到着すると前のObservableから購読を解除し、新しいObservableに購読します。switchMapは検索で使用され、新しいリクエストが前のリクエストをキャンセルします。

RxJavaからKotlin Flowに移行すべきですか?

新しいKotlinプロジェクトでは、コルーチン統合とサイズが小さいためFlowが推奨されます。既存のRxJavaプロジェクトでは、コードベース全体がコルーチンに移行する場合にのみ移行が正当化されます — 両方のライブラリを中間的に使用するとアーキテクチャが複雑になります。

まとめ

  • RxJavaは、さまざまなシナリオ向けのObservable、Flowable、Single、Completable、Maybe型を備えたJVM向けReactiveXライブラリです
  • Flowableは速度不一致時のオーバーフローを防ぐためにReactive Streamsを介したバックプレッシャーをサポートします
  • 演算子map、flatMap、switchMap、zip、combineLatest、debounceは宣言的なストリーム処理を提供します
  • Schedulerio()、computation()、mainThread()はUIをブロックせずに実行スレッドを管理します
  • RxAndroidはAndroidSchedulers.mainThread()を提供し、UI更新を簡素化してRxJavaをAndroidと統合します
  • Kotlin Flowはコルーチン統合を備えたネイティブの代替手段ですが、RxJavaは演算子エコシステムで優位性を維持しています
  • MVVM + RxJavaはUIから分離されたViewModelとリアクティブ購読を備えた標準的なAndroid開発パターンです

ターンキー方式のモバイルアプリケーションを開発します

IT Sectrは2017年からスタートアップや企業向けにiOS・Androidアプリケーションを開発しています。私たちがご相談に乗り、最適なソリューションをご提案します。

プロジェクトについて相談

こちらもお読みください