RxJava 与响应式编程
RxJava 仍常见于存量 Android 项目. 面试重点不是背操作符, 而是能解释 RxJava 2/3 版本边界, 流语义, 背压, 资源释放, 错误终止以及向协程/Flow 的渐进迁移.
前置知识与示例说明: 先理解线程, 异常和集合基础 (见 Java 与 JVM 基础); 与协程 Flow 的冷热流, 取消语义对照见协程与 Flow. 本文代码均为上下文片段: 须按项目的 RxJava 2 或 RxJava 3 主版本补齐依赖, Android scheduler, api 和 UI 方法, 不能混用包名后声称可运行.
学习目标: 能根据事件是否可丢, 是否有序, 消费是否跟得上来选择 Rx 类型与操作符, 且能给出订阅释放和错误终止的验证方法.
一, 版本与类型选择
RxJava 2 与 RxJava 3 的核心模型相近, 但包名, 依赖坐标和生态 adapter 不兼容. 维护项目时先确认实际主版本及 RxAndroid/Retrofit adapter 的匹配关系, 不要在同一业务链路中无计划混用.
| 类型 | 发射约束 | 常见场景 |
|---|---|---|
Observable<T> | 0..N, 无背压协议 | UI 事件, 低频流 |
Flowable<T> | 0..N, 有背压协议 | 高频或跨异步边界的数据流 |
Single<T> | 1 个成功值或错误 | 必须有结果的请求 |
Maybe<T> | 0 或 1 个成功值或错误 | 可为空的缓存查询 |
Completable | 完成或错误, 无值 | 写入, 删除, 同步动作 |
类型选择应表达业务基数, 不要把所有接口都包装为 Observable.
二, 变换操作符的语义差异
map: 一对一同步变换; mapper 返回普通值.flatMap: 一对多并发合并, 结果顺序通常不保证与输入一致.concatMap: 按上游顺序串行订阅内部流, 以吞吐换顺序.switchMap: 新值到来时取消 / 切换旧内部流, 适合搜索联想等 “只关心最新结果” 的场景.zip: 按位置配对, 任一源无法继续配对时可能结束.combineLatest: 各源至少发射一次后, 任一源更新都与其他源最新值组合.
面试回答要给出并发, 顺序, 取消和错误传播差异, 不能只说 “都是把一个流变成另一个流”.
三, 调度器与线程边界
subscribeOn 影响订阅及其上游执行位置, 多次调用时通常由最靠近源且实际生效的调度决定; observeOn 从出现位置起切换下游, 可以多次使用. 不要把 “第一个 / 最后一个生效” 背成脱离具体 operator 的绝对口诀.
val disposable = api.getUser()
.subscribeOn(Schedulers.io())
.map(::toUiModel)
.observeOn(AndroidSchedulers.mainThread())
.doFinally { hideLoading() }
.subscribe(::render, ::showError)
doFinally 在完成, 错误或 dispose 后都会执行, 适合对称收尾; doOnComplete 不覆盖错误和取消. 重 CPU 计算应使用合适 scheduler, 不能因为上游是网络请求就把所有变换都放在主线程.
四, 背压不是 “加个 BUFFER”
当生产速度持续高于消费速度时, 无界缓存会把速度问题转成内存和延迟问题. Flowable 策略要按业务语义选择:
BUFFER: 不丢数据的诉求不等于天然有界; 必须明确容量, 上限和溢出策略.DROP: 丢弃无法及时消费的数据, 适合允许采样的场景.LATEST: 只保留最新值, 适合状态刷新, 不适合交易或消息流水.ERROR: 让过载显式失败, 由上层降级或限流.
还应优先考虑上游限速, 批处理, 窗口, 采样和减少无效工作. 背压策略必须与 “是否允许丢失, 是否要求顺序, 最大延迟” 一起回答.
Subject 家族: 热源, 状态与生命周期
Subject 同时是 Observer 和 Observable, 常被用作桥接回调的热源; 它不是默认的全局事件总线. 应优先把状态放在单一 owner, 避免任意位置都能 onNext 造成来源不可追踪. 大多数 Subject 的 onNext/onError/onComplete 不可由多个线程并发调用; 若确实有多个生产线程, 创建后立即用 subject.toSerialized() 暴露串行化入口, 并仍由一个 owner 定义事件顺序. 终止后不得再发射: onComplete/onError 是终止信号, 后续 onNext 无效; 未被订阅者消费的错误还可能进入全局错误处理器, 因此必须在 source 边界处理.
| Subject | 新订阅者看到什么 | 适合什么 | 主要风险 |
|---|---|---|---|
PublishSubject | 订阅后的新事件 | 纯瞬时且允许错过的局部信号 | 订阅前事件丢失 |
BehaviorSubject | 最新一个值 (若尚未终止) | 有当前状态的热源 | 不接受 null; 把事件误当状态会重放旧动作 |
ReplaySubject | 历史缓存 | 明确需要回放的有限历史 | 无界 replay 会增长内存 |
AsyncSubject | 完成前最后一个值 | 只关心最终结果的特殊桥接 | 未完成就不会发射 |
// 上下文片段:RxJava 3 写法;RxJava 2 的包名不同.
private val searchInput = PublishSubject.create<String>().toSerialized()
val results = searchInput
.map(String::trim)
.debounce(300, TimeUnit.MILLISECONDS)
.distinctUntilChanged()
.switchMapSingle { keyword -> api.search(keyword) }
这里 debounce 的时间线是 a(0ms) -> an(100ms) -> and(250ms) -> 550ms 发射 and; throttleLatest(300ms, unit, emitLast) 按窗口尽快发射一个值并保留窗口内最新值, emitLast = true 会在上游完成时补发最后暂存值, false 则不补发; 它适合允许周期性刷新. sample(300ms) 在采样时刻取最近值, 可能一个值也不发射; throttleFirst(300ms) 保留每个窗口第一值, 适合防连点. 300ms 只是示例: 用输入间隔, 请求数和可感知等待时间测量后调整. 搜索用错 flatMap 的症状是旧请求晚返回覆盖新结果; 证据是请求 ID 与返回顺序相反; 定位为没有取消旧 inner source; 修复为 switchMap 并确认 source 支持 dispose; 验证为故意让旧请求更慢, 检查 UI 只显示最新关键词结果.
背压策略的可操作选择
Observable 没有下游 request 协议; 从它转为 Flowable 时必须声明过载语义. Flowable.create(source, BackpressureStrategy.BUFFER) 本身使用无界缓冲, 不能称为有限 BUFFER. 需要容量上限时, 选择能提供有界队列的 source, 或在合适边界调用 onBackpressureBuffer(capacity, onOverflow, strategy); 也要区分 observeOn 等异步 operator 可能拥有自己的有界预取缓冲, 这不等于 source 已具备端到端背压策略. 示例: 传感器 UI 仪表可以丢中间读数, 审计流水不能丢也必须保序.
// 上下文片段:仅当业务允许丢弃中间仪表值.
val latestReading = sensorObservable
.toFlowable(BackpressureStrategy.LATEST)
.observeOn(AndroidSchedulers.mainThread())
.subscribe(::renderReading, ::showError)
先记录生产频率, 下游处理耗时, 队列深度, 允许的最大陈旧时间和丢失比例; 若所有记录必须处理, 应采取上游限速/批处理/持久化队列或 ERROR 后明确降级, 而不是无界 BUFFER. 典型失败是内存增长和 UI 显示数秒前状态; 证据是队列长度, 堆快照和事件时间戳; 定位为生产速率长期大于消费速率; 修复按上述业务语义选择限速, 批处理或策略; 验证在压测 / 可控 fake 高频源下观察队列上限和延迟, 不虚构结果.
五, 生命周期与资源释放
页面或 owner 应集中管理订阅:
private val disposables = CompositeDisposable()
fun load() {
disposables.add(repository.observeUser()
.observeOn(AndroidSchedulers.mainThread())
.subscribe(::render, ::showError))
}
fun clear() {
disposables.clear() // owner 仍可能复用; dispose() 后容器不可再接收新订阅
}
注意:
clear()取消当前集合但允许继续 add;dispose()会永久标记容器已释放.- 长生命周期单例订阅短生命周期 View/Activity 容易泄漏.
- dispose 是取消请求继续向下游交付, 不等于底层工作一定瞬时停止; 取决于 source 是否支持取消.
- 生命周期库只能帮助触发 dispose, 不能修复错误的 owner 和共享链设计.
六, 错误, 重试与终止
Rx 流一旦向下游发送 terminal error 就结束. onErrorReturn/onErrorResumeNext 会把错误转换为值或替代流, 需要避免吞掉不可恢复错误.
retryWhen 必须有终止条件:
- 仅重试可恢复错误, 不重试鉴权失败, 参数错误等确定性失败.
- 指数退避并加入 jitter, 避免客户端同时重试.
- 设置最大次数或总时限.
- owner 销毁, 网络条件不满足或用户取消时停止.
- 重试写操作前确认服务端幂等键和业务语义.
七, 向协程与 Flow 渐进迁移
不要一次性重写整个存量项目. 推荐从边界开始:
- 固定 RxJava 主版本并补行为测试.
- 新 Retrofit API 优先使用
suspend; DAO/SDK 边界按真实支持选择适配方式. - 在 repository 边界集中做 Rx/Flow/
suspend转换, 避免 UI 层同时理解两套模型. - 明确取消, 热/冷流, 背压/缓冲, 错误终止和线程语义的差异.
- 逐模块迁移并比较稳定性, 复杂度和性能, 不以 “新技术” 本身作为改造收益.
Flow 的挂起可形成自然背压, 但 buffer, conflate, SharedFlow 等仍可能丢失或积压; 不能把 Flow 简化为 “自动解决所有背压”.
高频面试题
Q1: flatMap, concatMap, switchMap 怎么选?
看是否允许并发, 是否要求顺序, 旧任务是否应取消. 并发合并用 flatMap, 严格顺序用 concatMap, 只关心最新结果用 switchMap.
Q2: CompositeDisposable.clear() 和 dispose() 的区别?
二者都释放当前订阅; clear() 后可继续 add, dispose() 后容器保持 disposed, 新加入项会立即被释放.
Q3: 重试为什么危险?
无限重试会放大故障和功耗, 写操作还可能重复扣款或写入. 必须分类错误, 限次, 退避, 加 jitter, 并依赖服务端幂等.
练习与掌握检查
- 为搜索, 下载进度, 审计流水分别选择
Subject/Rx 类型和时间操作符, 写出是否允许丢失, 重放和乱序的理由. - 用一个可控 fake source 让上游快于下游, 记录队列或时间戳, 比较
LATEST与有限BUFFER的业务后果. - 在页面销毁后调用
CompositeDisposable.clear(), 验证订阅不再更新 View; 若仍更新, 按 “owner→订阅位置→dispose 支持→验证” 排查. - 能解释上述搜索时间线, 旧请求覆盖问题和背压选择, 即达到本篇掌握标准.
版本与参考资料
- 最后核验: 2026-08-07.
- 适用基线: 先确认项目使用 RxJava 2 还是 RxJava 3, 以及 RxAndroid/Retrofit/Room 等 adapter 的对应版本.
- 官方资料: RxJava 项目, RxAndroid 项目.