HelloWorld RxJava 集成教程

这篇教程用最直观的方式示范如何在一个简单的 HelloWorld Android 项目中集成 RxJava:从依赖配置、线程调度、可观察数据流、常用操作符,到资源回收和错误处理,逐步示例代码并解释每一步的原理与注意点,帮助你在实际项目中安全、清晰地采用响应式编程。以及实践中的坑位与优化建议。适合新手上手。

HelloWorld RxJava 集成教程

HelloWorld RxJava 集成教程

先说结论(快速上手思路)

把 RxJava 当成「把异步流写成同步样子」的工具。先在 Gradle 加依赖,理解 Observable/FlowableSubscriber/ObserverSchedulerDisposable,然后用简单例子把回调改成链式操作,最后用 CompositeDisposable 管理生命周期、用 Flowable 处理背压。下面一步步来,边写边解释。

为什么要用 RxJava?

先简单说优点:

  • 把异步逻辑抽象为数据流,代码更可组合、可测试。
  • 链式操作符(map、flatMap 等)让转换与组合变得表达性强。
  • 统一线程切换(subscribeOn/observeOn),简化 UI 与后台线程交互。
  • 配合 Retrofit、数据库、事件总线等能显著减少回调地狱。

缺点也要知道:学习曲线、调试困难、滥用会导致不可维护的链。权衡后在复杂异步场景引入是值得的。

核心概念一遍讲清楚(费曼式解释)

Observable / Flowable / Single / Maybe / Completable

把它们想像成不同口径的水管:

  • Observable:常规流,适合无限或中等量数据。
  • Flowable:带背压的流,适合高频生产者,防止消费者被淹没。
  • Single:只会发一次成功值或错误(比如网络请求返回一个对象)。
  • Maybe:可能有值也可能无值或错误。
  • Completable:只关心完成或错误(比如写数据库,只关心成功/失败)。

Observer / Subscriber / Disposable

订阅者就是消费者,订阅后会拿到数据或错误。Disposable 用来取消订阅,避免内存泄漏。把多个 Disposable 放进 CompositeDisposable 一并清理。

Scheduler(线程调度)

两个常用概念:subscribeOn 决定数据生产在哪个线程,observeOn 决定观察者在哪个线程收到数据。常用:Schedulers.io()、Schedulers.computation()、AndroidSchedulers.mainThread()(需要 RxAndroid)。

在 HelloWorld Android 项目中一步步集成

1. 新建项目与添加依赖

这是最简单也是最容易出错的地方:版本不匹配。下面给出常用组合(示例版本,随时间升级):

示例版本(按需替换)
io.reactivex.rxjava3:rxjava 3.1.5
io.reactivex.rxjava3:rxandroid 3.0.0
com.squareup.retrofit2:adapter-rxjava3 2.9.0(若用 Retrofit)

在 app/build.gradle(Kotlin 或 Groovy)里加入:

implementation 'io.reactivex.rxjava3:rxjava:3.1.5'
implementation 'io.reactivex.rxjava3:rxandroid:3.0.0'

2. 第一个 HelloWorld 示例(最小可运行)

思路:Observable 发出 “Hello RxJava”,订阅者在主线程更新 TextView。

// 假想 Activity 代码片段(概念示意)
private val disposables = CompositeDisposable()

override fun onCreate(...) { // ... val obs = io.reactivex.rxjava3.core.Observable.just("Hello RxJava") .subscribeOn(io.reactivex.rxjava3.schedulers.Schedulers.io()) .observeOn(io.reactivex.rxjava3.android.schedulers.AndroidSchedulers.mainThread()) .map { it + " — from HelloWorld" }

val d = obs.subscribe({ text -> textView.text = text }, { err -> textView.text = "Error: ${err.message}" }) disposables.add(d) }

override fun onDestroy() { super.onDestroy() disposables.clear() }

要点:把订阅结果放到 CompositeDisposable,Activity 销毁时 clear,避免泄漏。

常用操作符与直观示例(把用途讲清楚)

  • map:一对一转换,像给每个元素做一次变换。
  • flatMap:一对多或异步映射,常用于把一个请求映射成另一个网络请求(结果乱序)。
  • concatMap:类似 flatMap,但保持顺序,按序发出。
  • switchMap:如果新的内流来了,丢弃旧的内流,适合搜索联想场景。
  • filter:过滤不需要的数据。
  • zip:合并多个流,按位置配对输出。

示例:两个网络请求串联(伪代码):

api.getUser(userId)                    // Single
  .flatMap { user -> api.getPosts(user.id) } // Single>
  .subscribeOn(Schedulers.io())
  .observeOn(AndroidSchedulers.mainThread())
  .subscribe({ posts -> /* 更新 UI */ }, { err -> /* 处理错误 */ })

