A
第 29 章JAVA45 分钟

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 关键作用:用户连续输入时,新关键词会取消上一次未完成的搜索请求,避免结果乱序。

常见坑与最佳实践

  1. 忘记 dispose:未 dispose 的订阅持有 Activity 引用,是 RxJava 内存泄漏的最常见原因。CompositeDisposable 配合 onDestroy / onCleared 是必备手段。
  2. 在主线程调用 fromCallable 做 IO:默认上游在 subscribe 调用线程,未指定 subscribeOn(Schedulers.io()) 时会在主线程执行网络/磁盘操作,触发 ANR。
  3. 链式调用顺序错误:observeOn 只影响 下游 操作,所以放在错误的中间位置会让后续操作意外换线程。建议在每次 observeOn 后用注释标注当前线程。
  4. Observable 用于大数据流:高速数据源用 Observable 而非 Flowable,会引发 MissingBackpressureException,明确用 Flowable + 背压策略。
  5. flatMap 乱序:依赖顺序的场景必须用 concatMap 而非 flatMap。
  6. subscribe 时异常未处理:未提供 onError 的 subscribe 会因 RxJavaPlugins 全局处理抛 OnErrorNotImplementedException 直接崩溃。生产环境务必提供错误回调。
  7. 在 ViewModel 中持有 Activity 引用:ViewModel 比 Activity 生命周期长,传入 Activity 一定泄漏,所有结果都通过 LiveData 传给 UI。
  8. 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 介绍与测试覆盖率分析。