RxJava: istota, komponenty i programowanie reaktywne

Autor: IT Sectr Opublikowano: 2026-05-03 Czas czytania: 10 min

RxJava — to biblioteka programowania reaktywnego dla JVM, implementująca asynchroniczne strumienie danych poprzez wzorzec Observable z funkcyjnymi operatorami transformacji. Portuje koncepcje ReactiveX do Javy i Kotlina, zapewniając jednolite API do pracy z żądaniami sieciowymi, bazami danych, zdarzeniami UI i zadaniami w tle. Według danych ReactiveX, 2025, biblioteka jest używana w ponad 120 000 projektów na GitHub i jest standardem programowania reaktywnego dla Android aż do pojawienia się Kotlin Flow. RxJava zastępuje AsyncTask, Loader i callbacki jednolitym łańcuchem przetwarzania danych.

Najważniejsze

  • RxJava — implementacja ReactiveX dla Java/Kotlin z typami Observable, Flowable, Single, Completable i Maybe
  • Observable reprezentuje strumień danych z zarządzaniem backpressure przez Flowable przy subskrypcji na wolnym consumerze
  • Operatory map, flatMap, switchMap, zip i combineLatest transformują i łączą asynchroniczne strumienie bez blokad
  • Scheduler — Schedulers.io(), computation(), mainThread() zarządzają na którym wątku wykonywana jest praca i subskrypcja
  • RxAndroid dodaje AndroidSchedulers.mainThread() do aktualizacji UI z reaktywnych łańcuchów

Czym jest RxJava?

RxJava — implementacja biblioteki ReactiveX (Reactive Extensions) dla maszyny wirtualnej Javy. Pierwsza wersja RxJava została wydana przez firmę Netflix w 2013 roku do zarządzania asynchronicznymi wywołaniami w aplikacjach serwerowych. W momencie powstania główną alternatywą w Javie były Future i Callback — oba podejścia prowadziły do callback-hell i skomplikowanego zarządzania wątkami. RxJava zaproponowała kompozycję operacji asynchronicznych przez Observable z łańcuchami funkcyjnych operatorów.

Architektura RxJava opiera się na specyfikacji Reactive Streams — standardzie dla asynchronicznego przetwarzania strumieni z nieblokującym backpressure. Specyfikacja definiuje cztery interfejsy: Publisher, Subscriber, Subscription i Processor. RxJava 2+ w pełni implementuje Reactive Streams poprzez typ Flowable, przestrzegając kontraktów backpressure w przeciwieństwie do RxJava 1. Observable w RxJava 2 nie obsługuje backpressure — jest przeznaczony dla strumieni z małą liczbą zdarzeń lub zdarzeń UI.

Według ankiety JetBrains, 2025, RxJava znajduje się w top-3 bibliotek dla programowania Android. Główne scenariusze użycia: obsługa żądań sieciowych przez Retrofit (zintegrowany z RxJava przez CallAdapter), praca z Room (reaktywne zapytania zwracają Flowable lub Maybe), animacje i zdarzenia UI przez RxBinding oraz debounce-wyszukiwanie przy wprowadzaniu tekstu. Wszystkie te scenariusze łączy jednolity łańcuch: źródło (Observable) → transformacja (operatory) → subskrypcja (subscribe).

Historia wersji RxJava

RxJava 1 (2013) zapoczątkował koncepcję Observable i operatorów, ale cierpiał na problemy z backpressure — w szybkich strumieniach dane gromadziły się w pamięci, powodując OutOfMemoryError. RxJava 2 (2016) naprawił architekturę, dzieląc Observable (bez backpressure) i Flowable (z backpressure). RxJava 3 (2020) dodał obsługę Java 8 Stream API, dodatkowe operatory i poprawioną wydajność przy subskrypcji. Obecnie RxJava 3 — zalecana wersja dla nowych projektów.

Typy reaktywnych strumieni w RxJava

RxJava zapewnia pięć głównych typów reaktywnych źródeł, z których każdy jest przeznaczony do określonego scenariusza. Observable i Flowable emitują wiele wartości, Single — jedną wartość lub błąd, Completable — tylko fakt zakończenia bez danych, Maybe — jedną wartość, zero lub błąd. Wybór odpowiedniego typu redukuje ilość kodu i czyni łańcuch samodokumentującym się.

