RxJava 响应式编程详解:原理、操作符与实战

举报
shenlan9755 发表于 2026/09/01 09:12:57 2026/09/01
【摘要】 RxJava 响应式编程详解:原理、操作符与实战 前言RxJava 是 Reactive Extensions 在 JVM 上的实现,它把异步数据流抽象为可观察序列,用统一的操作符编排同步与异步逻辑。在安卓、后端网关、实时数据处理等场景,RxJava 是处理复杂异步流的利器。本文将从响应式编程思想、核心机制、操作符体系、背压、线程调度、实战模式等维度,系统讲解 RxJava。 一、响应式编...

RxJava 响应式编程详解:原理、操作符与实战

前言

RxJava 是 Reactive Extensions 在 JVM 上的实现,它把异步数据流抽象为可观察序列,用统一的操作符编排同步与异步逻辑。在安卓、后端网关、实时数据处理等场景,RxJava 是处理复杂异步流的利器。本文将从响应式编程思想、核心机制、操作符体系、背压、线程调度、实战模式等维度,系统讲解 RxJava。


一、响应式编程思想

1.1 命令式 vs 响应式

命令式:主动拉取数据,控制流由调用者决定。

List<User> users = userRepository.findAll();
for (User user : users) {
    System.out.println(user.getName());
}

响应式:数据源主动推送,消费者被动响应。

userRepository.findAll()
    .subscribe(user -> System.out.println(user.getName()));

1.2 核心特征

特征 说明
数据流 一切皆流,事件、UI 操作、网络响应都可建模为流
异步非阻塞 通过回调或协程传递结果,不阻塞调用线程
声明式组合 用操作符链式声明数据变换,而非命令式循环
推模式 数据源主动推送,消费者订阅
背压支持 消费慢于生产时,有机制协调速率

1.3 应用场景

  • 多个异步源合并(网络 + 缓存 + 数据库)。
  • UI 事件流(点击、输入、滚动)防抖、节流。
  • 实时数据流(WebSocket、传感器、股票行情)。
  • 复杂状态机(重试、超时、条件分支)。
  • 数据流水线(ETL、过滤、映射、聚合)。

二、核心概念

2.1 Observable 与 Observer

Observable 是数据源,Observer 是消费者。订阅时建立关系:

Observable  ──onNext(t)──→  Observer
            ──onNext(t)──→
            ──onError(e)─→  (终止)
            ──onComplete() (终止)

2.2 基本示例

Observable<String> observable = Observable.create(emitter -> {
    emitter.onNext("Hello");
    emitter.onNext("RxJava");
    emitter.onComplete();
});

observable.subscribe(
    item -> System.out.println("收到: " + item),
    error -> System.err.println("错误: " + error),
    () -> System.out.println("完成")
);

输出:

收到: Hello
收到: RxJava
完成

2.3 Observable 的变体

类型 发射 背压 适用
Observable 0…N 事件流、UI 事件
Flowable 0…N 大数据量、生产快消费慢
Single 1 单次请求(网络调用)
Maybe 0 或 1 可能无结果的查询
Completable 0(仅完成/错误) 副作用操作(写入、删除)

2.4 Single 示例

Single<User> single = Single.fromCallable(() -> fetchUserFromNetwork());
single.subscribe(
    user -> System.out.println("用户: " + user),
    error -> System.err.println("失败: " + error)
);

2.5 Completable 示例

Completable completable = Completable.fromAction(() -> saveToDatabase(user));
completable.subscribe(
    () -> System.out.println("保存成功"),
    error -> System.err.println("保存失败: " + error)
);

三、生命周期与 Disposable

3.1 订阅返回 Disposable

Disposable disposable = observable.subscribe(item -> handle(item));

// 在合适时机取消订阅,避免内存泄漏
disposable.dispose();

3.2 CompositeDisposable 批量管理

public class UserManager {
    private final CompositeDisposable disposables = new CompositeDisposable();

    public void init() {
        disposables.add(loadProfile());
        disposables.add(loadSettings());
        disposables.add(loadNotifications());
    }

