RxJava 是一个用于JVM的响应式编程库,通过Observable模式配合函数式转换操作符实现异步数据流。它将ReactiveX的概念移植到Java和Kotlin,为网络请求、数据库、UI事件和后台任务提供统一的API。根据ReactiveX, 2025的数据,该库在GitHub上被超过12万个项目使用,并且在Kotlin Flow出现之前一直是Android响应式编程的标准。RxJava用统一的数据处理链取代了AsyncTask、Loader和回调。
要点
RxJava 是ReactiveX(响应式扩展)库在Java虚拟机上的实现。RxJava的第一个版本由Netflix公司于2013年发布,用于管理服务器应用程序中的异步调用。在创建时,Java中的主要替代方案是Future和Callback — 这两种方法都导致了回调地狱和复杂的线程管理。RxJava通过带有函数式操作符链的Observable提出了异步操作的组合。
RxJava 的架构基于Reactive Streams规范 — 一种具有非阻塞背压的异步流处理标准。该规范定义了四个接口:Publisher、Subscriber、Subscription和Processor。RxJava 2+通过Flowable类型完全实现了Reactive Streams,与RxJava 1不同,它遵守背压契约。RxJava 2中的Observable不支持背压 — 它适用于事件数量较少的流或UI事件。
根据JetBrains, 2025的调查,RxJava是Android开发的前三大库之一。主要使用场景:通过Retrofit处理网络请求(通过CallAdapter与RxJava集成)、使用Room(响应式查询返回Flowable或Maybe)、通过RxBinding实现动画和UI事件,以及在文本输入时进行防抖搜索。所有这些场景都由相同类型的链连接:源(Observable)→转换(操作符)→订阅(subscribe)。
RxJava 1(2013年)奠定了Observable和操作符的概念,但存在背压问题 — 在快速流中,数据会积累在内存中,导致OutOfMemoryError。RxJava 2(2016年)修复了架构,将Observable(无背压)和Flowable(有背压)分开。RxJava 3(2020年)增加了对Java 8 Stream API的支持、额外的操作符以及改进的订阅性能。目前,RxJava 3是新项目的推荐版本。
RxJava 提供了五种主要的响应式源类型,每种针对特定场景。Observable和Flowable发出多个值,Single — 一个值或错误,Completable — 仅完成的事实而不带数据,Maybe — 一个值、零个或错误。选择正确的类型可以减少代码量并使链自文档化。
| 类型 | 事件数量 | 背压 | 场景 |
|---|---|---|---|
| Observable | 0..N,然后完成 | 否 | UI事件,短流 |
| Flowable | 0..N,然后完成 | 是 | 网络响应,数据库流 |
| Single | 恰好1个或错误 | 否 | HTTP请求,读取一条记录 |
| Completable | 0(仅完成) | 否 | 写入数据库,发送事件 |
| Maybe | 0、1或错误 | 否 | 缓存:有值或无值 |
Flowable 是处理大数据流最灵活的类型。它实现了带有背压支持的Reactive Streams Publisher:消费者可以通过Subscription.request(n)请求特定数量的元素。这防止了生产者和消费者速度不匹配时缓冲区溢出。如果背压不重要 — 请使用Observable,由于没有请求机制,它的开销更小。
Single 是HTTP请求的最佳选择。带有RxJava CallAdapter的Retrofit 2为每个请求返回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("错误", 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 确定Observable源在哪个Scheduler上执行。如果链中有多个subscribeOn — 离源最近的具有优先级。observeOn 将下游切换到指定的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() 是来自RxAndroid库的Scheduler,它在Android主线程上执行代码。它对响应式链中的任何UI更新都是必需的。该库在内部使用Handler,并保证即使在高负载下也能在UI线程中执行。对于后台操作,Schedulers.io()支持无限制的线程池,适用于任何阻塞操作。Schedulers.computation()使用固定池,等于处理器核心数。
RxJava 在Android中用于三个主要场景:对Room的响应式查询、与Retrofit的集成以及通过RxBinding的响应式UI绑定。每个场景都有自己的一套类型:Room为可观察查询返回Flowable,Retrofit — 为HTTP请求返回Single,RxBinding — 为UI事件返回Observable。
Room 是Google的数据持久化库。从Room 2.1开始,数据库支持响应式返回类型:Flowable和Observable。当表中的任何记录发生变化时,Room会自动向流中发送新值。开发者在ViewModel中订阅Flowable,并在每次更改时无需手动查询即可获取最新数据。
// 带响应式查询的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、通过Transformations的LiveData),Activity或Fragment订阅它们。这提供了可测试性:ViewModel在没有UI的情况下进行测试,通过RxJavaPlugins.setComputationScheduler替换Scheduler。ViewModel中的CompositeDisposable管理订阅的生命周期 — 在onCleared()时,所有订阅都被取消。
Kotlin Flow 是Kotlin中冷流的原生实现,内置于协程中,并在Kotlin 1.3中引入。Flow解决了与RxJava相同的任务,但具有根本性的区别:内置的协程支持(挂起函数)、通过协程取消进行取消,以及没有背压问题 — Flow使用挂起而不是缓冲。Flow是Kotlin标准库的一部分,不需要额外的依赖。
RxJava 仍然是Java项目、支持Java 7-8的项目以及现有RxJava代码库的首选。RxJava生态系统明显更丰富:>400个操作符对比Flow的约50个,通过内置CallAdapter与Retrofit集成,通过Flowable支持背压,以及RxBinding、RxPermissions、RxLocation在Android上的可用性。Kotlin Flow正在快速追赶,但RxJava在复杂流组合场景中的灵活性仍然更高。
| 特性 | RxJava | Kotlin Flow |
|---|---|---|
| 语言 | Java / Kotlin | 仅Kotlin |
| 取消 | Disposable / CompositeDisposable | 协程取消 |
| 背压 | Flowable(BUFFER、DROP、LATEST策略) | 通过conflate / buffer |
| 操作符 | 400+ | 约50(可扩展) |
| Room集成 | Flowable, Observable | Flow, StateFlow |
| ViewModel | CompositeDisposable | viewModelScope + Flow |
常见问题
Observable 不支持背压 — 如果生产者快于消费者,事件会积累在内存中。Flowable 通过Subscription.request()实现带有背压的Reactive Streams,防止速度不匹配时缓冲区溢出。
Single 用于返回恰好一个值或错误的操作:HTTP请求、从数据库读取一条记录、计算结果。Single在语义上等同于Future,通过删除未使用的onComplete来缩短代码。
Disposable上的dispose()方法取消订阅。对于组管理,使用CompositeDisposable — 它收集所有Disposable并在调用clear()时同时取消它们。典型位置 — ViewModel中的onCleared()或Activity中的onPause()。
flatMap 订阅所有嵌套的Observable,并以任意顺序组合它们的事件。switchMap 在收到新元素时取消订阅前一个Observable并订阅新的。switchMap用于搜索 — 每个新请求取消前一个。
对于Kotlin新项目,Flow 因与协程集成和更小的体积而更受青睐。对于现有的RxJava项目,迁移只有当整个代码库切换到协程时才合理 — 同时使用两个库会使架构复杂化。
总结
我们将开发一款交钥匙移动应用程序
IT Sectr自2017年以来为初创企业和企业打造iOS和Android应用程序。我们将为您提供咨询并提出最佳解决方案。