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 остаётся в production-коде тысяч приложений и считается зрелой, проверенной технологией. ReactiveX — это кроссплатформенная спецификация, реализованная также для JavaScript (RxJS), .NET (Rx.NET), Swift (RxSwift) и других языков.

Паттерн Observer в RxJava

ReactiveX расширяет классический паттерн Observer двумя механизмами: цепочка операторов (operator chaining) и управление потоками (schedulers). Observable не начинает эмиттить данные, пока на него не подпишется Observer (lazy evaluation). Это позволяет строить 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 мс паузы. 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 vs Kotlin Coroutines: сравнение подходов

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

  • RxJava — реактивный, поток данных, >200 операторов, push-based, сложный learning curve
  • Coroutines — последовательный, suspend/await, ~40 функций, pull-based, простой синтаксис
  • RxJava — зрелый (2016), огромная экосистема, но тяжёлый learning curve
  • 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 vs Coroutines — корутины рекомендованы Google для нового кода, RxJava для legacy
  • CompositeDisposable — безопасное управление подписками с отменой при уничтожении экрана

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

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

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

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