    public void destroy() {
        disposables.dispose();  // 一次性取消所有订阅
    }
}

关键:在 Activity/Controller 销毁时务必 dispose,否则异步回调持有引用导致泄漏。


四、创建操作符

4.1 just / fromArray / fromIterable

Observable.just(1, 2, 3).subscribe(System.out::println);

Observable.fromArray("A", "B", "C").subscribe(System.out::println);

List<Integer> list = Arrays.asList(1, 2, 3);
Observable.fromIterable(list).subscribe(System.out::println);

4.2 range

Observable.range(1, 5)  // 1,2,3,4,5
    .subscribe(System.out::println);

4.3 interval(周期发射)

Observable.interval(1, TimeUnit.SECONDS)
    .subscribe(i -> System.out.println("tick: " + i));

Thread.sleep(5000);  // 输出 tick: 0 ~ tick: 4

4.4 timer(延迟发射一次)

Observable.timer(3, TimeUnit.SECONDS)
    .subscribe(i -> System.out.println("3 秒后触发"));

4.5 defer(延迟创建)

每次订阅时才创建 Observable,保证拿到最新状态:

AtomicInteger counter = new AtomicInteger(0);

Observable<Integer> deferred = Observable.defer(() ->
    Observable.just(counter.get())
);

counter.set(10);
deferred.subscribe(System.out::println);  // 10

counter.set(20);
deferred.subscribe(System.out::println);  // 20

4.6 fromCallable(包装阻塞调用)

Observable<User> userObs = Observable.fromCallable(() -> {
    Thread.sleep(1000);
    return new User("张三");
});

fromCallable 会捕获异常并经 onError 传递,比 create 更安全。


五、变换操作符

5.1 map(一对一变换)

Observable.just("1", "2", "3")
    .map(Integer::parseInt)
    .map(i -> i * 10)
    .subscribe(System.out::println);  // 10, 20, 30

5.2 flatMap(一对多展开,无序)

Observable.just("A", "B")
    .flatMap(letter -> Observable.just(letter + "1", letter + "2"))
    .subscribe(System.out::println);
// 可能输出: A1 A2 B1 B2 或交错

5.3 concatMap(一对多展开,保序)

Observable.just("A", "B")
    .concatMap(letter -> Observable.just(letter + "1", letter + "2"))
    .subscribe(System.out::println);
// 一定输出: A1 A2 B1 B2

5.4 flatMap vs concatMap vs switchMap

操作符 顺序 内部流切换行为 适用
flatMap 不保证 旧流继续发射 并行、不关心顺序
concatMap 保证 等旧流完成再订阅新流 严格顺序
switchMap 保证 取消旧流,只关心最新 搜索联想、最新数据

5.5 switchMap 搜索联想示例

searchInputEmitter
    .debounce(300, TimeUnit.MILLISECONDS)
    .switchMap(keyword -> searchService.search(keyword))
    .subscribe(results -> updateUI(results));

用户快速输入时,switchMap 自动取消旧请求,只保留最新结果,避免乱序覆盖。

5.6 scan(累加)

Observable.just(1, 2, 3, 4)
    .scan(0, Integer::sum)
    .subscribe(System.out::println);  // 0, 1, 3, 6, 10

5.7 groupBy(分组)

Observable.just(1, 2, 3, 4, 5, 6)
    .groupBy(i -> i % 2 == 0 ? "偶" : "奇")
    .flatMap(group -> group.collect(Collectors.toList()).toObservable()
        .map(list -> group.getKey() + ": " + list))
    .subscribe(System.out::println);
// 偶: [2, 4, 6]
// 奇: [1, 3, 5]

六、过滤操作符

6.1 filter

Observable.range(1, 10)
    .filter(i -> i % 2 == 0)
    .subscribe(System.out::println);  // 2, 4, 6, 8, 10

6.2 take / skip

Observable.range(1, 10)
    .skip(3)
    .take(2)
    .subscribe(System.out::println);  // 4, 5

6.3 distinct / distinctUntilChanged

Observable.just(1, 1, 2, 2, 3, 1, 1)
    .distinct()
    .subscribe(System.out::println);  // 1, 2, 3

