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 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).
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.
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ę.
| Typ | Liczba zdarzeń | Backpressure | Scenariusz |
|---|---|---|---|
| Observable | 0..N, następnie zakończenie | Nie | Zdarzenia UI, krótkie strumienie |
| Flowable | 0..N, następnie zakończenie | Tak | Odpowiedzi sieciowe, strumienie z BD |
| Single | Dokładnie 1 lub błąd | Nie | Żądanie HTTP, odczyt jednego rekordu |
| Completable | 0 (tylko zakończenie) | Nie | Zapis do BD, wysyłka zdarzenia |
| Maybe | 0, 1 lub błąd | Nie | Pamięć 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ć.
// 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 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.
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.
// 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.
| Kategoria | Operator | Zachowanie |
|---|---|---|
| Transformacja | map / flatMap / switchMap | Przekształcenie pojedynczej wartości lub strumienia |
| Filtracja | filter / distinct / take | Selekcja wartości według warunku |
| Łączenie | zip / combineLatest / merge | Połączenie 2+ strumieni |
| Błędy | onErrorResumeNext / retry | Odzyskiwanie po awarii |
| Narzędzia | delay / timeout / debounce | Zarządzanie czasem w strumieniu |
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 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.
// 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 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 — 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.
// 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.
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.
| Charakterystyka | RxJava | Kotlin Flow |
|---|---|---|
| Język | Java / Kotlin | Kotlin tylko |
| Anulowanie | Disposable / CompositeDisposable | Coroutine cancellation |
| Backpressure | Flowable (strategie BUFFER, DROP, LATEST) | Przez conflate / buffer |
| Operatory | 400+ | ~50 (rozszerzalna) |
| Room integracja | Flowable, Observable | Flow, StateFlow |
| ViewModel | CompositeDisposable | viewModelScope + Flow |
Często zadawane pytania
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.
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.
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.
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.
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
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.
Przeczytaj również