RxJava: основи, ReactiveX и работа с потоци от данни

Автор: IT Sectr Публикувано: 2026-03-16 Време за четене: 8 мин

RxJava е библиотека за реактивно програмиране за Java и Android, която имплементира модела Observer чрез Observable и Observer. Според ReactiveX GitHub, 2026, RxJava позволява обработка на асинхронни потоци от данни и събития с помощта на вериги от оператори. Основната единица е Observable, който излъчва данни към Observer чрез верига от трансформации. RxJava 3 е текущата стабилна версия с поддръжка на Java 8 lambda, Reactive Streams и интеграция с Android чрез RxAndroid.

Основни точки

  • RxJava — Java имплементация на ReactiveX за асинхронна обработка на потоци от данни
  • Observable — източник на данни, който излъчва елементи към Observer
  • Observer — абонат, получаващ onNext, onError и onComplete известия
  • Оператори — верига от функции за трансформация, филтрация и комбиниране на потоци
  • Schedulers — компонент за управление на нишките за изпълнение на Observable и Observer

Какво са RxJava и ReactiveX

RxJava — Java имплементация на спецификацията ReactiveX, библиотека за асинхронно програмиране чрез наблюдаеми потоци (Observable). RxJava 2 беше пусната през 2016 г. с поддръжка на Reactive Streams (Flowable) и разделение на rx.Observable и io.reactivex.Observable. RxJava 3 (2019) — текущата основна версия с обратна съвместимост с RxJava 2.

Основната идея на RxJava — всичко е поток: поток от данни, поток от събития, поток от състояния. Всяка асинхронна операция може да бъде представена като Observable, който излъчва данни, грешка или сигнал за завършване. Observer се абонира за Observable и получава известия в реално време.

Според данни на Badoo (2024), преди прехода към корутини, 76% от Android приложенията от топ-200 на Google Play използваха RxJava за асинхронни операции. Сега делът намалява в полза на корутините, но RxJava остава в производствения код на хиляди приложения и се счита за зряла, доказана технология. ReactiveX — междуплатформена спецификация, имплементирана също за JavaScript (RxJS), .NET (Rx.NET), Swift (RxSwift) и други езици.

Моделът Observer в RxJava

ReactiveX разширява класическия модел Observer с два механизма: верига от оператори (operator chaining) и управление на нишки (schedulers). Observable не започва да излъчва данни, докато Observer не се абонира (мързелива оценка). Това позволява изграждане на pipeline от данни, който се активира само при наличие на абонамент.

Типове Observable: Observable, Flowable, Single, Maybe, Completable

Observable — основен тип, който излъчва 0..N елемента с onError или onComplete. Подходящ за потоци от данни с неограничена дължина — например събития на кликвания или актуализации на геолокация. Observable не поддържа backpressure.

Flowable — версията Reactive Streams на Observable с поддръжка на backpressure. Използва се, когато източникът на данни може да генерира елементи по-бързо, отколкото Observer успява да обработва. Flowable поддържа стратегиите BACKPRESSURE_BUFFER, DROP, LATEST и ERROR.

ТипЕлементиBackpressureПриложение
Observable0..NНеUI събития, малки потоци
Flowable0..NДаГолеми данни, реално време
Single1 (onSuccess/onError)Единичен отговор (мрежа)
Maybe0..1Опционална стойност (кеш)
Completable0 (onComplete/onError)Операция без данни (запис)

Single, Maybe и Completable

Single излъчва точно един елемент или грешка — идеален за мрежови заявки. Maybe — 0 или 1 елемент, подходящ за кеш, където данните може да липсват. Completable — само onComplete или onError, без данни, удобен за операции за запис или изтриване. Тези типове опростяват API, стеснявайки договора до конкретен случай. Retrofit (популярният HTTP клиент за Android) поддържа и петте типа RxJava директно, позволявайки избор на най-подходящия тип за връщане за всяка крайна точка без излишна обвивка.

Оператори на RxJava: трансформация и филтрация на потоци

Операторите са функции, които преобразуват един Observable в друг. Веригата от оператори (operator chain) описва pipeline от данни: всеки оператор приема потока от предишния, трансформира го и го предава на следващия. RxJava съдържа над 200 оператора, разделени в категории.

  • map — преобразува всеки елемент (Integer → String)
  • flatMap — преобразува елемента в Observable и обединява всички в един поток
  • filter — пропуска елементи според условие
  • zip — комбинира елементи на N Observable по индекс
  • merge — обединява няколко Observable в един, запазвайки времевия ред
  • debounce — пропуска елементи, ако интервалът между тях е по-малък от зададения

flatMap — един от най-мощните оператори на RxJava. Позволява изпълнение на асинхронна заявка за всеки елемент и събиране на резултатите в общ поток. Например flatMap се използва за зареждане на детайли по списък с ID: всяко ID → мрежова заявка → обединяване на резултатите. За разлика от map, който просто преобразува елемент, flatMap може да излъчва множество елементи или да превключи към друг Observable, което го прави основа за изграждане на асинхронни pipeline-и.

Управление на грешки чрез оператори

onErrorResumeNext — при грешка превключва към резервен Observable. retry — повтаря абонамента при грешка N пъти. onErrorReturn — връща стойност по подразбиране вместо грешка. doOnError — изпълнява странично действие при грешка без да променя потока (логване или аналитика). Комбинирането на тези оператори позволява изграждане на надеждни pipeline-и с ясна стратегия за обработка на откази без ръчен try/catch.

Schedulers: управление на нишки в RxJava