Observable.just(1, 1, 2, 2, 3, 1, 1)
    .distinctUntilChanged()
    .subscribe(System.out::println);  // 1, 2, 3, 1

6.4 debounce(去抖动)

clickEmitter
    .debounce(500, TimeUnit.MILLISECONDS)
    .subscribe(event -> saveDraft());

500ms 内的连续点击只保留最后一次,适合保存草稿、搜索触发。

6.5 throttleFirst / throttleLast(节流)

clickEmitter
    .throttleFirst(1, TimeUnit.SECONDS)
    .subscribe(event -> doAction());

每秒只响应第一次点击,防连点。

6.6 sample(采样)

sensorEmitter
    .sample(100, TimeUnit.MILLISECONDS)
    .subscribe(value -> record(value));

每 100ms 取一个最新值,适合高频传感器降采样。


七、组合操作符

7.1 merge(交错合并)

Observable<String> obs1 = Observable.interval(1, TimeUnit.SECONDS)
    .map(i -> "A" + i);
Observable<String> obs2 = Observable.interval(500, TimeUnit.MILLISECONDS)
    .map(i -> "B" + i);

Observable.merge(obs1, obs2)
    .subscribe(System.out::println);

两个流的事件按到达时间交错输出。

7.2 concat(顺序连接)

Observable.concat(
    Observable.just("A1", "A2"),
    Observable.just("B1", "B2")
).subscribe(System.out::println);
// A1 A2 B1 B2

第一个流完成后再订阅第二个,保证顺序。

7.3 zip(按位置配对)

Observable.zip(
    Observable.just("张三", "李四", "王五"),
    Observable.just(25, 30, 28),
    (name, age) -> name + "(" + age + "岁)"
).subscribe(System.out::println);
// 张三(25岁) 李四(30岁) 王五(28岁)

zip 等待所有流都发射第 N 项后才组合,慢的流会拖慢整体。

7.4 combineLatest(最新值组合)

BehaviorProcessor<Integer> obs1 = BehaviorProcessor.createDefault(0);
BehaviorProcessor<Integer> obs2 = BehaviorProcessor.createDefault(0);

Observable.combineLatest(obs1, obs2, (a, b) -> a + b)
    .subscribe(sum -> System.out.println("sum: " + sum));

obs1.onNext(1);  // sum: 1
obs2.onNext(2);  // sum: 3
obs1.onNext(5);  // sum: 7

任一流发射新值,就用各流最新值组合,适合表单实时校验。

7.5 startWith

networkUserObservable
    .startWith(cacheUserObservable)
    .subscribe(user -> render(user));

先显示缓存,再用网络数据覆盖。


八、错误处理操作符

8.1 onErrorReturn(出错返回默认值)

fetchUserObservable
    .onErrorReturn(error -> User.EMPTY)
    .subscribe(user -> render(user));

8.2 onErrorResumeNext(出错切换到备用流)

networkObservable
    .onErrorResumeNext(error -> cacheObservable)
    .subscribe(user -> render(user));

网络失败回退到缓存。

8.3 retry(重试)

fetchObservable
    .retry(3)  // 最多重试 3 次
    .subscribe(user -> render(user));

8.4 retryWhen(条件重试)

fetchObservable
    .retryWhen(errors -> errors
        .zipWith(Observable.range(1, 3), (error, retryCount) -> retryCount)
        .flatMap(retryCount -> Observable.timer((long) Math.pow(2, retryCount), TimeUnit.SECONDS))
    )
    .subscribe(user -> render(user));

指数退避重试:1s、2s、4s。

8.5 timeout

fetchObservable
    .timeout(5, TimeUnit.SECONDS, fallbackObservable)
    .subscribe(user -> render(user));

5 秒未返回则切换到 fallback。


九、线程调度

9.1 Schedulers

