RxJava(ReactiveX Java)是 响应式编程(Reactive Programming)框架,用于处理 异步数据流。它提供了丰富的 API,可以方便地进行事件驱动编程、异步任务管理、线程调度等。

核心概念

  1. Observable: 数据源,发出数据流。
  2. Observer: 订阅 Observable,处理数据、错误和完成事件。
  3. Subscriber: Observer 的实现类,通常用于订阅 Observable。
  4. Operators: 用于转换、过滤、组合数据流的操作符,如 mapfilterflatMap 等。
  5. Schedulers: 控制线程调度,如 Schedulers.io() 用于 I/O 操作,Schedulers.computation() 用于计算任务。

主要特点

  • 异步处理: 简化异步编程,避免回调地狱。
  • 链式调用: 通过操作符链式处理数据流。
  • 线程调度: 灵活控制任务执行线程。
  • 丰富的操作符: 提供多种操作符处理数据流。

Demo

package com.markz.rxjava.real
 
import android.os.Bundle
import android.text.Editable
import android.text.TextWatcher
import android.util.Log
import androidx.activity.enableEdgeToEdge
import androidx.appcompat.app.AppCompatActivity
import androidx.core.view.ViewCompat
import androidx.core.view.WindowInsetsCompat
import androidx.core.widget.addTextChangedListener
import com.markz.rxjava.R
import com.markz.rxjava.databinding.ActivityWeatherBinding
import io.reactivex.rxjava3.android.schedulers.AndroidSchedulers
import io.reactivex.rxjava3.core.Observable
import io.reactivex.rxjava3.core.Observer
import io.reactivex.rxjava3.disposables.CompositeDisposable
import io.reactivex.rxjava3.disposables.Disposable
import io.reactivex.rxjava3.schedulers.Schedulers
import java.util.concurrent.TimeUnit
 
class WeatherActivity : AppCompatActivity() {
 
    private lateinit var binding: ActivityWeatherBinding
 
    private lateinit var weatherApiService: WeatherApiService
 
    private val disposable = CompositeDisposable() // 统一管理 RxJava 订阅,避免内存泄漏
 
    override fun onCreate(savedInstanceState: Bundle?) {
        super.onCreate(savedInstanceState)
        enableEdgeToEdge()
 
        binding = ActivityWeatherBinding.inflate(layoutInflater)
        setContentView(binding.root)
 
        ViewCompat.setOnApplyWindowInsetsListener(findViewById(R.id.main)) { v, insets ->
            val systemBars = insets.getInsets(WindowInsetsCompat.Type.systemBars())
            v.setPadding(systemBars.left, systemBars.top, systemBars.right, systemBars.bottom)
            insets
        }
 
        val searchInput = binding.searchPlaceEdit
        val resultText = binding.resultText
 
        weatherApiService = ServiceCreator.create()
 
 
        // 创建 RxJava 观察流
        val searchObservable = Observable.create<String> { emitter ->
            searchInput.addTextChangedListener(object : TextWatcher {
                override fun beforeTextChanged(
                    s: CharSequence?, start: Int, count: Int, after: Int
                ) {
                    // do nothing
                }
 
                override fun onTextChanged(s: CharSequence?, start: Int, before: Int, count: Int) {
                    // do noting
                }
 
                override fun afterTextChanged(s: Editable?) {
                    emitter.onNext(s.toString()) // 发射输入的文本
                }
            })
        }
 
        // 订阅 RxJava 观察流
        searchObservable.debounce(500, TimeUnit.MILLISECONDS) // 防抖,减少 API 调用
            .filter { it.isNotEmpty() } // 过滤空输入
            .distinctUntilChanged() // 只有输入内容发生变化才触发请求
            .switchMap { city ->
                weatherApiService.searchPlaces(city)
                ?.subscribeOn(Schedulers.io()) // 网络请求在 IO 线程
                ?.observeOn(AndroidSchedulers.mainThread())!! // UI 更新在主线程
            }
            .observeOn(AndroidSchedulers.mainThread()) // UI 更新在主线程
            .subscribe(object : Observer<PlaceResponse> {
                override fun onSubscribe(d: Disposable) {
                    disposable.add(d)
                }
 
                override fun onError(e: Throwable) {
                    Log.d("Network", "错误信息:${e.message}")
                    resultText.text = "❌ 请求失败,请稍候重试!"
                }
 
                override fun onComplete() {
                }
 
                override fun onNext(res: PlaceResponse) {
                    Log.d("ThreadCheck", "Current thread: ${Thread.currentThread().name}")
                    if (res.status == "ok") {
                        val placeList: List<Place> = res.placeList
                        val info = "🌍 ${placeList[0].toString()}"
                        resultText.text = info
                    }
                }
            })
    }
 
