RxJava:基础、ReactiveX 以及数据流处理

作者: IT Sectr 发布日期: 2026-03-16 阅读时间: 8 分钟

RxJava 是 Java 和 Android 的响应式编程库,通过 Observable 和 Observer 实现观察者模式。根据 ReactiveX GitHub, 2026RxJava 允许使用运算符链处理异步数据流和事件。基本单元是 Observable,它通过转换链向 Observer 发射数据。RxJava 3 是当前稳定版本,支持 Java 8 lambda、Reactive Streams 以及与 Android 的集成(通过 RxAndroid)。

要点

  • RxJava — ReactiveX 的 Java 实现,用于异步数据处理
  • Observable — 数据源,向 Observer 发射元素
  • Observer — 订阅者,接收 onNext、onError 和 onComplete 通知
  • 运算符 — 用于转换、过滤和组合流的函数链
  • Schedulers — 用于管理 Observable 和 Observer 执行线程的组件

什么是 RxJava 和 ReactiveX

RxJava — ReactiveX 规范的 Java 实现,一个使用可观察流(Observable)进行异步编程的库。RxJava 2 于 2016 年发布,支持 Reactive Streams(Flowable)并分为 rx.Observable 和 io.reactivex.Observable。RxJava 3(2019)— 当前主要版本,向后兼容 RxJava 2。

RxJava 的核心思想 — 一切都是流:数据流事件流状态流。任何异步操作都可以表示为发射数据、错误或完成信号的 Observable。Observer 订阅 Observable 并实时接收通知。

根据 Badoo 的数据(2024),在转向协程之前,Google Play 前 200 名中的 76% Android 应用使用 RxJava 进行异步操作。现在份额正在向协程倾斜,但 RxJava 仍存在于数千个应用的生产代码中,被认为是成熟、经过验证的技术。ReactiveX — 一个跨平台规范,也适用于 JavaScript(RxJS)、.NET(Rx.NET)、Swift(RxSwift)和其他语言。

RxJava 中的观察者模式

ReactiveX 通过两种机制扩展了经典的观察者模式:运算符链(operator chaining)和线程管理(schedulers)。Observable 在 Observer 订阅之前不会开始发射数据(惰性求值)。这使得可以构建仅在订阅存在时才会激活的数据处理管道。

Observable 的类型:Observable、Flowable、Single、Maybe、Completable

Observable — 基本类型,发射 0..N 个元素,带 onError 或 onComplete。适用于无限长度的数据流 — 例如点击事件或地理定位更新。Observable 不支持背压。

Flowable — 支持背压的 Reactive Streams 版本。当数据源生成元素的速度快于 Observer 的处理能力时使用。Flowable 支持 BACKPRESSURE_BUFFER、DROP、LATEST 和 ERROR 策略。

类型元素背压用途
Observable0..NUI 事件,小流
Flowable0..N大数据,实时
Single1 (onSuccess/onError)单个响应(网络)
Maybe0..1可选值(缓存)
Completable0 (onComplete/onError)无数据操作(写入)

Single、Maybe 和 Completable

Single 精确发射一个元素或错误 — 非常适合网络请求。Maybe — 0 或 1 个元素,适用于数据可能不存在的缓存。Completable — 只有 onComplete 或 onError,没有数据,方便进行写入或删除操作。这些类型简化了 API,将契约缩小到特定情况。Retrofit(流行的 Android HTTP 客户端)直接支持所有五种 RxJava 类型,允许为每个端点选择最合适的返回类型而无需不必要的包装。

RxJava 运算符:流的转换和过滤

运算符是将一个 Observable 转换为另一个的函数。运算符链(operator chain)描述了数据处理管道:每个运算符从前一个接收流,进行转换并传递给下一个。RxJava 包含 200 多个运算符,分为不同类别。

  • map — 转换每个元素(Integer → String)
  • flatMap — 将元素转换为 Observable 并将所有内容合并到一个流中
  • filter — 根据条件过滤元素
  • zip — 按索引组合 N 个 Observable 的元素
  • merge — 将多个 Observable 合并为一个,保持时间顺序
  • debounce — 如果元素之间的间隔小于指定间隔则通过

