Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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 渐进迁移

不要一次性重写整个存量项目. 推荐从边界开始:

  1. 固定 RxJava 主版本并补行为测试.
  2. 新 Retrofit API 优先使用 suspend; DAO/SDK 边界按真实支持选择适配方式.
  3. 在 repository 边界集中做 Rx/Flow/suspend 转换, 避免 UI 层同时理解两套模型.
  4. 明确取消, 热/冷流, 背压/缓冲, 错误终止和线程语义的差异.
  5. 逐模块迁移并比较稳定性, 复杂度和性能, 不以 “新技术” 本身作为改造收益.

Flow 的挂起可形成自然背压, 但 buffer, conflate, SharedFlow 等仍可能丢失或积压; 不能把 Flow 简化为 “自动解决所有背压”.

高频面试题

Q1: flatMap, concatMap, switchMap 怎么选?
看是否允许并发, 是否要求顺序, 旧任务是否应取消. 并发合并用 flatMap, 严格顺序用 concatMap, 只关心最新结果用 switchMap.

Q2: CompositeDisposable.clear() 和 dispose() 的区别?
二者都释放当前订阅; clear() 后可继续 add, dispose() 后容器保持 disposed, 新加入项会立即被释放.

Q3: 重试为什么危险?
无限重试会放大故障和功耗, 写操作还可能重复扣款或写入. 必须分类错误, 限次, 退避, 加 jitter, 并依赖服务端幂等.

练习与掌握检查

  1. 为搜索, 下载进度, 审计流水分别选择 Subject/Rx 类型和时间操作符, 写出是否允许丢失, 重放和乱序的理由.
  2. 用一个可控 fake source 让上游快于下游, 记录队列或时间戳, 比较 LATEST 与有限 BUFFER 的业务后果.
  3. 在页面销毁后调用 CompositeDisposable.clear(), 验证订阅不再更新 View; 若仍更新, 按 “owner→订阅位置→dispose 支持→验证” 排查.
  4. 能解释上述搜索时间线, 旧请求覆盖问题和背压选择, 即达到本篇掌握标准.

版本与参考资料

  • 最后核验: 2026-08-07.
  • 适用基线: 先确认项目使用 RxJava 2 还是 RxJava 3, 以及 RxAndroid/Retrofit/Room 等 adapter 的对应版本.
  • 官方资料: RxJava 项目, RxAndroid 项目.