    override fun onDestroy() {
        super.onDestroy()
        disposable.clear() // Activity 销毁时取消订阅,避免内存泄漏
    }
}

方法解析

Observable.create** 创建可观察流**

val searchObservable = Observable.create<String> { emitter ->
    searchInput.addTextChangedListener(object : TextWatcher {
        override fun afterTextChanged(s: Editable?) {
            emitter.onNext(s.toString())  // 每次输入变化都会发射新值
        }
        override fun beforeTextChanged(s: CharSequence?, start: Int, count: Int, after: Int) {}
        override fun onTextChanged(s: CharSequence?, start: Int, before: Int, count: Int) {}
    })
}
  • Observable.create { emitter -> } 这个方法创建了一个 可观察的数据流
  • emitter.onNext(s.toString())
    • 每次用户输入变化,都会发射(emit)一个新的字符串
    • 这个字符串会流向 ObserverSubscriber 进行订阅处理

debounce** 作用 —— 防抖**

.debounce(500, TimeUnit.MILLISECONDS) // 防抖,减少 API 调用

为什么要防抖?

如果你连续输入B -> Be -> Bei -> Beij -> Beiji -> Beijin -> Beijing

每次输入都会触发请求,会造成API 调用太频繁,影响性能。

所以 debounce(500ms) 让它等 500 毫秒,如果用户继续输入,计时器重置,直到用户停下输入 500ms 后才执行请求。

distinctUntilChanged** 作用 —— 去重**

.distinctUntilChanged() // 只有输入内容发生变化才触发请求

示例

输入触发请求?
“Bei”✅ 发起请求
”Bei”❌ 跳过(和上次一样)
“Beij”✅ 发起请求
”Beij”❌ 跳过

这个操作避免了重复的 API 调用,如果用户输入了一次 “Bei”,然后停留 10 秒,但不修改内容,RxJava 不会再触发请求

4. switchMap 作用

.switchMap { city ->
    RetrofitClient.apiService.getWeather(city)
        .subscribeOn(Schedulers.io())  // 让网络请求在 IO 线程执行
        .observeOn(AndroidSchedulers.mainThread())  // UI 更新在主线程
}

问题:为什么要用 switchMap

如果你输入:

  1. "B"(API 请求 1)
  2. "Be"(API 请求 2)
  3. "Bei"(API 请求 3)
  4. "Beij"(API 请求 4)
    网络请求 1,2,3 可能会并发,最后 UI 显示的是 请求 1 先完成的数据,而不是最新的 “Beij”

switchMap** 解决了这个问题**:

  • 新请求到来时,取消前一个请求,只保留最后一个。
  • 避免早期的无效数据覆盖最新数据

总结

可以这么理解:

  • Observable 负责发射数据(就像水龙头流出水)。
  • Observer 负责接收数据(就像接水的杯子)。
  • subscribe() 就是连接水龙头和杯子
  • switchMap()debounce()distinctUntilChanged() 这些操作符就是控制水流的阀门,让水流得更高效。

Question

Observable.just()Observable.create() 的区别?

这两个方法都能创建 Observable,但有以下关键区别:

方法适用场景特点
Observable.just(T...)已知固定数据- 直接发送一个或多个数据项(同步) - 不能执行异步操作
Observable.create(ObservableOnSubscribe<T>)自定义逻辑(异步/手动触发)- 需要手动调用onNext() 发送数据 - 可以执行异步操作(如数据库、网络请求)

示例 1:Observable.just()(适用于简单数据流)

val observable = Observable.just("Hello", "RxJava")
 
observable.subscribe { item -> 
    println("Received: $item") 
}

输出:

Received: Hello
Received: RxJava
  • just() 只能发送固定数据(不能动态创建)。
  • 同步执行,立即发送所有数据。