错误处理与重试策略

  • onErrorReturn:出错时返回一个默认值继续流。
  • onErrorResumeNext:出错时切换到另一个流。
  • retry / retryWhen:重试策略,retryWhen 可结合延迟与条件。

不要滥用无限重试——要有上限和退避(backoff)策略来防止雪崩。

背压(Backpressure)说明与实战

当生产速度远高于消费速度时,需要背压保护。RxJava 2/3 使用 Flowable 来表达带背压的流。创建高频源时优先用 Flowable.create,并选择合适的 BackpressureStrategy(BUFFER、DROP、LATEST、ERROR)。

// 简单示例
Flowable.create({ emitter ->
  for (i in 1..10000) {
    if (emitter.isCancelled) return@create
    emitter.onNext(i)
  }
  emitter.onComplete()
}, BackpressureStrategy.BUFFER)
  .observeOn(Schedulers.computation())
  .subscribe({ v -> process(v) }, { e -> log(e) })

与 Android 生命周期和架构组件协作

实务中常见三种做法:

  • 在 Activity/Fragment 用 CompositeDisposable,并在 onDestroy/onDestroyView 清理。
  • 在 ViewModel 持有 CompositeDisposable,Activity/Fragment 观察 ViewModel,ViewModel 在 onCleared 清理。
  • 使用 AutoDispose 或 Lifecycle-aware 适配器(额外库)实现自动解绑。

与 Retrofit 集成(常见场景)

Retrofit 支持返回 Single/Observable/Flowable,配合 RxJava 能把网络调用变成链式操作。要在 Retrofit 中加入 CallAdapter(依赖示例见上表)。示例接口:

interface ApiService {
  @GET("user/{id}")
  fun getUser(@Path("id") id: String): io.reactivex.rxjava3.core.Single
}

要注意:网络错误会落到 onError,别忘了做友好的错误提示与重试。

调试与性能优化小贴士(实用)

  • 用 doOnNext / doOnError 插入日志,调试链条数据。
  • 避免在操作符里进行耗时阻塞,耗时任务放到 Schedulers.io() 或 computation。
  • 限制并发:flatMap 有并发参数,避免同时发太多请求。
  • 尽量使用具体类型(Single/Completable)表达语义,代码可读性更好。

常见坑位(真的会踩的)

  • 忘记 clear/ dispose,导致 Activity 泄漏。
  • subscribeOn/observeOn 顺序理解错,引发 UI 在后台线程更新的崩溃。
  • 把 Observable 用在高频场景却不考虑背压,导致 OOM 或丢帧。
  • 在链条里捕获错误却未重新抛出,导致上层永远不知道异常发生。

把前面内容串成一个小实战流程(一步步做)

  1. 在项目里添加 rxjava + rxandroid 依赖并 sync。
  2. 先写一个最简单的 Observable.just 测试端到端:发、转、订阅,确保主线程能更新 UI。
  3. 把耗时任务放到 Schedulers.io(),用 observeOn(AndroidSchedulers.mainThread()) 更新 UI。
  4. 把 Disposable 添加到 CompositeDisposable,在 onDestroy() 清理。
  5. 遇到高频数据改用 Flowable 并选择合适 BackpressureStrategy。
  6. 当需要多个并发或串联请求时,选择合适的操作符(flatMap、concatMap、zip)。

举几个容易记住的小范例(便于带回工作中复用)

1) 把回调改为 Single:

fun fetchOnce(): Single {
  return Single.create { emitter ->
    networkCall { result, err ->
      if (err != null) emitter.onError(err)
      else emitter.onSuccess(result)
    }
  }
}

2) 搜索联想防抖(switchMap + debounce):

textChanges // Observable from EditText
  .debounce(300, TimeUnit.MILLISECONDS)
  .filter { it.isNotBlank() }
  .switchMapSingle { query -> api.search(query).toSingle() }
  .observeOn(AndroidSchedulers.mainThread())
  .subscribe({ showResults(it) }, { log(it) })

最后说点生活化的建议(这个很管用)

刚开始不要试图把项目全盘改成 Rx,选一个清晰的入口(网络请求或事件流)先做改造,跑通了再慢慢推广。写链的时候刻意把每一步拆成小的、可测试的函数,方便排查问题。遇到奇怪的线程问题先回头确认 subscribeOn/observeOn 的顺序。

如果你跟我一样喜欢边写边想,就会发现 Rx 的很多力量来自将复杂的时间维度统一成流来处理——一开始有点抽象,但养成习惯后你会觉得代码干净很多。试试看上面的小步骤,别忘了版本和生命周期的那两个“坑”。