调度器 说明
Schedulers.io() IO 密集(网络、磁盘),线程池缓存复用
Schedulers.computation() CPU 密集(计算、编解码),固定线程数 = CPU 核数
Schedulers.newThread() 每次新建线程
Schedulers.single() 单线程串行
Schedulers.trampoline() 当前线程排队执行
AndroidSchedulers.mainThread() 安卓主线程(需 RxAndroid)

9.2 subscribeOn / observeOn

userRepository.getUser()
    .subscribeOn(Schedulers.io())       // 上游在 IO 线程执行
    .observeOn(AndroidSchedulers.mainThread())  // 下游在主线程接收
    .subscribe(user -> renderUI(user));
  • subscribeOn:指定数据源发射所在的线程,只第一次生效。
  • observeOn:指定其后操作符执行的线程,可多次切换。

9.3 多次切换线程

Observable.fromCallable(() -> readFromDisk())       // IO
    .subscribeOn(Schedulers.io())
    .map(data -> parseAndValidate(data))             // IO
    .observeOn(Schedulers.computation())
    .map(parsed -> heavyCompute(parsed))             // computation
    .observeOn(Schedulers.io())
    .map(result -> writeToCache(result))             // IO
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(finalResult -> renderUI(finalResult)); // main

十、背压(Backpressure)

10.1 问题

当生产者发射速度远超消费者处理速度时,未消费的事件在内部队列堆积,最终 OOM。

Observable.range(1, 1_000_000)
    .observeOn(Schedulers.io())
    .map(i -> {
        Thread.sleep(100);  // 消费慢
        return i;
    })
    .subscribe();
// 内部队列堆积 100 万项,可能 OOM

10.2 Flowable 与背压策略

Flowable 是支持背压的 Observable,订阅时通过 Subscription.request(n) 主动拉取。

Flowable.range(1, 1_000_000)
    .observeOn(Schedulers.io())
    .map(i -> {
        Thread.sleep(100);
        return i;
    })
    .subscribe(new Subscriber<Integer>() {
        private Subscription subscription;

        @Override
        public void onSubscribe(Subscription s) {
            this.subscription = s;
            s.request(100);  // 先请求 100 项
        }

        @Override
        public void onNext(Integer i) {
            handle(i);
            subscription.request(100);  // 处理完再请求下一批
        }

        @Override
        public void onError(Throwable e) {}

        @Override
        public void onComplete() {}
    });

10.3 背压策略

ObservableFlowable 时需指定策略:

策略 行为
onBackpressureBuffer() 无限缓冲(默认,可能 OOM)
onBackpressureBuffer(capacity) 有界缓冲,超出调用 onOverflow
onBackpressureDrop() 丢弃无法及时处理的事件
onBackpressureLatest() 只保留最新一个事件
hotObservable
    .onBackpressureBuffer(1000, () -> log.warn("溢出"), BackpressureOverflowStrategy.DROP_OLDEST)
    .observeOn(Schedulers.io())
    .subscribe(this::process);

10.4 何时用 Flowable

  • 数据源是"冷"流,且数据量大(如数据库游标、文件行)。
  • 生产速率可能超过消费速率。
  • 需要显式控制拉取节奏。

反例:UI 事件、网络单次请求用 Observable/Single 即可,无需 Flowable。


十一、Subject / Processor:既是源又是消费者

11.1 PublishSubject

广播给所有订阅者,订阅后才能收到事件。

PublishSubject<String> subject = PublishSubject.create();

subject.subscribe(s -> System.out.println("订阅者1: " + s));
subject.onNext("Hello");  // 订阅者1: Hello

subject.subscribe(s -> System.out.println("订阅者2: " + s));
subject.onNext("World");  // 订阅者1: World / 订阅者2: World

11.2 BehaviorSubject

缓存最新一个值,新订阅者立即收到。

BehaviorSubject<Integer> subject = BehaviorSubject.createDefault(0);

subject.onNext(1);
subject.subscribe(v -> System.out.println("新订阅者: " + v));  // 新订阅者: 1

适合表示"状态"(如当前用户、当前选中项)。

11.3 ReplaySubject

缓存所有历史事件,新订阅者回放全部。

ReplaySubject<String> subject = ReplaySubject.create();
subject.onNext("A");
subject.onNext("B");
subject.subscribe(s -> System.out.println(s));  // A B