示例 2:Observable.create()(适用于异步或复杂逻辑)

val observable = Observable.create<String> { emitter ->
    emitter.onNext("Loading...")
    Thread.sleep(1000)  // 模拟网络请求
    emitter.onNext("Data loaded!")
    emitter.onComplete()
}
 
observable.subscribe(
    { item -> println("Received: $item") },
    { error -> println("Error: ${error.message}") },
    { println("Complete!") }
)

输出:

Received: Loading...
Received: Data loaded!
    Complete!

🔹 特点

  • 需要手动调用 onNext() 发送数据。
  • 可以执行异步操作(如 Thread.sleep() 或网络请求)。
  • 可以控制 onComplete(),表示数据流结束。

RxJava 和 Jetpack 基础组件+协程的区别?

RxJava 和 Jetpack 组件(如ViewModelLiveData)确实有一些重叠的功能,但它们的设计目标和适用场景不同。RxJava 并没有被淘汰,但在 Android 开发中,Jetpack 组件(尤其是ViewModelLiveData)已经成为更推荐的方式,特别是在处理 UI 相关的数据流时。

RxJava

优点

  • 强大的异步操作:RxJava 提供了丰富的操作符(如 mapflatMapswitchMap 等),可以轻松处理复杂的异步数据流。
  • 线程切换:通过 subscribeOnobserveOn,可以灵活地控制任务执行的线程。
  • 事件流处理:适合处理连续的事件流(如搜索框输入、实时数据更新等)。

缺点

  • 学习曲线陡峭:RxJava 的概念(如 Observable、Observer、操作符等)对初学者来说较难掌握。
  • 容易引发内存泄漏:如果订阅没有及时取消,可能会导致内存泄漏。
  • 代码复杂度高:对于简单的 UI 数据绑定,RxJava 可能会显得过于复杂。

适用场景

  • 复杂的异步数据流处理(如多个网络请求的组合、实时数据更新等)。
  • 需要灵活线程调度的场景。
  • 事件驱动型应用(如聊天应用、实时数据监控等)。

Jetpack

Jetpack 是 Google 推出的一套 Android 开发工具库,旨在简化开发流程并提高代码质量。其中与 RxJava 功能重叠的主要是 ViewModelLiveData

**ViewModel **

  • 作用:用于管理 UI 相关的数据,并在配置更改(如屏幕旋转)时保持数据的一致性。
  • 生命周期感知ViewModel 的生命周期与 Activity/Fragment 分离,避免内存泄漏。
  • 数据持久化:在 Activity/Fragment 重建时,ViewModel 中的数据不会丢失。

LiveData

  • 作用:用于观察数据的变化,并在数据更新时自动通知 UI。
  • 生命周期感知LiveData 只会通知处于活跃生命周期状态的观察者,避免内存泄漏。
  • UI 数据绑定:与 ViewModel 结合使用,可以轻松实现数据驱动 UI。

优点

  • 简单易用ViewModelLiveData 的概念简单,学习成本低。
  • 生命周期感知:自动管理生命周期,避免内存泄漏。
  • 与 Android 架构深度集成:Jetpack 组件是 Android 官方推荐的开发方式,与 Android Studio 和其他工具集成良好。

缺点

  • 功能相对简单:相比 RxJava,LiveData 的功能较为简单,不适合处理复杂的异步数据流。
  • 缺乏操作符LiveData 没有 RxJava 那样丰富的操作符,处理复杂逻辑时需要额外代码。

适用场景

  • UI 数据绑定(如将网络请求结果绑定到 UI)。
  • 简单的异步任务(如数据库查询、网络请求等)。
  • 需要生命周期感知的场景。

RxJava 和 Jetpack 组件的区别

