RxJava 3 异步编程
RxJava 3 异步编程:基础类型 Observable / Flowable / Single / Maybe / Completable、常用操作符、线程调度 subscribeOn / observeOn、背压策略、与 Android Lifecycle 互操作、内存泄漏与 dispose / CompositeDisposable。
学习目标
- 掌握 RxJava 3 五种基础类型的语义与适用场景
- 熟练使用 map / flatMap / filter / concatMap 等常用操作符
- 理解 subscribeOn / observeOn 线程调度规则
- 了解背压策略与 Flowable 的使用时机
- 能够使用 RxJavaAndroid 与 Lifecycle 互操作并避免内存泄漏
本章定位:Kotlin 协程(coroutines)是 Google 官方推荐的异步方案,新项目首选。但 Java 项目无法使用协程,RxJava 仍是 Java + Views 项目异步处理的主流最佳实践。存量 Java 项目可继续使用 RxJava;当项目迁移到 Kotlin 后,可评估用协程替代 RxJava,或在过渡期两者并存(同一项目混用 RxJava 与协程是常见做法)。
与 Part 8-10 的对应关系:Kotlin 协程与 Flow 参见 第 43 章 协程与 Flow。
学习目标
- 掌握 RxJava 3 五种基础类型的语义与适用场景;
- 熟练使用
map/flatMap/filter/concatMap等常用操作符; - 理解
subscribeOn/observeOn线程调度规则; - 了解背压策略与
Flowable的使用时机; - 能够使用 RxJavaAndroid 与 Lifecycle 互操作并避免内存泄漏。
RxJava 3 异步编程
依赖配置
dependencies {
implementation "androidx.appcompat:appcompat:1.7.0"
// RxJava 3
implementation "io.reactivex.rxjava3:rxjava:3.1.8"
// RxJavaAndroid:提供 AndroidSchedulers.mainThread()
implementation "io.reactivex.rxjava3:rxandroid:3.0.2"
// 与 Lifecycle 互操作(可选)
implementation "androidx.lifecycle:lifecycle-runtime:2.8.4"
implementation "com.uber.rxdogtag:rxdogtag:2.0.0" // 全局未处理异常监控
}
五种基础类型
RxJava 3 提供 5 种基础可观察类型,按 数据数量 与 是否支持背压 划分:
| 类型 | 数据数量 | 背压 | 典型场景 |
|---|---|---|---|
Observable<T> |
0..N | 不支持 | UI 事件流、有限列表 |
Flowable<T> |
0..N | 支持 | 大数据流、IO 密集 |
Single<T> |
1(或错误) | 不支持 | 网络 GET 请求、单值查询 |
Maybe<T> |
0 或 1(或错误) | 不支持 | 可能为空的查询 |
Completable |
0(或错误) | 不支持 | 写入操作、删除操作 |
// Observable:多个事件
Observable<String> obs = Observable.just("a", "b", "c");
// Single:要么成功要么失败
Single<User> single = Single.fromCallable(() -> queryUser(1L));
// Maybe:可能为空
Maybe<User> maybe = Maybe.fromCallable(() -> {
User u = queryFromCache(id);
return u; // 可能为 null
});
// Completable:只关心成功 / 失败
Completable done = Completable.fromAction(() -> writeFile(path));
// Flowable:数据可能堆积,需要背压
Flowable<Integer> flow = Flowable.range(1, 1_000_000);
创建 Observable
// just
Observable<String> o1 = Observable.just("hello");
// fromIterable
Observable<Integer> o2 = Observable.fromIterable(Arrays.asList(1, 2, 3));
// range
Observable<Integer> o3 = Observable.range(1, 10);
// create(手写下游发射)
Observable<String> o4 = Observable.create(emitter -> {
emitter.onNext("a");
emitter.onNext("b");
emitter.onComplete();
});
// interval:周期发射
Observable<Long> ticker = Observable.interval(1, TimeUnit.SECONDS);
// fromCallable:包装同步调用
Observable<String> result = Observable.fromCallable(() -> fetchFromDb());
fromCallable 的好处是异常会被 try/catch 后封装为 onError,避免回调里抛出未受检异常导致崩溃。
常用操作符
map:一对一变换
Observable<Integer> lengths = Observable.just("hello", "world")
.map(String::length);
// 输出:5, 5
flatMap:一对多变平,并发不保序
Observable<String> result = Observable.just(1L, 2L)
.flatMap(id -> fetchUserAsync(id)); // 返回 Observable<User>,扁平为 User
// 注意:flatMap 的多个内部 Observable 是并发的,发射顺序不保证
concatMap:一对多变平,保序串行
Observable<User> users = Observable.just(1L, 2L, 3L)
.concatMap(id -> fetchUserAsync(id)); // 串行执行,保持 1->2->3 顺序
switchMap:切换上游时取消旧内部 Observable
// 典型场景:搜索框输入,新关键词覆盖旧请求
Disposable d = textChanges
.switchMap(this::search)
.subscribe(ui -> render(ui));
filter:过滤
Observable<Integer> evens = Observable.range(1, 10)
.filter(i -> i % 2 == 0);
其他常用
// distinct 去重
Observable<Integer> distinct = Observable.just(1, 2, 2, 3, 3).distinct();
// take 取前 N 个
Observable<Integer> first3 = Observable.range(1, 100).take(3);
// debounce 防抖(300ms 内多次发射只保留最后一个)
Observable<String> debounced = textChanges.debounce(300, TimeUnit.MILLISECONDS);
// combineLatest:多个源最新值组合
Observable.combineLatest(
sourceA.startWithItem(""),
sourceB.startWithItem(""),
(a, b) -> a + "/" + b
).subscribe(s -> render(s));
线程调度
subscribeOn / observeOn
Observable<User> pipeline = Observable.fromCallable(() -> fetchUser())
.subscribeOn(Schedulers.io()) // 上游在 IO 线程执行
.observeOn(AndroidSchedulers.mainThread()) // 下游在主线程
.map(this::toUiModel) // 主线程
.observeOn(Schedulers.computation()) // 再切换到计算线程
.map(this::heavyCompute);
关键规则:
subscribeOn只生效一次:决定onSubscribe/onNext的发射在哪个线程(多次调用只有第一次生效);observeOn可多次切换:每调用一次,下游就在该线程;- 默认所有操作在同一线程(即调用
subscribe的线程)。
AndroidSchedulers
RxAndroid 提供 AndroidSchedulers.mainThread(),是从 IO 切回 UI 线程的常用方式。Schedulers.io() 用于网络/磁盘 IO;Schedulers.computation() 用于 CPU 密集计算;Schedulers.single() 用于需要串行的场景。
完整异步加载示例
public class UserLoader {
public Observable<UserUiModel> loadUser(long id) {
return Observable.fromCallable(() -> api.fetchUser(id))
.subscribeOn(Schedulers.io())
.map(this::toUiModel)
.observeOn(AndroidSchedulers.mainThread());
}
private UserUiModel toUiModel(User user) {
return new UserUiModel(user.id, user.name, "用户:" + user.name);
}
}
背压策略与 Flowable
Observable 不支持背压。当上游高速发射、下游消费慢时,数据会堆积在内部队列里,最终抛 MissingBackpressureException 或 OOM。Flowable 引入 BackpressureStrategy 控制策略:
| 策略 | 行为 |
|---|---|
MISSING |
不缓冲也不报错,下游自行处理 |
ERROR |
内部队列溢出时抛 MissingBackpressureException |
BUFFER |
无限缓冲(小心 OOM) |
DROP |
队列满时丢弃新数据 |
LATEST |
仅保留最新一个 |
Flowable<Integer> f = Flowable.create(emitter -> {
for (int i = 0; i < 1_000_000; i++) {
emitter.onNext(i);
}
emitter.onComplete();
}, BackpressureStrategy.DROP);
f.observeOn(Schedulers.computation())
.subscribe(i -> System.out.println(i));
也可以用 onBackpressureBuffer(), onBackpressureDrop(), onBackpressureLatest() 在已有 Flowable 上调整策略。
判断是否需要 Flowable:上游连续发射大量数据且速度 > 下游消费速度,用 Flowable;其余典型 Android 场景(事件、单次网络请求)用 Observable / Single 即可。
与 Android Lifecycle 互操作
RxJavaAndroid
RxAndroid 的核心就是 AndroidSchedulers.mainThread()。配合 ViewModel / Activity 使用时要注意生命周期管理。
内存泄漏与 dispose
subscribe() 返回 Disposable,未 dispose 的订阅会持有外部引用,Activity 销毁后回调仍可能执行,导致:
- 持有已销毁的 Activity View 引用,内存泄漏;
- 回调里更新已销毁 UI 崩溃。
正确做法是在 Activity / Fragment 销毁时 dispose:
public class UserActivity extends AppCompatActivity {
private final CompositeDisposable disposables = new CompositeDisposable();
private UserLoader loader;
@Override
protected void onCreate(Bundle savedInstanceState) {
super.onCreate(savedInstanceState);
loader = new UserLoader();
Disposable d = loader.loadUser(1L)
.subscribe(ui -> render(ui),
err -> showError(err));
disposables.add(d);
}
private void render(UserUiModel ui) { /* ... */ }
private void showError(Throwable e) { /* ... */ }
@Override
protected void onDestroy() {
super.onDestroy();
disposables.clear(); // 取消所有订阅
}
}
CompositeDisposable
CompositeDisposable 是一个容器,可批量添加 / 清空:
CompositeDisposable cd = new CompositeDisposable();
cd.add(d1);
cd.add(d2);
cd.add(d3);
// 批量取消
cd.clear(); // 清空并 dispose 全部
cd.dispose(); // 同上但不可再用
ViewModel 中也应在 onCleared 中清理:
public class UserViewModel extends ViewModel {
private final CompositeDisposable disposables = new CompositeDisposable();
private final UserLoader loader = new UserLoader();
private final MutableLiveData<UserUiModel> data = new MutableLiveData<>();
public LiveData<UserUiModel> getData() { return data; }
public void load(long id) {
Disposable d = loader.loadUser(id)
.subscribe(data::postValue,
err -> data.postValue(null));
disposables.add(d);
}
@Override
protected void onCleared() {
super.onCleared();
disposables.clear();
}
}
错误处理
loader.loadUser(id)
.onErrorResumeNext(err -> {
if (err instanceof IOException) {
return loadFromCache(id); // 网络错误降级到缓存
}
return Observable.error(err);
})
.onErrorReturnItem(UserUiModel.empty())
.subscribe(this::render, this::showError);
常用错误处理算子:onErrorReturn / onErrorReturnItem / onErrorResumeNext / onErrorComplete / retry。
完整实战:搜索 + 防抖
下面是一个搜索框 + 防抖 + 网络请求 + 错误兜底的完整示例:
public class SearchViewModel extends ViewModel {
private final PublishProcessor<String> query = PublishProcessor.create();
private final MutableLiveData<List<String>> results = new MutableLiveData<>();
private final MutableLiveData<Boolean> loading = new MutableLiveData<>(false);
private final MutableLiveData<String> error = new MutableLiveData<>();
private final CompositeDisposable disposables = new CompositeDisposable();
private final SearchApi api = new SearchApi();
private final Scheduler io = Schedulers.io();
private final Scheduler main = AndroidSchedulers.mainThread();
public SearchViewModel() {
Disposable d = query
.debounce(300, TimeUnit.MILLISECONDS, io) // 防抖
.filter(s -> !s.isEmpty())
.switchMap(this::doSearch) // 切换上游取消旧请求
.subscribe(res -> {
loading.setValue(false);
results.setValue(res);
}, err -> {
loading.setValue(false);
error.setValue(err.getMessage());
});
disposables.add(d);
}
private Flowable<List<String>> doSearch(String keyword) {
loading.setValue(true);
return api.search(keyword)
.subscribeOn(io)
.observeOn(main)
.onErrorReturnItem(Collections.emptyList())
.toFlowable();
}
public void setQuery(String q) { query.onNext(q); }
public LiveData<List<String>> getResults() { return results; }
public LiveData<Boolean> getLoading() { return loading; }
public LiveData<String> getError() { return error; }
@Override
protected void onCleared() {
super.onCleared();
disposables.clear();
}
}
这里 switchMap 关键作用:用户连续输入时,新关键词会取消上一次未完成的搜索请求,避免结果乱序。
常见坑与最佳实践
- 忘记 dispose:未 dispose 的订阅持有 Activity 引用,是 RxJava 内存泄漏的最常见原因。
CompositeDisposable配合onDestroy/onCleared是必备手段。 - 在主线程调用 fromCallable 做 IO:默认上游在
subscribe调用线程,未指定subscribeOn(Schedulers.io())时会在主线程执行网络/磁盘操作,触发 ANR。 - 链式调用顺序错误:
observeOn只影响 下游 操作,所以放在错误的中间位置会让后续操作意外换线程。建议在每次observeOn后用注释标注当前线程。 - Observable 用于大数据流:高速数据源用 Observable 而非 Flowable,会引发
MissingBackpressureException,明确用 Flowable + 背压策略。 - flatMap 乱序:依赖顺序的场景必须用
concatMap而非flatMap。 - subscribe 时异常未处理:未提供
onError的subscribe会因 RxJavaPlugins 全局处理抛OnErrorNotImplementedException直接崩溃。生产环境务必提供错误回调。 - 在 ViewModel 中持有 Activity 引用:ViewModel 比 Activity 生命周期长,传入 Activity 一定泄漏,所有结果都通过 LiveData 传给 UI。
- subscribeOn 多次调用:只有第一次生效,多次写只会让代码迷惑;要在不同段切换线程用
observeOn。
章节小结
本章覆盖了 RxJava 3 的五种基础类型、常用操作符(map / flatMap / concatMap / switchMap / filter)、线程调度(subscribeOn 只生效一次、observeOn 可多次切换)、背压策略与 Flowable 的使用时机、与 Android Lifecycle 互操作以及内存泄漏防护(Disposable / CompositeDisposable)。核心要点:异步链路必须可取消,UI 回调必须切到主线程并随生命周期释放订阅。
下一章预告
下一章将进入测试领域:测试金字塔、JUnit4(@Test / @Before / @After / @Rule)、Mockito(mock / when / verify)、Espresso UI 测试(onView / perform / check)、Room 测试、Robolectric 介绍与测试覆盖率分析。