TypLiczba zdarzeńBackpressureScenariusz
Observable0..N, następnie zakończenieNieZdarzenia UI, krótkie strumienie
Flowable0..N, następnie zakończenieTakOdpowiedzi sieciowe, strumienie z BD
SingleDokładnie 1 lub błądNieŻądanie HTTP, odczyt jednego rekordu
Completable0 (tylko zakończenie)NieZapis do BD, wysyłka zdarzenia
Maybe0, 1 lub błądNiePamięć podręczna: jest wartość lub nie

Flowable — najbardziej elastyczny typ do pracy z dużymi strumieniami danych. Implementuje Reactive Streams Publisher z obsługą backpressure: consumer może zażądać określonej liczby elementów przez Subscription.request(n). Zapobiega to przepełnieniu bufora przy niezgodności prędkości producera i consumera. Jeśli backpressure nie jest krytyczny — użyj Observable, ma mniejsze narzuty z powodu braku mechanizmu request.

Single — optymalny wybór dla żądań HTTP. Retrofit 2 z RxJava CallAdapter zwraca Single<ResponseBody> dla każdego żądania. Single gwarantuje dokładnie jedno wywołanie onSuccess lub onError, co odpowiada semantyce żądania HTTP — jedna odpowiedź lub jeden błąd. Completable jest używany do operacji zapisu, które nie zwracają danych: insert, update, delete. Maybe jest wygodny przy sprawdzaniu pamięci podręcznej — może zwrócić wartość, może nie zwrócić.

kotlin
// Przykład użycia Single dla żądania HTTP
interface ApiService {
    @GET("users/{id}")
    fun getUser(@Path("id") userId: Int): Single<User>
}

// Subskrypcja z obsługą na głównym wątku
apiService.getUser(42)
    .subscribeOn(Schedulers.io())
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe({ user ->
        textView.text = user.name
    }, { error ->
        Log.e("API", "Error: ${error.message}")
    })
    .addTo(compositeDisposable)

Operatory transformacji i zarządzania strumieniami

Operatory RxJava — to funkcje wyższego rzędu, które przyjmują jeden reaktywny źródło i zwracają inne, transformując strumień danych. RxJava 3 zawiera ponad 400 operatorów podzielonych na kategorie: transformacja, filtracja, łączenie, obsługa błędów i zarządzanie czasem. Każdy operator jest leniwy — łańcuch budowany jest przy deklaracji, wykonuje się przy subskrypcji.

Operatory transformacji

map — podstawowy operator, przekształcający każdą wartość przez funkcję. flatMap przyjmuje funkcję zwracającą Observable dla każdego elementu i rozwiją wynik w jeden strumień. switchMap jest podobny do flatMap, ale przy otrzymaniu nowego elementu wypisuje się z poprzedniego Observable. concatMap zachowuje kolejność elementów — w przeciwieństwie do flatMap, sekwencyjnie subskrybuje każdy zagnieżdżony Observable.

kotlin
// Parsowanie JSON z transformacją i filtracją
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("Błąd", it.message) })

Łączenie strumieni — obszar, w którym RxJava jest szczególnie silny. zip łączy elementy z wielu Observable parami według indeksu: pierwszy z pierwszym, drugi z drugim. combineLatest emituje nową wartość przy zmianie dowolnego ze strumieni, łącząc ostatnie wartości wszystkich strumieni. merge łączy wiele Observable w jeden, zachowując kolejność nadejścia zdarzeń. concat sekwencyjnie subskrybuje każdy Observable i przekazuje wszystkie jego zdarzenia, zanim przejdzie do następnego.

Zarządzanie czasem obejmuje debounce (oczekiwanie na pauzę w strumieniu przed wysłaniem), throttleFirst (przepuszczenie pierwszego zdarzenia, ignorowanie pozostałych w oknie), timeout (błąd, jeśli zdarzenie nie nadeszło w ciągu interwału). Wyszukiwanie z debounce przy wprowadzaniu tekstu — najczęstszy scenariusz: searchObservable.debounce(300, MILLISECONDS).distinctUntilChanged() zapobiega zbędnym żądaniom przy szybkim pisaniu.