特性RxJavaJetpack 组件(ViewModel + LiveData)
异步处理强大,支持复杂异步数据流简单,适合 UI 数据绑定
线程调度灵活,支持多线程切换需要手动切换线程(如配合 Coroutine
生命周期感知需要手动管理(如 CompositeDisposable自动管理,避免内存泄漏
操作符丰富(如 mapflatMapswitchMap 等)无操作符,功能简单
学习曲线陡峭,概念复杂简单,易于上手
适用场景复杂异步数据流、事件驱动型应用UI 数据绑定、简单异步任务

RxJava 是否已经是过时的产物?

RxJava 并没有被淘汰,但在 Android 开发中,Jetpack 组件(尤其是 ViewModelLiveData)已经成为更推荐的方式,特别是在处理 UI 相关的数据流时。以下是一些原因:

  1. Jetpack 组件的普及
    • Jetpack 是 Google 官方推荐的开发方式,与 Android 生态深度集成。
    • ViewModelLiveData 的设计更符合 Android 的生命周期模型。
  2. Kotlin 协程的崛起
    • Kotlin 协程提供了更简单、更直观的异步编程方式,逐渐取代了 RxJava 的部分功能。
    • 协程与 LiveDataViewModel 结合使用,可以轻松实现复杂的异步任务。
  3. 开发效率
    • 对于大多数应用场景,Jetpack 组件和协程已经足够使用,且代码更简洁、更易维护。
    • RxJava 的复杂性和学习曲线使得它在简单场景中显得“杀鸡用牛刀”。

对比 Jetpack+协程 RxJava 有什么优势?

特性RxJavaJetpack + 协程 (LiveData/Flow)
异步流处理✅ 强大,支持事件流(ObservableFlowableFlow 也支持流式数据处理
背压处理Flowable 专门用于背压处理⚠️ Flow 默认不支持背压,但可以用 buffer()
组合多个异步数据源✅ 强大的 merge(), combineLatest(), zip()⚠️ Flow 也支持 combine(),但功能略弱
多线程调度subscribeOn()observeOn() 控制流切换withContext() 控制协程切换
生命周期感知❌ 需要手动管理 Disposable,容易内存泄漏LiveDataFlow 自动感知生命周期
延迟计算defer() 可以动态创建 Observablesuspend 方法 + lazy
错误处理onErrorResumeNext()retry()try-catch / catch()
响应式 UI 更新✅ 适合事件流,比如点击流、搜索框输入流⚠️ LiveData 更适合 UI 层,但 Flow 也能实现
学习曲线⚠️ 难度较高,API 复杂✅ Kotlin 原生支持,代码更简洁
官方支持✅ 社区维护良好✅ Google 官方推荐
RxJava 相比 Jetpack + 协程的优势

1. 更强大的数据流合并 & 变换

  • RxJava 提供 merge(), zip(), combineLatest(), flatMap(), switchMap() 等操作符,可以轻松处理多个数据流的合并和变换。
  • Jetpack 的 Flow 也提供 combine(),但功能比 RxJava 弱。

🔹 RxJava 示例(多个 API 合并)

Observable.zip(
    api.getUser(),
    api.getUserPosts(),
    BiFunction<User, List<Post>, UserProfile> { user, posts ->
        UserProfile(user, posts)
    }
).subscribe { profile -> println(profile) }

🔹 Flow 也可以做到,但写法更冗长

suspend fun getUserProfile(): UserProfile {
    val user = async { api.getUser() }
    val posts = async { api.getUserPosts() }
    return UserProfile(user.await(), posts.await())
}

2. 背压处理更成熟

  • Flowable 通过 onBackpressureDrop()onBackpressureBuffer() 处理 高频数据流(如传感器、WebSocket)。
  • Flow 默认不支持背压,但可以用 buffer()

🔹 RxJava 背压示例

Flowable.interval(10, TimeUnit.MILLISECONDS)
    .onBackpressureDrop() // 丢弃多余的数据
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe { println(it) }

🔹 Flow 需要 buffer() 来避免数据丢失

flow {
    for (i in 1..1000) {
        emit(i)
    }
}.buffer() // 避免丢数据
    .collect { println(it) }

RxJava 在高频数据流(如股票行情、按键监听)方面比 Flow 更稳定。
3. 事件流(UI 交互、搜索框防抖)
RxJava 适用于 UI 交互事件流,如:

  • 搜索框输入防抖
  • 组合点击事件(双击、长按等)
  • 动态表单验证

🔹 RxJava 防抖搜索

val searchObservable = RxTextView.textChanges(searchBox)
    .debounce(300, TimeUnit.MILLISECONDS) // 300ms 防抖
    .distinctUntilChanged() // 过滤重复输入
    .switchMap { query -> api.search(query).toObservable() }
    .subscribe { result -> showResult(result) }

🔹 Flow 实现(稍微复杂)

searchBox.textChanges()
    .debounce(300.milliseconds)
    .distinctUntilChanged()
    .flatMapLatest { query -> flow { emit(api.search(query)) } }
    .collect { result -> showResult(result) }

虽然 Flow 也能实现,但 RxJava 现成的操作符更丰富

Jetpack + 协程的优势

虽然 RxJava 在流式处理和事件合并上更强,但 Jetpack + 协程有以下优势
1. 代码更简洁

  • RxJava 代码复杂,而 suspend 方法更直观:
suspend fun fetchData() = api.getData()

相比 RxJava:

Observable.fromCallable { api.getData() }
    .subscribeOn(Schedulers.io())
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe { data -> show(data) }

协程**不需要 subscribeOn() 和 **observeOn(),更简洁。
2. 自动感知生命周期

  • LiveData & Flow自动在 ViewModel 销毁时停止任务,避免内存泄漏。
  • RxJava 需要 **手动管理 **Disposable,否则会有泄漏风险:
val disposable = api.getData()
    .subscribe { data -> show(data) }
disposable.dispose() // 需要手动清理

ViewModelScope.launch会自动在 ViewModel 销毁时取消任务

viewModelScope.launch {
    val data = api.getData()
    show(data)
} // 不需要手动取消!

3. Google 官方推荐
未来 Android 发展方向是 Kotlin + Flow,RxJava 主要由社区维护,不再是官方推荐方案。

总结
适用场景推荐方案
简单异步请求(API 调用)✅ 协程 + suspend
UI 层数据管理(生命周期感知)✅ ViewModel + LiveData
流式数据处理(事件、合并流)✅ RxJava
高频数据流(传感器、WebSocket)✅ RxJava (Flowable)
复杂数据变换✅ RxJava (combineLatest() 更强)

💡 如果你的项目是新项目,优先用 Jetpack + 协程

💡 如果项目已有 RxJava 代码,或者需要高级流式处理,RxJava 仍然是优秀选择

推荐策略

  • 90% 情况下:Jetpack + 协程(Flow / LiveData)
  • 特殊情况(高频数据、复杂事件流):RxJava 仍然是强大工具

所以什么是响应式编程?RxJava 和响应式编程的关系?

响应式编程(Reactive Programming,简称 RP) 是一种 异步、数据驱动 的编程范式,核心思想是:

  • 一切皆流(Everything is a stream):数据、用户事件、API 响应等都可以看作是 数据流(Stream)。
  • 声明式编程(Declarative Programming):你描述 数据如何流动,而不是手动管理线程或回调。
  • 自动传播变化(Reactive):当数据变化时,订阅者(Subscriber) 会自动收到更新。

🔹 举个例子

假设你在 Excel 里输入 C1 = A1 + B1,当 A1B1 变化时,C1 会自动更新。这就是 响应式编程 的思想。

RxJava(Reactive Extensions for Java) 是响应式编程在 Java / Android 上的实现之一。它提供了一套 强大的 API 来处理:

  • 异步任务(API 调用)
  • 事件流(用户输入、点击、传感器数据)
  • 数据变换map()flatMap()
  • 多数据源合并merge()zip()

RxJava 让开发者 不用写复杂的回调,就能优雅地处理 异步数据流

RxJava 例子:响应式搜索:

val searchObservable = RxTextView.textChanges(searchBox)
    .debounce(300, TimeUnit.MILLISECONDS)  // 防抖
    .distinctUntilChanged()  // 过滤重复输入
    .switchMap { query -> api.search(query).toObservable() }  // 取消旧请求,发起新请求
    .subscribe { result -> showResult(result) }

🚀 这段代码的 核心思想

  1. 监听搜索框输入(流)
  2. 防止高频 API 请求(防抖)
  3. 避免重复查询(去重)
  4. 自动切换最新请求switchMap
  5. 订阅并更新 UI
    这就是 RxJava + 响应式编程 的典型用法,它让 异步数据流管理变得优雅高效

更新: 2025-03-24 10:43:25
原文: https://www.yuque.com/dongpozhouzi-mshe3/zhm85g/hn0tcyl47wwb446k


相关笔记