RxJava: baze, ReactiveX și lucrul cu fluxuri de date

Autor: IT Sectr Publicat: 2026-03-16 Timp de citire: 8 min

RxJava este o bibliotecă de programare reactivă pentru Java și Android, care implementează modelul Observer prin Observable și Observer. Potrivit ReactiveX GitHub, 2026, RxJava permite procesarea fluxurilor de date asincrone și a evenimentelor cu ajutorul lanțurilor de operatori. Unitatea de bază este Observable, care emite date către Observer printr-un lanț de transformări. RxJava 3 este versiunea stabilă curentă cu suport pentru Java 8 lambda, Reactive Streams și integrare cu Android prin RxAndroid.

Principalele puncte

  • RxJava — implementarea Java a ReactiveX pentru procesarea asincronă a fluxurilor de date
  • Observable — sursa de date care emite elemente către Observer
  • Observer — abonatul care primește notificări onNext, onError și onComplete
  • Operatori — lanț de funcții pentru transformarea, filtrarea și combinarea fluxurilor
  • Schedulers — componenta pentru gestionarea firelor de execuție ale Observable și Observer

Ce este RxJava și ReactiveX

RxJava — implementarea Java a specificației ReactiveX, o bibliotecă pentru programare asincronă folosind fluxuri observabile (Observable). RxJava 2 a fost lansată în 2016 cu suport pentru Reactive Streams (Flowable) și separarea în rx.Observable și io.reactivex.Observable. RxJava 3 (2019) — versiunea majoră curentă cu compatibilitate inversă cu RxJava 2.

Ideea principală a RxJava — totul este un flux: flux de date, flux de evenimente, flux de stări. Orice operație asincronă poate fi reprezentată ca un Observable care emite date, o eroare sau un semnal de finalizare. Observer se abonează la Observable și primește notificări în timp real.

Conform datelor Badoo (2024), înainte de trecerea la corutine, 76% din aplicațiile Android din top-200 Google Play foloseau RxJava pentru operații asincrone. Acum ponderea scade în favoarea corutinelor, dar RxJava rămâne în codul de producție a miilor de aplicații și este considerată o tehnologie matură și testată. ReactiveX — o specificație cross-platformă, implementată și pentru JavaScript (RxJS), .NET (Rx.NET), Swift (RxSwift) și alte limbaje.

Modelul Observer în RxJava

ReactiveX extinde modelul clasic Observer cu două mecanisme: lanț de operatori (operator chaining) și gestionarea firelor (schedulers). Observable nu începe să emită date până când Observer nu se abonează (evaluare leneșă). Acest lucru permite construirea unui pipeline de date care se activează doar la prezența unei abonări.

Tipuri de Observable: Observable, Flowable, Single, Maybe, Completable

Observable — tipul de bază care emite 0..N elemente cu onError sau onComplete. Potrivit pentru fluxuri de date de lungime nelimitată — de exemplu, evenimente de click sau actualizări de geolocație. Observable nu suportă backpressure.

Flowable — versiunea Reactive Streams a Observable cu suport pentru backpressure. Se folosește când sursa de date poate genera elemente mai repede decât Observer poate procesa. Flowable suportă strategiile BACKPRESSURE_BUFFER, DROP, LATEST și ERROR.

TipElementeBackpressureUtilizare
Observable0..NNuEvenimente UI, fluxuri mici
Flowable0..NDaDate mari, timp real
Single1 (onSuccess/onError)Răspuns unic (rețea)
Maybe0..1Valoare opțională (cache)
Completable0 (onComplete/onError)Operație fără date (scriere)

Single, Maybe și Completable

Single emite exact un element sau o eroare — ideal pentru cereri de rețea. Maybe — 0 sau 1 element, potrivit pentru cache unde datele pot lipsi. Completable — doar onComplete sau onError, fără date, convenabil pentru operații de scriere sau ștergere. Aceste tipuri simplifică API-ul, restrângând contractul la un caz specific. Retrofit (clientul HTTP popular pentru Android) suportă toate cele cinci tipuri RxJava direct, permițând alegerea celui mai potrivit tip de returnare pentru fiecare endpoint fără îmbrăcăminte suplimentară.