11.4 Processor(支持背压的 Subject)

FlowableProcessor 子类(如 BehaviorProcessorPublishProcessor)兼具背压与 Subject 特性,适合做状态总线。


十二、实战模式

12.1 缓存 + 网络数据合并

public Observable<User> getUserWithCache(String id) {
    return Observable.concat(
        cacheRepository.getUser(id).toObservable(),  // 先缓存
        networkRepository.getUser(id)
            .doOnNext(user -> cacheRepository.save(user))
            .toObservable()
    ).firstOrError().toObservable();
}

firstOrError 保证缓存命中就不再请求网络。

12.2 轮询 + 条件停止

public Observable<Task> pollTaskUntilDone(String taskId) {
    return Observable.interval(0, 5, TimeUnit.SECONDS)
        .flatMap(i -> taskService.getTask(taskId).toObservable())
        .takeUntil(task -> task.isDone())
        .filter(task -> task.isDone());
}

每 5 秒轮询,任务完成后停止并发射最终结果。

12.3 多接口并行合并

public Observable<HomePage> loadHomePage() {
    return Observable.zip(
        apiService.getBanners(),
        apiService.getProducts(),
        apiService.getRecommendations(),
        (banners, products, recommendations) ->
            new HomePage(banners, products, recommendations)
    ).subscribeOn(Schedulers.io())
     .observeOn(AndroidSchedulers.mainThread());
}

三个接口并行请求,全部完成后合并渲染,比串行快 3 倍。

12.4 级联请求

public Observable<OrderDetail> loadOrderDetail(String orderId) {
    return apiService.getOrder(orderId)
        .flatMap(order -> apiService.getUser(order.getUserId()),
            (order, user) -> new OrderDetail(order, user))
        .flatMap(detail -> apiService.getLogistics(detail.getOrder().getId()),
            (detail, logistics) -> {
                detail.setLogistics(logistics);
                return detail;
            });
}

订单 → 用户 → 物流,串行级联,每步用上一步结果。

12.5 表单实时校验

BehaviorProcessor<String> nameEmitter = BehaviorProcessor.createDefault("");
BehaviorProcessor<String> emailEmitter = BehaviorProcessor.createDefault("");

Observable.combineLatest(
    nameEmitter.map(this::validateName),
    emailEmitter.map(this::validateEmail),
    (nameValid, emailValid) -> nameValid && emailValid
).distinctUntilChanged()
 .subscribe(valid -> submitButton.setEnabled(valid));

任一字段变化都重新校验,按钮状态自动跟随。

12.6 防抖搜索 + 节流点击 + 生命周期

Disposable disposable = searchView.textChanges()
    .debounce(300, TimeUnit.MILLISECONDS)
    .map(CharSequence::toString)
    .filter(text -> text.length() >= 2)
    .switchMap(text -> searchService.search(text)
        .subscribeOn(Schedulers.io())
        .onErrorResumeNext(Observable.empty()))
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(results -> adapter.update(results));

// Activity 销毁时
disposable.dispose();

完整链路:去抖 → 过滤 → 取消旧请求 → 网络请求 → 错误吞掉 → 主线程渲染 → 销毁取消。


十三、自定义操作符

13.1 用 lift 自定义

Observable<Integer> observable = Observable.just(1, 2, 3);

observable.lift(observer -> new DisposableObserver<Integer>() {
    @Override
    public void onNext(Integer value) {
        observer.onNext(value * 100);  // 自定义变换
    }

    @Override
    public void onError(Throwable e) {
        observer.onError(e);
    }

    @Override
    public void onComplete() {
        observer.onComplete();
    }
}).subscribe(System.out::println);  // 100, 200, 300

13.2 用 compose 复用逻辑

Transformer<Integer, Integer> doubleAndLog = upstream -> upstream
    .map(i -> i * 2)
    .doOnNext(i -> System.out.println("处理后: " + i));

Observable.just(1, 2, 3)
    .compose(doubleAndLog)
    .subscribe();
