RxJava egy reaktív programozási könyvtár a JVM számára, amely aszinkron adatfolyamokat valósít meg az Observable minta segítségével funkcionális transzformációs operátorokkal. A ReactiveX koncepcióit Java-ra és Kotlin-ra portolja, egységes API-t biztosítva hálózati kérések, adatbázisok, UI-események és háttérfeladatok kezeléséhez. A ReactiveX, 2025 adatai szerint a könyvtárat több mint 120 000 projekt használja a GitHub-on, és ez a reaktív programozás szabványa Android-hoz a Kotlin Flow megjelenéséig. Az RxJava leváltja az AsyncTask-ot, Loader-t és callback-eket egy egységes adatfeldolgozási lánccal.
Főbb pontok
RxJava a ReactiveX (Reactive Extensions) könyvtár implementációja a Java virtuális gépre. Az RxJava első verzióját a Netflix cég adta ki 2013-ban aszinkron hívások kezelésére szerveralkalmazásokban. A létrehozáskor a fő alternatíva Java-ban a Future és a Callback volt — mindkét megközelítés callback-hell-hez és bonyolult szálkezeléshez vezetett. Az RxJava az aszinkron műveletek kompozícióját javasolta Observable-en keresztül funkcionális operátorok láncaival.
Az RxJava architektúrája a Reactive Streams specifikáción alapul — ez egy szabvány az aszinkron folyamfeldolgozásra nem blokkoló backpressure-rel. A specifikáció négy interfészt határoz meg: Publisher, Subscriber, Subscription és Processor. Az RxJava 2+ teljes mértékben implementálja a Reactive Streams-et a Flowable típuson keresztül, betartva a backpressure szerződéseket, ellentétben az RxJava 1-gyel. Az Observable az RxJava 2-ben nem támogatja a backpressure-t — kis számú eseményt tartalmazó folyamokhoz vagy UI-eseményekhez készült.
A JetBrains, 2025 felmérése szerint az RxJava a top-3 könyvtár közé tartozik Android-fejlesztéshez. Fő használati forgatókönyvek: hálózati kérések feldolgozása Retrofit-en keresztül (integrálva RxJava-val CallAdapter-en keresztül), munka Room-mal (reaktív lekérdezések Flowable-t vagy Maybe-t adnak vissza), animációk és UI-események RxBinding-en keresztül, és debounce keresés szövegbevitelkor. Mindezeket a forgatókönyveket egy azonos típusú lánc egyesíti: forrás (Observable) → transzformáció (operátorok) → feliratkozás (subscribe).
RxJava 1 (2013) lefektette az Observable és az operátorok koncepcióját, de backpressure problémákkal küzdött — gyors folyamokban az adatok felhalmozódtak a memóriában, OutOfMemoryError-t okozva. RxJava 2 (2016) kijavította az architektúrát, szétválasztva az Observable-t (backpressure nélkül) és a Flowable-t (backpressure-rel). RxJava 3 (2020) hozzáadta a Java 8 Stream API támogatását, további operátorokat és javított teljesítményt feliratkozáskor. Jelenleg az RxJava 3 az ajánlott verzió új projektekhez.
RxJava öt fő reaktív forrástípust biztosít, mindegyik egy adott forgatókönyvre összpontosítva. Az Observable és a Flowable több értéket bocsát ki, a Single — egy értéket vagy hibát, a Completable — csak a befejezés tényét adatok nélkül, a Maybe — egy értéket, nullát vagy hibát. A megfelelő típus kiválasztása csökkenti a kód mennyiségét és öndokumentálóvá teszi a láncot.
| Típus | Események száma | Backpressure | Forgatókönyv |
|---|---|---|---|
| Observable | 0..N, majd befejezés | Nem | UI-események, rövid folyamok |
| Flowable | 0..N, majd befejezés | Igen | Hálózati válaszok, folyamok adatbázisból |
| Single | Pontosan 1 vagy hiba | Nem | HTTP-kérés, egy rekord olvasása |
| Completable | 0 (csak befejezés) | Nem | Írás adatbázisba, esemény küldése |
| Maybe | 0, 1 vagy hiba | Nem | Gyorsítótár: van érték vagy nincs |
Flowable a legrugalmasabb típus nagy adatfolyamok kezelésére. Megvalósítja a Reactive Streams Publisher-t backpressure támogatással: a consumer meghatározott számú elemet kérhet a Subscription.request(n)-en keresztül. Ez megakadályozza a puffer túlcsordulását a producer és consumer sebességének eltérésekor. Ha a backpressure nem kritikus — használja az Observable-t, amely kisebb terheléssel jár a request mechanizmus hiánya miatt.
Single az optimális választás HTTP-kérésekhez. A Retrofit 2 RxJava CallAdapter-rel Single<ResponseBody>-t ad vissza minden kéréshez. A Single pontosan egy onSuccess vagy onError hívást garantál, ami megfelel a HTTP-kérés szemantikájának — egy válasz vagy egy hiba. Completable olyan írási műveletekhez használatos, amelyek nem adnak vissza adatokat: insert, update, delete. Maybe kényelmes a gyorsítótár ellenőrzésekor — visszaadhat értéket, de nem is.
// Példa a Single használatára HTTP-kéréshez
interface ApiService {
@GET("users/{id}")
fun getUser(@Path("id") userId: Int): Single<User>
}
// Feliratkozás feldolgozással a főszálon
apiService.getUser(42)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe({ user ->
textView.text = user.name
}, { error ->
Log.e("API", "Error: ${error.message}")
})
.addTo(compositeDisposable)
Operátorok RxJava magasabb rendű függvények, amelyek egy reaktív forrást fogadnak és egy másikat adnak vissza, transzformálva az adatfolyamot. Az RxJava 3 több mint 400 operátort tartalmaz kategóriákra bontva: transzformáció, szűrés, kombinálás, hibakezelés és időkezelés. Minden operátor lusta — a lánc deklarációkor épül fel, feliratkozáskor hajtódik végre.
map az alap operátor, amely minden értéket egy függvényen keresztül transzformál. A flatMap olyan függvényt fogad, amely minden elemhez Observable-t ad vissza, és az eredményt egyetlen folyammá bontja ki. A switchMap hasonlít a flatMap-re, de új elem érkezésekor leiratkozik az előző Observable-ről. A concatMap megőrzi az elemek sorrendjét — a flatMap-től eltérően szekvenciálisan iratkozik fel minden beágyazott Observable-re.
// JSON feldolgozás transzformációval és szűréssel
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("Hiba", it.message) })
Folyamok kombinálása az a terület, ahol az RxJava különösen erős. A zip több Observable elemeit páronként kombinálja index szerint: elsőt az elsővel, másodikat a másodikkal. A combineLatest új értéket bocsát ki bármely folyam változásakor, kombinálva az összes folyam legutóbbi értékeit. A merge több Observable-t egyesít egyetlen folyammá, megőrizve az események érkezési sorrendjét. A concat szekvenciálisan iratkozik fel minden Observable-re, és továbbítja az összes eseményét, mielőtt a következőre lépne.
Időkezelés magában foglalja a debounce-t (szünet várása a folyamban küldés előtt), a throttleFirst-t (első esemény átengedése, a többi figyelmen kívül hagyása az ablakban), a timeout-ot (hiba, ha az esemény nem érkezett meg az intervallumon belül). A debounce keresés szövegbevitelkor a leggyakoribb forgatókönyv: searchObservable.debounce(300, MILLISECONDS).distinctUntilChanged() megakadályozza a felesleges kéréseket gyors gépeléskor.
| Kategória | Operátor | Viselkedés |
|---|---|---|
| Transzformáció | map / flatMap / switchMap | Egyetlen érték vagy folyam transzformálása |
| Szűrés | filter / distinct / take | Értékek kiválasztása feltétel alapján |
| Kombinálás | zip / combineLatest / merge | 2+ folyam egyesítése |
| Hibák | onErrorResumeNext / retry | Helyreállítás hiba után |
| Segédprogramok | delay / timeout / debounce | Időkezelés a folyamban |
Scheduler az RxJava-ban egy absztrakció a szálpool felett. A könyvtár öt beépített Schedulert biztosít: Schedulers.io() I/O műveletekhez (hálózat, fájlok), Schedulers.computation() CPU-intenzív feladatokhoz, Schedulers.newThread() minden új szálhoz, Schedulers.single() egyszálú végrehajtáshoz és Schedulers.trampoline() azonnali végrehajtáshoz az aktuális szálon.
subscribeOn határozza meg, hogy melyik Scheduler-en hajtódik végre az Observable forrás. Ha több subscribeOn van a láncban — a forráshoz legközelebbi élvez elsőbbséget. observeOn átkapcsolja a downstream-et a megadott Scheduler-re — minden observeOn használat megváltoztatja a szálat a következő operátorok számára. Tipikus Android minta: subscribeOn(Schedulers.io()) hálózati munkához, observeOn(AndroidSchedulers.mainThread()) UI frissítéséhez.
// Többszálú feldolgozás kontextusváltással
Observable.fromCallable(() -> database.getItems())
.subscribeOn(Schedulers.io()) // Adatbázis io-n
.map(items -> processItems(items)) // transzformáció io-n
.observeOn(Schedulers.computation()) // átkapcsolunk computation-re
.map(processed -> compressImages(processed))
.observeOn(AndroidSchedulers.mainThread())
.subscribe(result -> ui.showResult(result))
AndroidSchedulers.mainThread() egy Scheduler az RxAndroid könyvtárból, amely az Android főszálán hajtja végre a kódot. Kötelező bármilyen UI-frissítéshez a reaktív láncban. A könyvtár belsőleg Handler-t használ, és garantálja a végrehajtást a UI szálon még magas terhelés mellett is. Háttérműveletekhez a Schedulers.io() korlátlan szálpool-t támogat, és bármilyen blokkoló művelethez alkalmas. A Schedulers.computation() rögzített pool-t használ, amely megegyezik a processzormagok számával.
RxJava Androidban három fő forgatókönyvhöz használatos: reaktív lekérdezések Room-hoz, integráció Retrofit-tel és a UI reaktív kötése RxBinding-en keresztül. Minden forgatókönyvre jellemző a saját típuskészlete: Room Flowable-t ad vissza megfigyelhető lekérdezésekhez, Retrofit — Single-t HTTP-kérésekhez, RxBinding — Observable-t UI-eseményekhez.
Room egy adatperzisztencia könyvtár a Google-tól. A Room 2.1-től kezdve az adatbázis támogatja a reaktív visszatérési típusokat: Flowable és Observable. Amikor bármely rekord megváltozik a táblában, a Room automatikusan új értéket küld a folyamba. A fejlesztő feliratkozik a Flowable-re a ViewModel-ben, és naprakész adatokat kap kézi lekérdezések nélkül minden változáskor.
// Room DAO reaktív lekérdezéssel
@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 kompozíció
class UserViewModel(private val dao: UserDao) : ViewModel() {
val users: Flowable<List<User>> = dao.getAllUsers()
.subscribeOn(Schedulers.io())
}
Az MVVM + RxJava minta azon alapul, hogy a ViewModel-nek nincs hivatkozása a View-ra. A ViewModel reaktív forrásokat tesz közzé (Flowable, LiveData Transformations-en keresztül), és az Activity vagy Fragment feliratkozik rájuk. Ez tesztelhetőséget biztosít: a ViewModel UI nélkül tesztelhető, a Scheduler RxJavaPlugins.setComputationScheduler-en keresztüli helyettesítésével. A CompositeDisposable a ViewModel-ben kezeli a feliratkozások életciklusát — az onCleared() hívásakor minden feliratkozás megszűnik.
Kotlin Flow a hideg folyamatok natív implementációja Kotlin-ban, beépítve a korutinokba és bevezetve a Kotlin 1.3-ban. A Flow ugyanazokat a feladatokat oldja meg, mint az RxJava, de alapvető különbségekkel: beépített korutin támogatás (suspend függvények), megszakítás coroutine cancellation-en keresztül, és a backpressure problémáinak hiánya — a Flow suspend-et használ a pufferelés helyett. A Flow a Kotlin szabványos könyvtárának része, nem igényel további függőségeket.
RxJava továbbra is előnyben részesített választás Java projektekhez, Java 7-8 támogatással rendelkező projektekhez és meglévő RxJava kódbázisokhoz. Az RxJava ökoszisztéma jelentősen gazdagabb: >400 operátor a Flow ~50 operátorával szemben, integráció Retrofit-tel a beépített CallAdapter-en keresztül, backpressure támogatás Flowable-n keresztül, valamint az RxBinding, RxPermissions, RxLocation elérhetősége Androidhoz. A Kotlin Flow gyorsan felzárkózik, de az RxJava rugalmassága összetett folyamkombinációs forgatókönyvekben még mindig magasabb.
| Jellemző | RxJava | Kotlin Flow |
|---|---|---|
| Nyelv | Java / Kotlin | Csak Kotlin |
| Megszakítás | Disposable / CompositeDisposable | Coroutine cancellation |
| Backpressure | Flowable (BUFFER, DROP, LATEST stratégiák) | Conflate / buffer segítségével |
| Operátorok | 400+ | ~50 (bővíthető) |
| Room integráció | Flowable, Observable | Flow, StateFlow |
| ViewModel | CompositeDisposable | viewModelScope + Flow |
Gyakran Ismételt Kérdések
Observable nem támogatja a backpressure-t — ha a producer gyorsabb, mint a consumer, az események felhalmozódnak a memóriában. Flowable megvalósítja a Reactive Streams-t backpressure-rel a Subscription.request()-en keresztül, ami megakadályozza a puffer túlcsordulását sebességeltérés esetén.
Single olyan műveletekhez használatos, amelyek pontosan egy értéket vagy hibát adnak vissza: HTTP-kérések, egy rekord olvasása az adatbázisból, eredmény kiszámítása. A Single szemantikailag megfelel a Future-nak, és lerövidíti a kódot a nem használt onComplete eltávolításával.
A dispose() metódus a Disposable-on megszünteti a feliratkozást. Csoportos kezeléshez CompositeDisposable használatos — összegyűjti az összes Disposable-t, és egyidejűleg megszünteti őket a clear() hívásakor. Tipikus hely — onCleared() a ViewModel-ben vagy onPause() az Activity-ben.
flatMap feliratkozik az összes beágyazott Observable-re, és tetszőleges sorrendben kombinálja az eseményeiket. switchMap új elem érkezésekor leiratkozik az előző Observable-ről, és feliratkozik az újra. A switchMap kereséshez használatos — minden új kérés megszakítja az előzőt.
Új Kotlin projektekhez a Flow előnyösebb a korutinokkal való integráció és a kisebb méret miatt. Meglévő RxJava projekteknél a migráció csak akkor indokolt, ha a teljes kódbázis áttér korutinokra — mindkét könyvtár köztes használata bonyolítja az architektúrát.
Összefoglalás
Kulcsrakész mobilalkalmazást fejlesztünk
Az IT Sectr 2017 óta készít iOS és Android alkalmazásokat induló vállalkozásoknak és vállalkozásoknak. Tanácsot adunk, és a legjobb megoldást javasoljuk.
Olvassa el is