RxJava: суть, компоненты и реактивное программирование

Автор: IT Sectr Опубликовано: 2026-05-03 Время чтения: 10 мин

RxJava — это библиотека реактивного программирования для JVM, реализующая асинхронные потоки данных через паттерн Observable с функциональными операторами трансформации. Она портирует концепции ReactiveX на Java и Kotlin, предоставляя единый API для работы с сетевыми запросами, базами данных, UI-событиями и фоновыми задачами. По данным ReactiveX, 2025, библиотека используется более чем в 120 000 проектах на GitHub и является стандартом реактивного программирования для Android до появления Kotlin Flow. RxJava заменяет AsyncTask, Loader и callback-и единой цепочкой обработки данных.

Главное

  • RxJava — ReactiveX-реализация для Java/Kotlin с типами Observable, Flowable, Single, Completable и Maybe
  • Observable представляет поток данных с backpressure-управлением через Flowable при подписке на медленном consumer
  • Операторы map, flatMap, switchMap, zip и combineLatest трансформируют и комбинируют асинхронные потоки без блокировок
  • Scheduler — Schedulers.io(), computation(), mainThread() управляют на каком потоке выполняется работа и подписка
  • RxAndroid добавляет AndroidSchedulers.mainThread() для обновления UI из реактивных цепочек

Что такое RxJava?

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

RxJava 1 (2013) заложил концепцию Observable и операторов, но страдал от проблем с backpressure — в быстрых потоках данные накапливались в памяти, вызывая OutOfMemoryError. RxJava 2 (2016) исправил архитектуру, разделив Observable (без backpressure) и Flowable (с backpressure). RxJava 3 (2020) добавил поддержку Java 8 Stream API, дополнительные операторы и улучшенную производительность при subscription. На текущий момент RxJava 3 — рекомендованная версия для новых проектов.

Типы реактивных потоков в RxJava

RxJava предоставляет пять основных типов реактивных источников, каждый из которых ориентирован на определённый сценарий. Observable и Flowable испускают множество значений, Single — одно значение или ошибку, Completable — только факт завершения без данных, Maybe — одно значение, ноль или ошибка. Выбор правильного типа сокращает количество кода и делает цепочку самодокументированной.

ТипКоличество событийBackpressureСценарий
Observable0..N, затем завершениеНетUI-события, короткие потоки
Flowable0..N, затем завершениеДаСетевые ответы, потоки из БД
SingleРовно 1 или ошибкаНетHTTP-запрос, чтение одной записи
Completable0 (только завершение)НетЗапись в БД, отправка события
Maybe0, 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 удобен при проверке кэша — может вернуть значение, может не вернуть.

kotlin
// Пример использования 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.

kotlin
// Парсинг 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Управление временем в потоке

Schedulers и многопоточность

Scheduler в RxJava — это абстракция над пулом потоков. Библиотека предоставляет пять встроенных Scheduler: Schedulers.io() для I/O-операций (сеть, файлы), Schedulers.computation() для CPU-интенсивных задач, Schedulers.newThread() для каждого нового потока, Schedulers.single() для однопоточного исполнения и Schedulers.trampoline() для немедленного выполнения в текущем потоке.

subscribeOn и observeOn

subscribeOn определяет, на каком Scheduler выполняется источник Observable. Если в цепочке несколько subscribeOn — приоритет у ближайшего к источнику. observeOn переключает downstream на указанный Scheduler — каждое использование observeOn меняет поток для последующих операторов. Типичный Android-паттерн: subscribeOn(Schedulers.io()) для работы с сетью, observeOn(AndroidSchedulers.mainThread()) для обновления UI.

java
// Многопоточная обработка с переключением контекста
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: практическое применение

RxJava в Android используется для трёх основных сценариев: реактивные запросы к Room, интеграция с Retrofit и реактивное связывание UI через RxBinding. Для каждого сценария характерен свой набор типов: Room возвращает Flowable для наблюдаемых запросов, Retrofit — Single для HTTP-запросов, RxBinding — Observable для UI-событий.

Room + RxJava

Room — библиотека персистентности данных от Google. Начиная с Room 2.1, база данных поддерживает реактивные возвращаемые типы: Flowable и Observable. При изменении любой записи в таблице Room автоматически отправляет новое значение в поток. Разработчик подписывается на Flowable в ViewModel и получает актуальные данные без ручных запросов при каждом изменении.

kotlin
// 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() все подписки отменяются.

RxJava vs Kotlin Flow

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 в сложных сценариях комбинирования потоков пока выше.

ХарактеристикаRxJavaKotlin Flow
ЯзыкJava / KotlinKotlin только
ОтменаDisposable / CompositeDisposableCoroutine cancellation
BackpressureFlowable (стратегии BUFFER, DROP, LATEST)Через conflate / buffer
Операторы400+~50 (расширяется)
Room интеграцияFlowable, ObservableFlow, StateFlow
ViewModelCompositeDisposableviewModelScope + Flow

Часто задаваемые вопросы

В чём разница между Observable и Flowable в RxJava?

Observable не поддерживает backpressure — если producer быстрее consumer, события накапливаются в памяти. Flowable реализует Reactive Streams с backpressure через Subscription.request(), что предотвращает переполнение буфера при несоответствии скоростей.

Когда использовать Single вместо Observable?

Single используется для операций, которые возвращают ровно одно значение или ошибку: HTTP-запросы, чтение одной записи из БД, вычисление результата. Single семантически соответствует Future и сокращает код, убирая неиспользуемые onComplete.

Как отменить подписку в RxJava?

Метод dispose() на Disposable отменяет подписку. Для группового управления используется CompositeDisposable — он собирает все Disposable и отменяет их одновременно при вызове clear(). Типичное место — onCleared() в ViewModel или onPause() в Activity.

Чем отличается flatMap от switchMap?

flatMap подписывается на все вложенные Observable и объединяет их события в произвольном порядке. switchMap при поступлении нового элемента отписывается от предыдущего Observable и подписывается на новый. switchMap используется при поиске — каждый новый запрос отменяет предыдущий.

Стоит ли мигрировать с RxJava на Kotlin Flow?

Для новых проектов на Kotlin Flow предпочтительнее благодаря интеграции с корутинами и меньшему размеру. Для существующих проектов на RxJava миграция оправдана только если вся кодовая база переходит на корутины — промежуточное использование обеих библиотек усложняет архитектуру.

Итоги

  • RxJava — ReactiveX-библиотека для JVM с типами Observable, Flowable, Single, Completable и Maybe для разных сценариев
  • Flowable поддерживает backpressure через Reactive Streams для предотвращения переполнения при несовпадении скоростей
  • Операторы map, flatMap, switchMap, zip, combineLatest, debounce обеспечивают декларативную обработку потоков
  • Schedulers io(), computation(), mainThread() управляют потоками выполнения без блокировки UI
  • RxAndroid интегрирует RxJava с Android, предоставляя AndroidSchedulers.mainThread() и упрощая обновление UI
  • Kotlin Flow — нативная альтернатива с интеграцией в корутины, но RxJava сохраняет преимущество в экосистеме операторов
  • MVVM + RxJava — стандартный паттерн Android-разработки с отделённой от UI ViewModel и реактивными подписками

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

IT Sectr создаёт приложения для iOS и Android для стартапов и бизнеса с 2017 года. Мы проконсультируем вас и предложим наилучшее решение.

Обсудить проект

Читайте также