RxJava este o bibliotecă de programare reactivă pentru JVM care implementează fluxuri de date asincrone prin modelul Observable cu operatori funcționali de transformare. Portează conceptele ReactiveX în Java și Kotlin, oferind o API unificată pentru lucrul cu cereri de rețea, baze de date, evenimente UI și sarcini de fundal. Conform datelor ReactiveX, 2025, biblioteca este utilizată în peste 120.000 de proiecte pe GitHub și este standardul programării reactive pentru Android până la apariția Kotlin Flow. RxJava înlocuiește AsyncTask, Loader și callback-urile cu un lanț unificat de procesare a datelor.
Principalele puncte
RxJava este o implementare a bibliotecii ReactiveX (Reactive Extensions) pentru mașina virtuală Java. Prima versiune RxJava a fost lansată de compania Netflix în 2013 pentru gestionarea apelurilor asincrone în aplicațiile server. La momentul creării, alternativa principală în Java era Future și Callback — ambele abordări duceau la callback-hell și gestionarea complexă a firelor. RxJava a propus compoziția operațiilor asincrone prin Observable cu lanțuri de operatori funcționali.
Arhitectura RxJava se bazează pe specificația Reactive Streams — un standard pentru procesarea asincronă a fluxurilor cu backpressure neblocant. Specificația definește patru interfețe: Publisher, Subscriber, Subscription și Processor. RxJava 2+ implementează complet Reactive Streams prin tipul Flowable, respectând contractele backpressure spre deosebire de RxJava 1. Observable în RxJava 2 nu suportă backpressure — este destinat fluxurilor cu un număr mic de evenimente sau evenimentelor UI.
Conform sondajului JetBrains, 2025, RxJava se află în top-3 biblioteci pentru dezvoltarea Android. Scenariile principale de utilizare: procesarea cererilor de rețea prin Retrofit (integrat cu RxJava prin CallAdapter), lucrul cu Room (interogările reactive returnează Flowable sau Maybe), animații și evenimente UI prin RxBinding și căutarea debounce la introducerea textului. Toate aceste scenarii sunt unite de un lanț de același tip: sursă (Observable) → transformare (operatori) → abonare (subscribe).
RxJava 1 (2013) a pus bazele conceptului Observable și operatorilor, dar suferea de probleme cu backpressure — în fluxurile rapide, datele se acumulau în memorie, cauzând OutOfMemoryError. RxJava 2 (2016) a reparat arhitectura, separând Observable (fără backpressure) și Flowable (cu backpressure). RxJava 3 (2020) a adăugat suport pentru Java 8 Stream API, operatori suplimentari și performanță îmbunătățită la abonare. În prezent, RxJava 3 este versiunea recomandată pentru proiecte noi.
RxJava oferă cinci tipuri principale de surse reactive, fiecare orientat către un scenariu specific. Observable și Flowable emit valori multiple, Single — o singură valoare sau eroare, Completable — doar faptul finalizării fără date, Maybe — o valoare, zero sau eroare. Alegerea tipului corect reduce cantitatea de cod și face lanțul auto-documentat.
| Tip | Număr de evenimente | Backpressure | Scenariu |
|---|---|---|---|
| Observable | 0..N, apoi finalizare | Nu | Evenimente UI, fluxuri scurte |
| Flowable | 0..N, apoi finalizare | Da | Răspunsuri de rețea, fluxuri din BD |
| Single | Exact 1 sau eroare | Nu | Cerere HTTP, citirea unei înregistrări |
| Completable | 0 (doar finalizare) | Nu | Scriere în BD, trimitere eveniment |
| Maybe | 0, 1 sau eroare | Nu | Cache: există valoare sau nu |
Flowable este cel mai flexibil tip pentru lucrul cu fluxuri mari de date. Implementează Reactive Streams Publisher cu suport pentru backpressure: consumer poate solicita un număr specific de elemente prin Subscription.request(n). Aceasta previne depășirea buffer-ului la nepotrivirea vitezei dintre producer și consumer. Dacă backpressure nu este critic — utilizați Observable, care are o supraîncărcare mai mică datorită absenței mecanismului request.
Single este alegerea optimă pentru cererile HTTP. Retrofit 2 cu RxJava CallAdapter returnează Single<ResponseBody> pentru fiecare cerere. Single garantează exact un apel onSuccess sau onError, ceea ce corespunde semanticii unei cereri HTTP — un răspuns sau o eroare. Completable este utilizat pentru operații de scriere care nu returnează date: insert, update, delete. Maybe este convenabil la verificarea cache-ului — poate returna o valoare, poate nu.
// Exemplu de utilizare a Single pentru o cerere HTTP
interface ApiService {
@GET("users/{id}")
fun getUser(@Path("id") userId: Int): Single<User>
}
// Abonare cu procesare pe firul principal
apiService.getUser(42)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe({ user ->
textView.text = user.name
}, { error ->
Log.e("API", "Error: ${error.message}")
})
.addTo(compositeDisposable)
Operatorii RxJava sunt funcții de ordin superior care primesc o sursă reactivă și returnează alta, transformând fluxul de date. RxJava 3 conține peste 400 de operatori împărțiți pe categorii: transformare, filtrare, combinare, gestionarea erorilor și gestionarea timpului. Fiecare operator este leneș — lanțul se construiește la declarare și se execută la abonare.
map este operatorul de bază care transformă fiecare valoare printr-o funcție. flatMap primește o funcție care returnează un Observable pentru fiecare element și desface rezultatul într-un flux unic. switchMap este similar cu flatMap, dar la primirea unui element nou se dezabonează de la Observable-ul anterior. concatMap păstrează ordinea elementelor — spre deosebire de flatMap, se abonează secvențial la fiecare Observable imbricat.
// Parsare JSON cu transformare și filtrare
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("Eroare", it.message) })
Combinarea fluxurilor este domeniul în care RxJava este deosebit de puternic. zip combină elemente din mai multe Observable-uri în perechi după index: primul cu primul, al doilea cu al doilea. combineLatest emite o nouă valoare la modificarea oricărui flux, combinând ultimele valori ale tuturor fluxurilor. merge combină mai multe Observable-uri într-unul singur, păstrând ordinea de sosire a evenimentelor. concat se abonează secvențial la fiecare Observable și transmite toate evenimentele sale înainte de a trece la următorul.
Gestionarea timpului include debounce (așteptarea unei pauze în flux înainte de trimitere), throttleFirst (trecerea primului eveniment, ignorarea restului în fereastră), timeout (eroare dacă evenimentul nu a sosit în interval). Căutarea debounce la introducerea textului este cel mai comun scenariu: searchObservable.debounce(300, MILLISECONDS).distinctUntilChanged() previne cererile inutile la tastarea rapidă.
| Categorie | Operator | Comportament |
|---|---|---|
| Transformare | map / flatMap / switchMap | Transformarea unei singure valori sau a unui flux |
| Filtrare | filter / distinct / take | Selectarea valorilor după condiție |
| Combinare | zip / combineLatest / merge | Unirea a 2+ fluxuri |
| Erori | onErrorResumeNext / retry | Recuperare după defect |
| Utilități | delay / timeout / debounce | Gestionarea timpului în flux |
Scheduler în RxJava este o abstracție peste un pool de fire. Biblioteca oferă cinci Scheduler-uri încorporate: Schedulers.io() pentru operații I/O (rețea, fișiere), Schedulers.computation() pentru sarcini intensive CPU, Schedulers.newThread() pentru fiecare fir nou, Schedulers.single() pentru execuție mono-fir și Schedulers.trampoline() pentru executare imediată în firul curent.
subscribeOn determină pe ce Scheduler se execută sursa Observable. Dacă în lanț sunt mai multe subscribeOn — prioritatea o are cel mai apropiat de sursă. observeOn comută downstream-ul pe Scheduler-ul specificat — fiecare utilizare a observeOn schimbă firul pentru operatorii următori. Un tipar tipic Android: subscribeOn(Schedulers.io()) pentru lucrul cu rețeaua, observeOn(AndroidSchedulers.mainThread()) pentru actualizarea UI.
// Procesare multi-thread cu comutare de context
Observable.fromCallable(() -> database.getItems())
.subscribeOn(Schedulers.io()) // BD pe io
.map(items -> processItems(items)) // transformare pe io
.observeOn(Schedulers.computation()) // comutăm pe computation
.map(processed -> compressImages(processed))
.observeOn(AndroidSchedulers.mainThread())
.subscribe(result -> ui.showResult(result))
AndroidSchedulers.mainThread() este un Scheduler din biblioteca RxAndroid care execută codul pe firul principal Android. Este obligatoriu pentru orice actualizare UI în lanțul reactiv. Biblioteca utilizează Handler intern și garantează execuția în firul UI chiar și la sarcină ridicată. Pentru operații de fundal, Schedulers.io() suportă un pool nelimitat de fire și este potrivit pentru orice operații blocante. Schedulers.computation() utilizează un pool fix, egal cu numărul de nuclee ale procesorului.
RxJava în Android este utilizat pentru trei scenarii principale: interogări reactive către Room, integrare cu Retrofit și legarea reactivă a UI prin RxBinding. Fiecare scenariu are propriul set de tipuri: Room returnează Flowable pentru interogări observabile, Retrofit — Single pentru cereri HTTP, RxBinding — Observable pentru evenimente UI.
Room este o bibliotecă de persistență a datelor de la Google. Începând cu Room 2.1, baza de date suportă tipuri de returnare reactive: Flowable și Observable. La modificarea oricărei înregistrări în tabel, Room trimite automat o nouă valoare în flux. Dezvoltatorul se abonează la Flowable în ViewModel și primește date actualizate fără interogări manuale la fiecare modificare.
// Room DAO cu interogare reactivă
@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 — compoziție Room + Network
class UserViewModel(private val dao: UserDao) : ViewModel() {
val users: Flowable<List<User>> = dao.getAllUsers()
.subscribeOn(Schedulers.io())
}
Tiparul MVVM + RxJava se bazează pe faptul că ViewModel nu are referințe la View. ViewModel publică surse reactive (Flowable, LiveData prin Transformations), iar Activity sau Fragment se abonează la ele. Aceasta asigură testabilitate: ViewModel se testează fără UI, înlocuind Scheduler-ul prin RxJavaPlugins.setComputationScheduler. CompositeDisposable în ViewModel gestionează ciclul de viață al abonărilor — la onCleared() toate abonările sunt anulate.
Kotlin Flow este o implementare nativă a fluxurilor reci în Kotlin, integrată în corutine și introdusă în Kotlin 1.3. Flow rezolvă aceleași sarcini ca și RxJava, dar cu diferențe fundamentale: suport încorporat pentru corutine (funcții suspend), anulare prin coroutine cancellation și absența problemelor cu backpressure — Flow folosește suspend în loc de bufferizare. Flow face parte din biblioteca standard Kotlin, nefiind necesare dependențe suplimentare.
RxJava rămâne alegerea preferată pentru proiectele în Java, proiectele cu suport Java 7-8 și bazele de cod existente pe RxJava. Ecosistemul RxJava este semnificativ mai bogat: >400 de operatori față de ~50 în Flow, integrare cu Retrofit prin CallAdapter încorporat, suport pentru backpressure prin Flowable și existența RxBinding, RxPermissions, RxLocation pentru Android. Kotlin Flow recuperează rapid, dar flexibilitatea RxJava în scenarii complexe de combinare a fluxurilor este încă superioară.
| Caracteristică | RxJava | Kotlin Flow |
|---|---|---|
| Limbaj | Java / Kotlin | Doar Kotlin |
| Anulare | Disposable / CompositeDisposable | Coroutine cancellation |
| Backpressure | Flowable (strategii BUFFER, DROP, LATEST) | Prin conflate / buffer |
| Operatori | 400+ | ~50 (extensibil) |
| Integrare Room | Flowable, Observable | Flow, StateFlow |
| ViewModel | CompositeDisposable | viewModelScope + Flow |
Întrebări frecvente
Observable nu suportă backpressure — dacă producer este mai rapid decât consumer, evenimentele se acumulează în memorie. Flowable implementează Reactive Streams cu backpressure prin Subscription.request(), ceea ce previne depășirea buffer-ului la nepotrivirea vitezelor.
Single este utilizat pentru operații care returnează exact o valoare sau o eroare: cereri HTTP, citirea unei singure înregistrări din BD, calcularea unui rezultat. Single corespunde semantic lui Future și scurtează codul eliminând onComplete neutilizat.
Metoda dispose() pe Disposable anulează abonarea. Pentru gestionarea în grup se utilizează CompositeDisposable — colectează toate Disposable-urile și le anulează simultan la apelarea clear(). Loc tipic — onCleared() în ViewModel sau onPause() în Activity.
flatMap se abonează la toate Observable-urile imbricate și combină evenimentele lor în ordine arbitrară. switchMap la primirea unui element nou se dezabonează de la Observable-ul anterior și se abonează la cel nou. switchMap este utilizat la căutare — fiecare cerere nouă o anulează pe cea anterioară.
Pentru proiecte noi în Kotlin, Flow este preferabil datorită integrării cu corutinele și dimensiunii mai mici. Pentru proiecte existente pe RxJava, migrarea este justificată doar dacă întreaga bază de cod trece la corutine — utilizarea intermediară a ambelor biblioteci complică arhitectura.
Concluzii
Vom dezvolta o aplicație mobilă la cheie
IT Sectr creează aplicații iOS și Android pentru startup-uri și afaceri din 2017. Vă vom consilia și vă vom propune cea mai bună soluție.
Citiți și