Operatori RxJava: transformarea și filtrarea fluxurilor

Operatori sunt funcții care transformă un Observable în altul. Lanțul de operatori (operator chain) descrie pipeline-ul de date: fiecare operator preia fluxul de la anteriorul, îl transformă și îl transmite următorului. RxJava conține peste 200 de operatori împărțiți pe categorii.

  • map — transformă fiecare element (Integer → String)
  • flatMap — transformă elementul într-un Observable și combină toate într-un singur flux
  • filter — lasă să treacă elementele după o condiție
  • zip — combină elementele a N Observable după index
  • merge — combină mai multe Observable într-unul singur, păstrând ordinea temporală
  • debounce — lasă să treacă elementele dacă între ele este mai puțin decât intervalul specificat

flatMap — unul dintre cei mai puternici operatori RxJava. Permite executarea unei cereri asincrone pentru fiecare element și colectarea rezultatelor într-un flux comun. De exemplu, flatMap este folosit pentru încărcarea detaliilor după o listă de ID-uri: fiecare ID → cerere de rețea → combinarea rezultatelor. Spre deosebire de map, care pur și simplu transformă un element, flatMap poate emite mai multe elemente sau poate comuta la un alt Observable, ceea ce îl face fundamentul construirii pipeline-urilor asincrone.

Gestionarea erorilor prin operatori

onErrorResumeNext — la eroare comută la un Observable de rezervă. retry — reîncearcă abonarea la eroare de N ori. onErrorReturn — returnează o valoare implicită în locul erorii. doOnError — execută o acțiune secundară la eroare fără a modifica fluxul (logare sau analitică). Combinarea acestor operatori permite construirea de pipeline-uri fiabile cu o strategie clară de gestionare a defecțiunilor fără try/catch manual.

Schedulers: gestionarea firelor în RxJava

Schedulers determină pe ce fir se execută Observable și Observer. subscribeOn setează firul pentru sursă, observeOn — firul pentru Observer și operatorii următori. Această separare — avantajul cheie al RxJava: sursa pe firul IO, procesarea pe computation, UI — pe firul principal.

Schedulers principale: Schedulers.io() — pentru operații I/O (rețea, disc), pool nelimitat. Schedulers.computation() — pentru calcule, pool fix după numărul de nuclee. Schedulers.newThread() — fir nou pentru fiecare sarcină. AndroidSchedulers.mainThread() — firul principal Android (RxAndroid). Există și Schedulers.trampoline() pentru executarea sarcinilor pe firul curent cu coadă FIFO, util pentru teste.

Conform datelor Google (2025), utilizarea corectă a Schedulers este cel mai dificil lucru în RxJava pentru începători. Greșeala tipică — apelarea subscribeOn după observeOn, care nu afectează sursa. subscribeOn trebuie să fie primul în lanț pentru sursă, observeOn — înainte de abonarea UI. Regula: subscribeOn afectează doar upstream-ul (sursa), observeOn comută downstream-ul (abonatul și toți operatorii după el).

Exemple de cod cu RxJava în Android

Să analizăm trei scenarii: cerere de rețea cu Single, cereri paralele cu zip și debounce pentru câmpul de căutare cu debounce.

Cerere de rețea cu Single

Single este ideal pentru cererile Retrofit: o cerere — un răspuns. Abonarea pe firul principal pentru actualizarea UI.

java
api.getUser(id)
    .subscribeOn(Schedulers.io())
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(new SingleObserver<User>() {
        @Override
        public void onSuccess(User user) { showUser(user); }
        @Override
        public void onError(Throwable e) { showError(e); }
    })

Cereri paralele cu zip

zip combină rezultatele a două Single independente într-unul singur. Se execută paralel, rezultatul — după finalizarea ambelor.

java
Single.zip(
    api.getProfile(),
    api.getSettings(),
    (profile, settings) -> new Dashboard(profile, settings)
)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(dashboard -> showDashboard(dashboard), e -> logError(e))

Debounce pentru câmpul de căutare

debounce ignoră schimbările rapide ale textului și trimite cererea doar după o pauză de 400 ms. distinctUntilChanged anulează cererea dacă textul nu s-a schimbat.

