RxJava е библиотека за реактивно програмиране за JVM, която имплементира асинхронни потоци от данни чрез модела Observable с функционални оператори за трансформация. Тя пренася концепциите на ReactiveX в Java и Kotlin, предоставяйки единен API за работа с мрежови заявки, бази данни, UI събития и фонови задачи. Според данни на ReactiveX, 2025, библиотеката се използва в над 120 000 проекта в GitHub и е стандартът за реактивно програмиране за Android до появата на Kotlin Flow. RxJava заменя AsyncTask, Loader и callback-овете с единна верига за обработка на данни.
Основни точки
RxJava е имплементация на библиотеката ReactiveX (Reactive Extensions) за виртуалната машина на Java. Първата версия на RxJava беше пусната от компанията Netflix през 2013 г. за управление на асинхронни извиквания в сървърни приложения. По време на създаването основната алтернатива в Java бяха Future и Callback — и двата подхода водеха до callback-hell и сложно управление на нишки. RxJava предложи композиция на асинхронни операции чрез Observable с вериги от функционални оператори.
Архитектурата на RxJava се основава на спецификацията Reactive Streams — стандарт за асинхронна обработка на потоци с неблокиращ backpressure. Спецификацията дефинира четири интерфейса: Publisher, Subscriber, Subscription и Processor. RxJava 2+ напълно имплементира Reactive Streams чрез типа Flowable, спазвайки договорите за backpressure за разлика от RxJava 1. Observable в RxJava 2 не поддържа backpressure — предназначен е за потоци с малък брой събития или UI събития.
Според проучване на JetBrains, 2025, RxJava е в топ-3 библиотеки за Android разработка. Основните сценарии за употреба: обработка на мрежови заявки чрез Retrofit (интегриран с RxJava чрез CallAdapter), работа с Room (реактивни заявки връщат Flowable или Maybe), анимации и UI събития чрез RxBinding и debounce търсене при въвеждане на текст. Всички тези сценарии са обединени от верига от един и същи тип: източник (Observable) → трансформация (оператори) → абониране (subscribe).
RxJava 1 (2013) постави концепцията за Observable и операторите, но страдаше от проблеми с backpressure — в бързи потоци данните се натрупваха в паметта, причинявайки OutOfMemoryError. RxJava 2 (2016) поправи архитектурата, разделяйки Observable (без backpressure) и Flowable (с backpressure). RxJava 3 (2020) добави поддръжка за Java 8 Stream API, допълнителни оператори и подобрена производителност при абониране. В момента RxJava 3 е препоръчваната версия за нови проекти.
RxJava предоставя пет основни типа реактивни източници, всеки от които е ориентиран към определен сценарий. Observable и Flowable излъчват множество стойности, Single — една стойност или грешка, Completable — само факта на завършване без данни, Maybe — една стойност, нула или грешка. Изборът на правилния тип намалява количеството код и прави веригата самодокументираща се.
| Тип | Брой събития | Backpressure | Сценарий |
|---|---|---|---|
| Observable | 0..N, след това завършване | Не | UI събития, кратки потоци |
| Flowable | 0..N, след това завършване | Да | Мрежови отговори, потоци от БД |
| Single | Точно 1 или грешка | Не | HTTP заявка, четене на един запис |
| Completable | 0 (само завършване) | Не | Запис в БД, изпращане на събитие |
| Maybe | 0, 1 или грешка | Не | Кеш: има стойност или не |
Flowable е най-гъвкавият тип за работа с големи потоци от данни. Той имплементира Reactive Streams Publisher с поддръжка на backpressure: consumer може да поиска определен брой елементи чрез Subscription.request(n). Това предотвратява препълване на буфера при несъответствие на скоростите на producer и consumer. Ако backpressure не е критичен — използвайте Observable, той има по-малко overhead поради липсата на механизъм request.
Single е оптималният избор за HTTP заявки. Retrofit 2 с RxJava CallAdapter връща Single<ResponseBody> за всяка заявка. Single гарантира точно едно извикване на onSuccess или onError, което съответства на семантиката на HTTP заявка — един отговор или една грешка. Completable се използва за операции за запис, които не връщат данни: insert, update, delete. Maybe е удобен при проверка на кеш — може да върне стойност, може и да не върне.
// Пример за използване на Single за HTTP заявка
interface ApiService {
@GET("users/{id}")
fun getUser(@Path("id") userId: Int): Single<User>
}
// Абониране с обработка на главната нишка
apiService.getUser(42)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe({ user ->
textView.text = user.name
}, { error ->
Log.e("API", "Error: ${error.message}")
})
.addTo(compositeDisposable)
Операторите RxJava са функции от по-висок ред, които приемат един реактивен източник и връщат друг, трансформирайки потока от данни. RxJava 3 съдържа над 400 оператора, разделени в категории: трансформация, филтрация, комбиниране, обработка на грешки и управление на времето. Всеки оператор е мързелив — веригата се изгражда при декларация, изпълнява се при абониране.
map е основният оператор, който трансформира всяка стойност чрез функция. flatMap приема функция, която връща Observable за всеки елемент, и разгръща резултата в единен поток. switchMap е подобен на flatMap, но при получаване на нов елемент се отписва от предишния Observable. concatMap запазва реда на елементите — за разлика от flatMap, той последователно се абонира за всеки вложен Observable.
// Парсване на JSON с трансформация и филтрация
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("Грешка", it.message) })
Комбиниране на потоци е областта, в която RxJava е особено силен. zip комбинира елементи от множество Observable по двойки според индекса: първи с първи, втори с втори. combineLatest излъчва нова стойност при промяна на който и да е от потоците, комбинирайки последните стойности на всички потоци. merge комбинира множество Observable в един, запазвайки реда на пристигане на събитията. concat последователно се абонира за всеки Observable и предава всички негови събития, преди да премине към следващия.
Управление на времето включва debounce (изчакване на пауза в потока преди изпращане), throttleFirst (пропускане на първото събитие, игнориране на останалите в прозореца), timeout (грешка, ако събитието не е пристигнало в интервала). Debounce търсене при въвеждане на текст е най-честият сценарий: searchObservable.debounce(300, MILLISECONDS).distinctUntilChanged() предотвратява ненужни заявки при бързо писане.
| Категория | Оператор | Поведение |
|---|---|---|
| Трансформация | map / flatMap / switchMap | Трансформиране на единична стойност или поток |
| Филтрация | filter / distinct / take | Избор на стойности по условие |
| Комбиниране | zip / combineLatest / merge | Обединяване на 2+ потока |
| Грешки | onErrorResumeNext / retry | Възстановяване след повреда |
| Инструменти | delay / timeout / debounce | Управление на времето в поток |
Scheduler в RxJava е абстракция над група от нишки. Библиотеката предоставя пет вградени Scheduler-а: Schedulers.io() за I/O операции (мрежа, файлове), Schedulers.computation() за CPU-интензивни задачи, Schedulers.newThread() за всяка нова нишка, Schedulers.single() за еднонишково изпълнение и Schedulers.trampoline() за незабавно изпълнение в текущата нишка.
subscribeOn определя на кой Scheduler се изпълнява източникът Observable. Ако във веригата има няколко subscribeOn — приоритет има най-близкият до източника. observeOn превключва downstream на указания Scheduler — всяко използване на observeOn променя нишката за следващите оператори. Типичен Android модел: subscribeOn(Schedulers.io()) за работа с мрежа, observeOn(AndroidSchedulers.mainThread()) за актуализиране на UI.
// Многонишкова обработка с превключване на контекст
Observable.fromCallable(() -> database.getItems())
.subscribeOn(Schedulers.io()) // БД на io
.map(items -> processItems(items)) // трансформация на io
.observeOn(Schedulers.computation()) // превключваме на computation
.map(processed -> compressImages(processed))
.observeOn(AndroidSchedulers.mainThread())
.subscribe(result -> ui.showResult(result))
AndroidSchedulers.mainThread() е Scheduler от библиотеката RxAndroid, който изпълнява код на главната нишка на Android. Той е задължителен за всякакви UI актуализации в реактивната верига. Библиотеката използва Handler вътрешно и гарантира изпълнение в UI нишката дори при високо натоварване. За фонови операции Schedulers.io() поддържа неограничена група от нишки и е подходящ за всякакви блокиращи операции. Schedulers.computation() използва фиксирана група, равна на броя на ядрата на процесора.
RxJava в Android се използва за три основни сценария: реактивни заявки към Room, интеграция с Retrofit и реактивно свързване на UI чрез RxBinding. Всеки сценарий има свой собствен набор от типове: Room връща Flowable за наблюдаеми заявки, Retrofit — Single за HTTP заявки, RxBinding — Observable за UI събития.
Room е библиотека за устойчивост на данни от Google. От Room 2.1 нататък базата данни поддържа реактивни типове за връщане: Flowable и Observable. При промяна на който и да е запис в таблицата, Room автоматично изпраща нова стойност в потока. Разработчикът се абонира за Flowable в ViewModel и получава актуални данни без ръчни заявки при всяка промяна.
// Room DAO с реактивна заявка
@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
class UserViewModel(private val dao: UserDao) : ViewModel() {
val users: Flowable<List<User>> = dao.getAllUsers()
.subscribeOn(Schedulers.io())
}
Моделът MVVM + RxJava се основава на факта, че ViewModel няма референции към View. ViewModel публикува реактивни източници (Flowable, LiveData чрез Transformations), а Activity или Fragment се абонират за тях. Това осигурява тестваемост: ViewModel се тества без UI, като се замества Scheduler чрез RxJavaPlugins.setComputationScheduler. CompositeDisposable в ViewModel управлява жизнения цикъл на абонаментите — при onCleared() всички абонаменти се отменят.
Kotlin Flow е естествена имплементация на студени потоци в Kotlin, вградена в корутините и представена в Kotlin 1.3. Flow решава същите задачи като RxJava, но с фундаментални разлики: вградена поддръжка за корутини (suspend функции), отмяна чрез coroutine cancellation и липса на проблеми с backpressure — Flow използва suspend вместо буфериране. Flow е част от стандартната библиотека на Kotlin и не изисква допълнителни зависимости.
RxJava остава предпочитаният избор за проекти на Java, проекти с поддръжка на Java 7-8 и съществуващи кодови бази на RxJava. Екосистемата на RxJava е значително по-богата: >400 оператора срещу ~50 във Flow, интеграция с Retrofit чрез вграден CallAdapter, поддръжка на backpressure чрез Flowable и наличие на RxBinding, RxPermissions, RxLocation за Android. Kotlin Flow бързо настига, но гъвкавостта на RxJava в сложни сценарии за комбиниране на потоци е все още по-висока.
| Характеристика | RxJava | Kotlin Flow |
|---|---|---|
| Език | Java / Kotlin | Само Kotlin |
| Отмяна | Disposable / CompositeDisposable | Coroutine cancellation |
| Backpressure | Flowable (стратегии BUFFER, DROP, LATEST) | Чрез conflate / buffer |
| Оператори | 400+ | ~50 (разширяем) |
| Room интеграция | Flowable, Observable | Flow, StateFlow |
| ViewModel | CompositeDisposable | viewModelScope + Flow |
Често задавани въпроси
Observable не поддържа backpressure — ако producer е по-бърз от consumer, събитията се натрупват в паметта. Flowable имплементира Reactive Streams с backpressure чрез Subscription.request(), което предотвратява препълване на буфера при несъответствие на скоростите.
Single се използва за операции, които връщат точно една стойност или грешка: HTTP заявки, четене на един запис от БД, изчисляване на резултат. Single семантично съответства на Future и съкращава кода, премахвайки неизползвания onComplete.
Методът dispose() на Disposable отменя абонамента. За групово управление се използва CompositeDisposable — събира всички Disposable и ги отменя едновременно при извикване на clear(). Типично място — onCleared() в ViewModel или onPause() в Activity.
flatMap се абонира за всички вложени Observable и комбинира техните събития в произволен ред. switchMap при получаване на нов елемент се отписва от предишния Observable и се абонира за новия. switchMap се използва при търсене — всяка нова заявка отменя предишната.
За нови проекти на Kotlin Flow е за предпочитане благодарение на интеграцията с корутините и по-малкия размер. За съществуващи проекти на RxJava миграцията е оправдана само ако цялата кодова база преминава към корутини — междинното използване на двете библиотеки усложнява архитектурата.
Резюме
Ще разработим мобилно приложение под ключ
IT Sectr създава iOS и Android приложения за стартъпи и бизнеси от 2017 г. Ще ви консултираме и ще предложим най-доброто решение.
Прочетете също