RxJava 是 Java 和 Android 的响应式编程库,通过 Observable 和 Observer 实现观察者模式。根据 ReactiveX GitHub, 2026,RxJava 允许使用运算符链处理异步数据流和事件。基本单元是 Observable,它通过转换链向 Observer 发射数据。RxJava 3 是当前稳定版本,支持 Java 8 lambda、Reactive Streams 以及与 Android 的集成(通过 RxAndroid)。
要点
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)和其他语言。
ReactiveX 通过两种机制扩展了经典的观察者模式:运算符链(operator chaining)和线程管理(schedulers)。Observable 在 Observer 订阅之前不会开始发射数据(惰性求值)。这使得可以构建仅在订阅存在时才会激活的数据处理管道。
Observable — 基本类型,发射 0..N 个元素,带 onError 或 onComplete。适用于无限长度的数据流 — 例如点击事件或地理定位更新。Observable 不支持背压。
Flowable — 支持背压的 Reactive Streams 版本。当数据源生成元素的速度快于 Observer 的处理能力时使用。Flowable 支持 BACKPRESSURE_BUFFER、DROP、LATEST 和 ERROR 策略。
| 类型 | 元素 | 背压 | 用途 |
|---|---|---|---|
| Observable | 0..N | 否 | UI 事件,小流 |
| Flowable | 0..N | 是 | 大数据,实时 |
| Single | 1 (onSuccess/onError) | — | 单个响应(网络) |
| Maybe | 0..1 | — | 可选值(缓存) |
| Completable | 0 (onComplete/onError) | 无数据操作(写入) |
Single 精确发射一个元素或错误 — 非常适合网络请求。Maybe — 0 或 1 个元素,适用于数据可能不存在的缓存。Completable — 只有 onComplete 或 onError,没有数据,方便进行写入或删除操作。这些类型简化了 API,将契约缩小到特定情况。Retrofit(流行的 Android HTTP 客户端)直接支持所有五种 RxJava 类型,允许为每个端点选择最合适的返回类型而无需不必要的包装。
运算符是将一个 Observable 转换为另一个的函数。运算符链(operator chain)描述了数据处理管道:每个运算符从前一个接收流,进行转换并传递给下一个。RxJava 包含 200 多个运算符,分为不同类别。
flatMap — RxJava 最强大的运算符之一。它允许为每个元素执行异步请求并将结果收集到一个公共流中。例如,flatMap 用于根据 ID 列表加载详细信息:每个 ID → 网络请求 → 合并结果。与简单地转换元素的 map 不同,flatMap 可以发射多个元素或切换到另一个 Observable,这使其成为构建异步管道的基础。
onErrorResumeNext — 出错时切换到备用 Observable。retry — 出错时重试订阅 N 次。onErrorReturn — 返回默认值而不是错误。doOnError — 在出错时执行副作用而不改变流(日志记录或分析)。组合这些运算符可以在无需手动 try/catch 的情况下构建具有清晰故障处理策略的可靠管道。
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 切换下游(订阅者及其之后的所有运算符)。
让我们考虑三种场景:使用 Single 进行网络请求、使用 zip 进行并行请求以及使用 debounce 对搜索字段进行防抖处理。
Single 非常适合 Retrofit 请求:一个请求 — 一个响应。在主线程上订阅以更新 UI。
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 将两个独立 Single 的结果合并为一个。并行执行,结果 — 两者都完成后。
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 如果文本未更改则取消请求。
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 解决相同的任务 — 异步编程 — 但采用根本不同的方法。RxJava 建立在观察者模式之上并且是基于推送的:源发送数据,Observer 做出反应。协程是拉取式的:代码通过 await 顺序获取数据。
根据 Google I/O 2024,Kotlin Coroutines 是 Android 中新的异步代码的推荐方法。RxJava 仍然支持现有项目。Google 提供过渡库(kotlinx-coroutines-rx3)用于逐步迁移。AndroidX(LiveData、Room、Paging 3)支持两种方法,允许在旧模块中使用 RxJava 并在新模块中使用协程而不会产生依赖冲突。
逐步过渡:每个新组件都用协程编写,旧的 RxJava 代码保持不变。RxJava → 协程 通过 awaitSingle() 或 awaitFirst()。协程 → RxJava 通过 future() 或 asFlowable()。对于大型项目,完全迁移需要 6–18 个月。
常见问题
Observable 不支持背压 — 如果源生成数据的速度快于处理器,将发生 MissingBackpressureException。Flowable 支持具有可配置缓冲策略的 Reactive Streams 背压。
subscribeOn 设置用于执行 Observable 源的 Scheduler。observeOn 设置用于 Observer 和链中所有后续运算符的 Scheduler。subscribeOn 影响上游,observeOn 影响下游。
对于新项目 — 是,Google 推荐协程。对于现有项目 — 通过 kotlinx-coroutines-rx3 逐步迁移。RxJava 对旧代码保持稳定和支持。
通过运算符:onErrorReturn(默认值)、onErrorResumeNext(备用 Observable)、retry(重试 N 次)。或通过 Observer.onError() 向用户显示。
CompositeDisposable — 用于管理多个订阅的容器。调用 dispose() 时将取消所有添加的订阅。在 Activity/Fragment 中用于在屏幕销毁时取消所有请求。
总结
我们将开发一款交钥匙移动应用程序
IT Sectr自2017年以来为初创企业和企业打造iOS和Android应用程序。我们将为您提供咨询并提出最佳解决方案。