这篇教程用最直观的方式示范如何在一个简单的 HelloWorld Android 项目中集成 RxJava:从依赖配置、线程调度、可观察数据流、常用操作符,到资源回收和错误处理,逐步示例代码并解释每一步的原理与注意点,帮助你在实际项目中安全、清晰地采用响应式编程。以及实践中的坑位与优化建议。适合新手上手。


先说结论(快速上手思路)
把 RxJava 当成「把异步流写成同步样子」的工具。先在 Gradle 加依赖,理解 Observable/Flowable、Subscriber/Observer、Scheduler 和 Disposable,然后用简单例子把回调改成链式操作,最后用 CompositeDisposable 管理生命周期、用 Flowable 处理背压。下面一步步来,边写边解释。
为什么要用 RxJava?
先简单说优点:
- 把异步逻辑抽象为数据流,代码更可组合、可测试。
- 链式操作符(map、flatMap 等)让转换与组合变得表达性强。
- 统一线程切换(subscribeOn/observeOn),简化 UI 与后台线程交互。
- 配合 Retrofit、数据库、事件总线等能显著减少回调地狱。
缺点也要知道:学习曲线、调试困难、滥用会导致不可维护的链。权衡后在复杂异步场景引入是值得的。
核心概念一遍讲清楚(费曼式解释)
Observable / Flowable / Single / Maybe / Completable
把它们想像成不同口径的水管:
- Observable:常规流,适合无限或中等量数据。
- Flowable:带背压的流,适合高频生产者,防止消费者被淹没。
- Single:只会发一次成功值或错误(比如网络请求返回一个对象)。
- Maybe:可能有值也可能无值或错误。
- Completable:只关心完成或错误(比如写数据库,只关心成功/失败)。
Observer / Subscriber / Disposable
订阅者就是消费者,订阅后会拿到数据或错误。Disposable 用来取消订阅,避免内存泄漏。把多个 Disposable 放进 CompositeDisposable 一并清理。
Scheduler(线程调度)
两个常用概念:subscribeOn 决定数据生产在哪个线程,observeOn 决定观察者在哪个线程收到数据。常用:Schedulers.io()、Schedulers.computation()、AndroidSchedulers.mainThread()(需要 RxAndroid)。
在 HelloWorld Android 项目中一步步集成
1. 新建项目与添加依赖
这是最简单也是最容易出错的地方:版本不匹配。下面给出常用组合(示例版本,随时间升级):
| 库 | 示例版本(按需替换) |
| io.reactivex.rxjava3:rxjava | 3.1.5 |
| io.reactivex.rxjava3:rxandroid | 3.0.0 |
| com.squareup.retrofit2:adapter-rxjava3 | 2.9.0(若用 Retrofit) |
在 app/build.gradle(Kotlin 或 Groovy)里加入:
implementation 'io.reactivex.rxjava3:rxjava:3.1.5'
implementation 'io.reactivex.rxjava3:rxandroid:3.0.0'
2. 第一个 HelloWorld 示例(最小可运行)
思路:Observable 发出 “Hello RxJava”,订阅者在主线程更新 TextView。
// 假想 Activity 代码片段(概念示意)
private val disposables = CompositeDisposable()
override fun onCreate(...) {
// ...
val obs = io.reactivex.rxjava3.core.Observable.just("Hello RxJava")
.subscribeOn(io.reactivex.rxjava3.schedulers.Schedulers.io())
.observeOn(io.reactivex.rxjava3.android.schedulers.AndroidSchedulers.mainThread())
.map { it + " — from HelloWorld" }
val d = obs.subscribe({ text ->
textView.text = text
}, { err ->
textView.text = "Error: ${err.message}"
})
disposables.add(d)
}
override fun onDestroy() {
super.onDestroy()
disposables.clear()
}
要点:把订阅结果放到 CompositeDisposable,Activity 销毁时 clear,避免泄漏。
常用操作符与直观示例(把用途讲清楚)
- map:一对一转换,像给每个元素做一次变换。
- flatMap:一对多或异步映射,常用于把一个请求映射成另一个网络请求(结果乱序)。
- concatMap:类似 flatMap,但保持顺序,按序发出。
- switchMap:如果新的内流来了,丢弃旧的内流,适合搜索联想场景。
- filter:过滤不需要的数据。
- zip:合并多个流,按位置配对输出。
示例:两个网络请求串联(伪代码):
api.getUser(userId) // Single
.flatMap { user -> api.getPosts(user.id) } // Single>
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe({ posts -> /* 更新 UI */ }, { err -> /* 处理错误 */ })
错误处理与重试策略
- onErrorReturn:出错时返回一个默认值继续流。
- onErrorResumeNext:出错时切换到另一个流。
- retry / retryWhen:重试策略,retryWhen 可结合延迟与条件。
不要滥用无限重试——要有上限和退避(backoff)策略来防止雪崩。
背压(Backpressure)说明与实战
当生产速度远高于消费速度时,需要背压保护。RxJava 2/3 使用 Flowable 来表达带背压的流。创建高频源时优先用 Flowable.create,并选择合适的 BackpressureStrategy(BUFFER、DROP、LATEST、ERROR)。
// 简单示例
Flowable.create({ emitter ->
for (i in 1..10000) {
if (emitter.isCancelled) return@create
emitter.onNext(i)
}
emitter.onComplete()
}, BackpressureStrategy.BUFFER)
.observeOn(Schedulers.computation())
.subscribe({ v -> process(v) }, { e -> log(e) })
与 Android 生命周期和架构组件协作
实务中常见三种做法:
- 在 Activity/Fragment 用 CompositeDisposable,并在 onDestroy/onDestroyView 清理。
- 在 ViewModel 持有 CompositeDisposable,Activity/Fragment 观察 ViewModel,ViewModel 在 onCleared 清理。
- 使用 AutoDispose 或 Lifecycle-aware 适配器(额外库)实现自动解绑。
与 Retrofit 集成(常见场景)
Retrofit 支持返回 Single/Observable/Flowable,配合 RxJava 能把网络调用变成链式操作。要在 Retrofit 中加入 CallAdapter(依赖示例见上表)。示例接口:
interface ApiService {
@GET("user/{id}")
fun getUser(@Path("id") id: String): io.reactivex.rxjava3.core.Single
}
要注意:网络错误会落到 onError,别忘了做友好的错误提示与重试。
调试与性能优化小贴士(实用)
- 用 doOnNext / doOnError 插入日志,调试链条数据。
- 避免在操作符里进行耗时阻塞,耗时任务放到 Schedulers.io() 或 computation。
- 限制并发:flatMap 有并发参数,避免同时发太多请求。
- 尽量使用具体类型(Single/Completable)表达语义,代码可读性更好。
常见坑位(真的会踩的)
- 忘记 clear/ dispose,导致 Activity 泄漏。
- subscribeOn/observeOn 顺序理解错,引发 UI 在后台线程更新的崩溃。
- 把 Observable 用在高频场景却不考虑背压,导致 OOM 或丢帧。
- 在链条里捕获错误却未重新抛出,导致上层永远不知道异常发生。
把前面内容串成一个小实战流程(一步步做)
- 在项目里添加 rxjava + rxandroid 依赖并 sync。
- 先写一个最简单的 Observable.just 测试端到端:发、转、订阅,确保主线程能更新 UI。
- 把耗时任务放到 Schedulers.io(),用 observeOn(AndroidSchedulers.mainThread()) 更新 UI。
- 把 Disposable 添加到 CompositeDisposable,在 onDestroy() 清理。
- 遇到高频数据改用 Flowable 并选择合适 BackpressureStrategy。
- 当需要多个并发或串联请求时,选择合适的操作符(flatMap、concatMap、zip)。
举几个容易记住的小范例(便于带回工作中复用)
1) 把回调改为 Single:
fun fetchOnce(): Single {
return Single.create { emitter ->
networkCall { result, err ->
if (err != null) emitter.onError(err)
else emitter.onSuccess(result)
}
}
}
2) 搜索联想防抖(switchMap + debounce):
textChanges // Observable from EditText
.debounce(300, TimeUnit.MILLISECONDS)
.filter { it.isNotBlank() }
.switchMapSingle { query -> api.search(query).toSingle() }
.observeOn(AndroidSchedulers.mainThread())
.subscribe({ showResults(it) }, { log(it) })
最后说点生活化的建议(这个很管用)
刚开始不要试图把项目全盘改成 Rx,选一个清晰的入口(网络请求或事件流)先做改造,跑通了再慢慢推广。写链的时候刻意把每一步拆成小的、可测试的函数,方便排查问题。遇到奇怪的线程问题先回头确认 subscribeOn/observeOn 的顺序。
如果你跟我一样喜欢边写边想,就会发现 Rx 的很多力量来自将复杂的时间维度统一成流来处理——一开始有点抽象,但养成习惯后你会觉得代码干净很多。试试看上面的小步骤,别忘了版本和生命周期的那两个“坑”。