flatMap — RxJava 最强大的运算符之一。它允许为每个元素执行异步请求并将结果收集到一个公共流中。例如,flatMap 用于根据 ID 列表加载详细信息:每个 ID → 网络请求 → 合并结果。与简单地转换元素的 map 不同,flatMap 可以发射多个元素或切换到另一个 Observable,这使其成为构建异步管道的基础。

通过运算符进行错误处理

onErrorResumeNext — 出错时切换到备用 Observable。retry — 出错时重试订阅 N 次。onErrorReturn — 返回默认值而不是错误。doOnError — 在出错时执行副作用而不改变流(日志记录或分析)。组合这些运算符可以在无需手动 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 初学者最困难的部分。典型错误 — 在 observeOn 之后调用 subscribeOn,这不会影响源。subscribeOn 应该在链中作为源的第一个,observeOn — 在 UI 订阅之前。规则:subscribeOn 只影响上游(源),observeOn 切换下游(订阅者及其之后的所有运算符)。

Android 中使用 RxJava 的代码示例

让我们考虑三种场景:使用 Single 进行网络请求、使用 zip 进行并行请求以及使用 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 忽略快速文本更改,仅在 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 与 Kotlin Coroutines:方法比较

RxJavaKotlin Coroutines 解决相同的任务 — 异步编程 — 但采用根本不同的方法。RxJava 建立在观察者模式之上并且是基于推送的:源发送数据,Observer 做出反应。协程是拉取式的:代码通过 await 顺序获取数据。

  • RxJava — 响应式,数据流,>200 个运算符,基于推送,学习曲线陡峭
  • Coroutines — 顺序式,suspend/await,约 40 个函数,基于拉取,语法简单
  • RxJava — 成熟(2016),生态系统庞大,但学习曲线陡峭
  • Coroutines — 现代(2018),Google 对新代码的首选
  • RxJava — 通过 Flowable 开箱即用的背压支持,完善的缓冲策略
  • Coroutines — Flow 近期才支持背压,但 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 不支持背压 — 如果源生成数据的速度快于处理器,将发生 MissingBackpressureException。Flowable 支持具有可配置缓冲策略的 Reactive Streams 背压。

什么是 subscribeOn 和 observeOn?

subscribeOn 设置用于执行 Observable 源的 Scheduler。observeOn 设置用于 Observer 和链中所有后续运算符的 Scheduler。subscribeOn 影响上游,observeOn 影响下游。

是否值得从 RxJava 迁移到协程?

对于新项目 — 是,Google 推荐协程。对于现有项目 — 通过 kotlinx-coroutines-rx3 逐步迁移。RxJava 对旧代码保持稳定和支持。

如何在 RxJava 中处理错误?

通过运算符:onErrorReturn(默认值)、onErrorResumeNext(备用 Observable)、retry(重试 N 次)。或通过 Observer.onError() 向用户显示。

什么是 CompositeDisposable?

CompositeDisposable — 用于管理多个订阅的容器。调用 dispose() 时将取消所有添加的订阅。在 Activity/Fragment 中用于在屏幕销毁时取消所有请求。

总结

  • RxJava — 基于观察者模式的 Java 和 Android 响应式编程库
  • Observable/Flowable — 分别带有和不带背压支持的数据源
  • Single、Maybe、Completable — 针对 1、0..1 和 0 个元素的专用类型
  • 运算符(map、flatMap、zip、filter)— 具有 200 多个函数的转换链
  • Schedulers — 源的 subscribeOn 和数据消费者的 observeOn
  • RxJava 与 Coroutines — Google 推荐将协程用于新代码,RxJava 用于遗留代码
  • CompositeDisposable — 安全的订阅管理,在屏幕销毁时取消

我们将开发一款交钥匙移动应用程序

IT Sectr自2017年以来为初创企业和企业打造iOS和Android应用程序。我们将为您提供咨询并提出最佳解决方案。

讨论项目

另请阅读