KategoriaOperatorZachowanie
Transformacjamap / flatMap / switchMapPrzekształcenie pojedynczej wartości lub strumienia
Filtracjafilter / distinct / takeSelekcja wartości według warunku
Łączeniezip / combineLatest / mergePołączenie 2+ strumieni
BłędyonErrorResumeNext / retryOdzyskiwanie po awarii
Narzędziadelay / timeout / debounceZarządzanie czasem w strumieniu

Schedulers i wielowątkowość

Scheduler w RxJava — to abstrakcja nad pulą wątków. Biblioteka udostępnia pięć wbudowanych Scheduler: Schedulers.io() dla operacji I/O (sieć, pliki), Schedulers.computation() dla zadań intensywnie wykorzystujących CPU, Schedulers.newThread() dla każdego nowego wątku, Schedulers.single() dla jednowątkowego wykonania i Schedulers.trampoline() dla natychmiastowego wykonania w bieżącym wątku.

subscribeOn i observeOn

subscribeOn określa, na którym Scheduler wykonywane jest źródło Observable. Jeśli w łańcuchu jest kilka subscribeOn — priorytet ma najbliższy źródła. observeOn przełącza downstream na określony Scheduler — każde użycie observeOn zmienia wątek dla kolejnych operatorów. Typowy wzorzec Android: subscribeOn(Schedulers.io()) do pracy z siecią, observeOn(AndroidSchedulers.mainThread()) do aktualizacji UI.

java
// Wielowątkowe przetwarzanie z przełączaniem kontekstu
Observable.fromCallable(() -> database.getItems())
    .subscribeOn(Schedulers.io())            // BD na io
    .map(items -> processItems(items))     // transformacja na io
    .observeOn(Schedulers.computation())    // przełączamy na computation
    .map(processed -> compressImages(processed))
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(result -> ui.showResult(result))

AndroidSchedulers.mainThread() — Scheduler z biblioteki RxAndroid, który wykonuje kod na głównym wątku Android. Jest obowiązkowy dla wszelkich aktualizacji UI w reaktywnym łańcuchu. Biblioteka używa Handler wewnętrznie i gwarantuje wykonanie w wątku UI nawet przy dużym obciążeniu. Dla operacji w tle Schedulers.io() obsługuje nieograniczoną pulę wątków i nadaje się do wszelkich operacji blokujących. Schedulers.computation() używa stałej puli, równej liczbie rdzeni procesora.

RxJava w Android: praktyczne zastosowanie

RxJava w Android jest używany do trzech głównych scenariuszy: reaktywne zapytania do Room, integracja z Retrofit i reaktywne wiązanie UI przez RxBinding. Dla każdego scenariusza charakterystyczny jest własny zestaw typów: Room zwraca Flowable dla obserwowanych zapytań, Retrofit — Single dla żądań HTTP, RxBinding — Observable dla zdarzeń UI.

Room + RxJava

Room — biblioteka trwałości danych od Google. Od wersji Room 2.1 baza danych obsługuje reaktywne typy zwracane: Flowable i Observable. Przy zmianie dowolnego rekordu w tabeli Room automatycznie wysyła nową wartość do strumienia. Programista subskrybuje Flowable w ViewModel i otrzymuje aktualne dane bez ręcznych zapytań przy każdej zmianie.

kotlin
// Room DAO z zapytaniem reaktywnym
@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 — kompozycja Room + Network
class UserViewModel(private val dao: UserDao) : ViewModel() {
    val users: Flowable<List<User>> = dao.getAllUsers()
        .subscribeOn(Schedulers.io())
}

Wzorzec MVVM + RxJava opiera się na tym, że ViewModel nie ma referencji do View. ViewModel publikuje reaktywne źródła (Flowable, LiveData przez Transformations), a Activity lub Fragment subskrybują się do nich. Zapewnia to testowalność: ViewModel jest testowany bez UI, podmieniając Scheduler przez RxJavaPlugins.setComputationScheduler. CompositeDisposable w ViewModel zarządza cyklem życia subskrypcji — przy onCleared() wszystkie subskrypcje są anulowane.

RxJava vs Kotlin Flow

