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, и дебаунс-поиск при вводе текста. Все эти сценарии объединяет однотипная цепочка: источник (Observable) → трансформация (операторы) → подписка (subscribe).
RxJava 1 (2013) заложил концепцию Observable и операторов, но страдал от проблем с backpressure — в быстрых потоках данные накапливались в памяти, вызывая OutOfMemoryError. RxJava 2 (2016) исправил архитектуру, разделив Observable (без backpressure) и Flowable (с backpressure). RxJava 3 (2020) добавил поддержку Java 8 Stream API, дополнительные операторы и улучшенную производительность при subscription. На текущий момент 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, он имеет меньше накладных расходов из-за отсутствия механизма 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("Err", it.message) })
Комбинирование потоков — область, где RxJava особенно силён. zip объединяет элементы из нескольких Observable попарно по индексу: первый с первым, второй со вторым. combineLatest испускает новое значение при изменении любого из потоков, комбинируя последние значения всех потоков. merge объединяет несколько Observable в один, сохраняя порядок поступления событий. concat последовательно подписывается на каждый Observable и передаёт все его события, прежде чем перейти к следующему.
Управление временем включает debounce (ожидание паузы в потоке перед отправкой), throttleFirst (пропуск первого события, игнорирование остальных в течение окна), timeout (ошибка, если событие не поступило в течение интервала). Дебаунс-поиск при вводе текста — самый распространённый сценарий: 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 года. Мы проконсультируем вас и предложим наилучшее решение.
Читайте также