java
RxTextView.textChanges(searchView)
    .debounce(400, TimeUnit.MILLISECONDS)
    .filter(text -> text.length() >= 3)
    .distinctUntilChanged()
    .switchMap(query -> api.search(query))
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(results -> showResults(results))

RxJava vs Kotlin Coroutines: compararea abordărilor

RxJava și Kotlin Coroutines rezolvă aceeași sarcină — programarea asincronă — dar cu abordări fundamental diferite. RxJava este construit pe modelul Observer și este push-based: sursa trimite date, Observer reacționează. Corutinele — pull-based: codul obține date secvențial prin await.

  • RxJava — reactiv, flux de date, >200 operatori, push-based, curbă de învățare dificilă
  • Coroutines — secvențial, suspend/await, ~40 funcții, pull-based, sintaxă simplă
  • RxJava — matur (2016), ecosistem vast, dar curbă de învățare dificilă
  • Coroutines — modern (2018), alegerea preferată Google pentru cod nou
  • RxJava — backpressure din cutie prin Flowable, strategii de buffer bine puse la punct
  • Coroutines — Flow cu backpressure recent, dar dezvoltat activ de JetBrains

Conform Google I/O 2024, Kotlin Coroutines este abordarea recomandată pentru codul asincron nou în Android. RxJava rămâne suportat pentru proiectele existente. Google oferă biblioteci de tranziție (kotlinx-coroutines-rx3) pentru migrare treptată. AndroidX (LiveData, Room, Paging 3) suportă ambele abordări, permițând utilizarea RxJava în modulele vechi și a corutinelor în modulele noi fără conflicte de dependențe.

Strategia de migrare de la RxJava la corutine

Trecerea treptată: fiecare componentă nouă se scrie cu corutine, codul RxJava vechi nu se atinge. RxJava → corutine prin awaitSingle() sau awaitFirst(). Corutine → RxJava prin future() sau asFlowable(). Migrarea completă durează 6–18 luni pentru proiecte mari.

Întrebări frecvente

Cu ce diferă Observable de Flowable?

Observable nu suportă backpressure — dacă sursa generează date mai repede decât procesorul, apare MissingBackpressureException. Flowable suportă Reactive Streams backpressure cu o strategie de buffer configurabilă.

Ce sunt subscribeOn și observeOn?

subscribeOn setează Scheduler-ul pentru executarea sursei Observable. observeOn setează Scheduler-ul pentru Observer și toți operatorii următori din lanț. subscribeOn afectează upstream-ul, observeOn — downstream-ul.

Merită să trec de la RxJava la corutine?

Pentru proiecte noi — da, Google recomandă corutinele. Pentru proiecte existente — migrare treptată prin kotlinx-coroutines-rx3. RxJava rămâne stabil și suportat pentru codul vechi.

Cum se gestionează erorile în RxJava?

Prin operatori: onErrorReturn (valoare implicită), onErrorResumeNext (Observable de rezervă), retry (reîncearcă de N ori). Sau prin Observer.onError() pentru afișarea utilizatorului.

Ce este CompositeDisposable?

CompositeDisposable — un container pentru gestionarea mai multor abonamente. La dispose() toate abonamentele adăugate sunt anulate. Se folosește în Activity/Fragment pentru anularea tuturor cererilor la distrugerea ecranului.

Rezumat

  • RxJava — bibliotecă de programare reactivă pentru Java și Android bazată pe modelul Observer
  • Observable/Flowable — surse de date cu și fără suport backpressure respectiv
  • Single, Maybe, Completable — tipuri specializate pentru 1, 0..1 și 0 elemente
  • Operatori (map, flatMap, zip, filter) — lanț de transformări cu peste 200 de funcții
  • Schedulers — subscribeOn pentru sursă și observeOn pentru consumatorul de date
  • RxJava vs Coroutines — corutinele recomandate de Google pentru cod nou, RxJava pentru legacy
  • CompositeDisposable — gestionarea sigură a abonamentelor cu anulare la distrugerea ecranului

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.

Discutați proiectul

Citiți și