Kotlin Flow — natywna implementacja zimnych strumieni w Kotlin, wbudowana w korutyny i przedstawiona w Kotlin 1.3. Flow rozwiązuje te same zadania co RxJava, ale z fundamentalnymi różnicami: wbudowana obsługa korutyn (funkcje suspend), anulowanie przez coroutine cancellation i brak problemów z backpressure — Flow używa suspend zamiast buforowania. Flow jest częścią standardowej biblioteki Kotlin, nie wymagając dodatkowych zależności.

RxJava pozostaje preferowanym wyborem dla projektów w Javie, projektów z obsługą Java 7-8 oraz istniejących baz kodu na RxJava. Ekosystem RxJava jest znacznie bogatszy: >400 operatorów wobec ~50 w Flow, integracja z Retrofit przez wbudowany CallAdapter, obsługa backpressure przez Flowable oraz obecność RxBinding, RxPermissions, RxLocation dla Android. Kotlin Flow szybko dogania, ale elastyczność RxJava w złożonych scenariuszach łączenia strumieni jest wciąż wyższa.

CharakterystykaRxJavaKotlin Flow
JęzykJava / KotlinKotlin tylko
AnulowanieDisposable / CompositeDisposableCoroutine cancellation
BackpressureFlowable (strategie BUFFER, DROP, LATEST)Przez conflate / buffer
Operatory400+~50 (rozszerzalna)
Room integracjaFlowable, ObservableFlow, StateFlow
ViewModelCompositeDisposableviewModelScope + Flow

Często zadawane pytania

Jaka jest różnica między Observable a Flowable w RxJava?

Observable nie obsługuje backpressure — jeśli producer jest szybszy niż consumer, zdarzenia gromadzą się w pamięci. Flowable implementuje Reactive Streams z backpressure przez Subscription.request(), co zapobiega przepełnieniu bufora przy niezgodności prędkości.

Kiedy używać Single zamiast Observable?

Single jest używany dla operacji, które zwracają dokładnie jedną wartość lub błąd: żądania HTTP, odczyt jednego rekordu z BD, obliczenie wyniku. Single semantycznie odpowiada Future i skraca kod, usuwając nieużywane onComplete.

Jak anulować subskrypcję w RxJava?

Metoda dispose() na Disposable anuluje subskrypcję. Do zarządzania grupowego używany jest CompositeDisposable — zbiera wszystkie Disposable i anuluje je jednocześnie przy wywołaniu clear(). Typowe miejsce — onCleared() w ViewModel lub onPause() w Activity.

Czym różni się flatMap od switchMap?

flatMap subskrybuje wszystkie zagnieżdżone Observable i łączy ich zdarzenia w dowolnej kolejności. switchMap przy nadejściu nowego elementu wypisuje się z poprzedniego Observable i subskrybuje nowy. switchMap jest używany przy wyszukiwaniu — każde nowe żądanie anuluje poprzednie.

Czy warto migrować z RxJava na Kotlin Flow?

Dla nowych projektów w Kotlin Flow jest preferowany dzięki integracji z korutynami i mniejszemu rozmiarowi. Dla istniejących projektów na RxJava migracja jest uzasadniona tylko jeśli cała baza kodu przechodzi na korutyny — pośrednie używanie obu bibliotek komplikuje architekturę.

Podsumowanie

  • RxJava — biblioteka ReactiveX dla JVM z typami Observable, Flowable, Single, Completable i Maybe dla różnych scenariuszy
  • Flowable obsługuje backpressure przez Reactive Streams, aby zapobiec przepełnieniu przy niezgodności prędkości
  • Operatory map, flatMap, switchMap, zip, combineLatest, debounce zapewniają deklaratywne przetwarzanie strumieni
  • Schedulers io(), computation(), mainThread() zarządzają wątkami wykonania bez blokowania UI
  • RxAndroid integruje RxJava z Android, udostępniając AndroidSchedulers.mainThread() i upraszczając aktualizację UI
  • Kotlin Flow — natywna alternatywa z integracją w korutyny, ale RxJava zachowuje przewagę w ekosystemie operatorów
  • MVVM + RxJava — standardowy wzorzec programowania Android z odseparowanym od UI ViewModel i reaktywnymi subskrypcjami

Opracujemy aplikację mobilną pod klucz

IT Sectr tworzy aplikacje na iOS i Androida dla startupów i firm od 2017 roku. Doradzimy Ci i zaproponujemy najlepsze rozwiązanie.

Omów projekt

Przeczytaj również