// 处理后: 2 / 处理后: 4 / 处理后: 6

compose 让操作符组合可复用,比在每处重复链式更清晰。


十四、测试

14.1 TestObserver

TestObserver<Integer> testObserver = Observable.just(1, 2, 3)
    .map(i -> i * 10)
    .test();

testObserver
    .assertNoErrors()
    .assertValues(10, 20, 30)
    .assertComplete();

14.2 TestScheduler 控制时间

TestScheduler scheduler = new TestScheduler();

TestObserver<Long> test = Observable.interval(1, TimeUnit.SECONDS, scheduler)
    .test();

test.assertEmpty();
scheduler.advanceTimeBy(2, TimeUnit.SECONDS);
test.assertValues(0L, 1L);

无需真实等待,单元测试可瞬时验证时间相关逻辑。


十五、性能与陷阱

15.1 常见陷阱

  1. 忘记 dispose:订阅未取消,异步回调持有 Activity/Controller 引用,内存泄漏。
  2. 滥用 flatMap 导致并行失控flatMap 默认并发 Integer.MAX_VALUE,可能压垮下游。用 flatMap(func, maxConcurrency) 限制:
ids.flatMap(id -> fetchDetail(id), 10)  // 最多并发 10
  1. 在 onError 里抛异常:会包装成 OnErrorNotImplementedException 导致崩溃。务必提供 onError 回调。
  2. subscribeOn 多次无效:只有最上游的 subscribeOn 生效,多次调用只第一次有效。
  3. hot 流用背压策略不当PublishSubject 转 Flowable 默认无界缓冲,可能 OOM,应显式 onBackpressureDrop
  4. 操作符链过长:可读性下降,建议用 compose 拆分为可命名 Transformer。

15.2 性能建议

  • IO 用 Schedulers.io(),CPU 密集用 Schedulers.computation(),避免线程池误用。
  • 大数据流用 Flowable + 背压,而非 Observable
  • flatMap 限制并发数,避免压垮下游与数据库连接池。
  • 频繁创建的小流可复用 Transformer,减少对象分配。
  • 监控订阅数,防止"幽灵订阅"长期占用资源。

十六、RxJava 与 Kotlin Flow 对比

维度 RxJava Kotlin Flow
抽象 Observable/Flowable Flow(协程基础)
背压 显式(request/策略) 隐式(suspend)
取消 Disposable 协程取消
错误 onError 回调 异常抛出
操作符 极丰富(200+) 较少但够用
学习曲线 平缓
协程集成 需桥接(await/asFlow 原生
调试 链式调用栈较难 协程栈较友好

建议:新项目优先 Kotlin Coroutine + Flow;存量 RxJava 项目可逐步迁移关键路径,或继续沿用 RxJava(生态成熟、操作符丰富)。


十七、总结

RxJava 的核心价值在于用统一的操作符语言编排异步数据流,把回调地狱、状态同步、错误重试等复杂逻辑转化为声明式的链式调用。

  • 核心抽象:Observable/Flowable/Single/Completable,对应不同数据形态。
  • 操作符体系:创建、变换、过滤、组合、错误处理、线程调度,覆盖绝大多数异步场景。
  • 背压:Flowable + request 机制协调生产消费速率,避免 OOM。
  • 实战模式:缓存合并、轮询、并行、级联、表单校验、搜索联想,皆可优雅实现。
  • 陷阱:忘记 dispose、并发失控、线程池误用、hot 流背压,需时刻警惕。

响应式编程不是银弹,但在异步密集、数据流复杂的场景,RxJava 仍是 Java 生态中表达力最强、生态最成熟的方案。掌握其思维模型,比记住操作符列表更重要。

【声明】本内容来自华为云开发者社区博主,不代表华为云及华为云开发者社区的观点和立场。转载时必须标注文章的来源(华为云社区)、文章链接、文章作者等基本信息,否则作者和本社区有权追究责任。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0

0/1000
抱歉,系统识别当前为高风险访问,暂不支持该操作

全部回复

上滑加载中

设置昵称

在此一键设置昵称,即可参与社区互动!

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。