Schedulers определят на коя нишка се изпълняват Observable и Observer. subscribeOn задава нишката за източника, observeOn — нишката за Observer и следващите оператори. Това разделение — ключовото предимство на RxJava: източник на IO нишка, обработка на computation, UI — на основната нишка.

Основни Schedulers: Schedulers.io() — за I/O операции (мрежа, диск), неограничен пул. Schedulers.computation() — за изчисления, фиксиран пул според броя ядра. Schedulers.newThread() — нова нишка за всяка задача. AndroidSchedulers.mainThread() — основната нишка на Android (RxAndroid). Съществува и Schedulers.trampoline() за изпълнение на задачи в текущата нишка с FIFO опашка, полезен за тестове.

Според данни на Google (2025), правилното използване на Schedulers е най-трудното в RxJava за начинаещи. Типична грешка — извикване на subscribeOn след observeOn, което не влияе на източника. subscribeOn трябва да бъде първи във веригата за източника, observeOn — преди UI абонамента. Правило: subscribeOn влияе само на upstream (източника), observeOn превключва downstream (абоната и всички оператори след него).

Примери за код с RxJava в Android

Нека разгледаме три сценария: мрежова заявка с Single, паралелни заявки с zip и debounce за полето за търсене с debounce.

Мрежова заявка с Single

Single е идеален за Retrofit заявки: една заявка — един отговор. Абонамент на основната нишка за актуализация на 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); }
    })

Паралелни заявки с zip

zip обединява резултатите от два независими Single в един. Изпълняват се паралелно, резултатът — след завършване и на двата.

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 за полето за търсене

debounce игнорира бързи промени на текста и изпраща заявката само след 400 ms пауза. distinctUntilChanged отменя заявката, ако текстът не се е променил.

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 срещу Kotlin Coroutines: сравнение на подходите

RxJava и Kotlin Coroutines решават една и съща задача — асинхронно програмиране — но с коренно различни подходи. RxJava е изграден върху модела Observer и е push-based: източникът изпраща данни, Observer реагира. Корутините — pull-based: кодът последователно получава данни чрез await.

  • RxJava — реактивен, поток от данни, >200 оператора, push-based, стръмна крива на обучение
  • Coroutines — последователен, suspend/await, ~40 функции, pull-based, прост синтаксис
  • RxJava — зрял (2016), огромна екосистема, но стръмна крива на обучение
  • Coroutines — модерен (2018), предпочитан избор на Google за нов код
  • RxJava — backpressure от кутията чрез Flowable, изпипани стратегии за буфериране
  • Coroutines — Flow с backpressure наскоро, но активно развиван от JetBrains

Според Google I/O 2024, Kotlin Coroutines са препоръчителният подход за нов асинхронен код в Android. RxJava остава поддържан за съществуващи проекти. Google предоставя преходни библиотеки (kotlinx-coroutines-rx3) за постепенна миграция. AndroidX (LiveData, Room, Paging 3) поддържа и двата подхода, позволявайки използването на RxJava в стари модули и корутини в нови модули без конфликти на зависимости.

Стратегия за миграция от RxJava към корутини

Постепенен преход: всеки нов компонент се пише с корутини, старият RxJava код не се пипа. RxJava → корутини чрез awaitSingle() или awaitFirst(). Корутини → RxJava чрез future() или asFlowable(). Пълната миграция отнема 6–18 месеца за големи проекти.

Често задавани въпроси

С какво Observable се различава от Flowable?

Observable не поддържа backpressure — ако източникът генерира данни по-бързо от процесора, възниква MissingBackpressureException. Flowable поддържа Reactive Streams backpressure с конфигурируема стратегия за буфериране.

Какво са subscribeOn и observeOn?

subscribeOn задава Scheduler за изпълнение на източника Observable. observeOn задава Scheduler за Observer и всички следващи оператори във веригата. subscribeOn влияе на upstream, observeOn — на downstream.

Струва ли си да преминем от RxJava към корутини?

За нови проекти — да, Google препоръчва корутини. За съществуващи проекти — постепенна миграция чрез kotlinx-coroutines-rx3. RxJava остава стабилен и поддържан за стар код.

Как се обработват грешки в RxJava?

Чрез оператори: onErrorReturn (стойност по подразбиране), onErrorResumeNext (резервен Observable), retry (повторение N пъти). Или чрез Observer.onError() за показване на потребителя.

Какво е CompositeDisposable?

CompositeDisposable — контейнер за управление на множество абонаменти. При dispose() всички добавени абонаменти се отменят. Използва се в Activity/Fragment за отмяна на всички заявки при унищожаване на екрана.

Резюме

  • RxJava — библиотека за реактивно програмиране за Java и Android, базирана на модела Observer
  • Observable/Flowable — източници на данни съответно с и без поддръжка на backpressure
  • Single, Maybe, Completable — специализирани типове за 1, 0..1 и 0 елемента
  • Оператори (map, flatMap, zip, filter) — верига за трансформация с над 200 функции
  • Schedulers — subscribeOn за източника и observeOn за консуматора на данни
  • RxJava срещу Coroutines — корутините се препоръчват от Google за нов код, RxJava за legacy
  • CompositeDisposable — безопасно управление на абонаменти с отмяна при унищожаване на екрана

Ще разработим мобилно приложение под ключ

IT Sectr създава iOS и Android приложения за стартъпи и бизнеси от 2017 г. Ще ви консултираме и ще предложим най-доброто решение.

Обсъдете проекта

Прочетете също