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 背压策略
将 Observable 转 Flowable 时需指定策略:
| 策略 | 行为 |
|---|---|
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 子类(如 BehaviorProcessor、PublishProcessor)兼具背压与 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 常见陷阱
- 忘记 dispose:订阅未取消,异步回调持有 Activity/Controller 引用,内存泄漏。
- 滥用 flatMap 导致并行失控:
flatMap默认并发Integer.MAX_VALUE,可能压垮下游。用flatMap(func, maxConcurrency)限制:
ids.flatMap(id -> fetchDetail(id), 10) // 最多并发 10
- 在 onError 里抛异常:会包装成
OnErrorNotImplementedException导致崩溃。务必提供onError回调。 - subscribeOn 多次无效:只有最上游的
subscribeOn生效,多次调用只第一次有效。 - hot 流用背压策略不当:
PublishSubject转 Flowable 默认无界缓冲,可能 OOM,应显式onBackpressureDrop。 - 操作符链过长:可读性下降,建议用
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 生态中表达力最强、生态最成熟的方案。掌握其思维模型,比记住操作符列表更重要。
- 点赞
- 收藏
- 关注作者
评论(0)