RxJava(ReactiveX Java)是 响应式编程(Reactive Programming)框架,用于处理 异步数据流。它提供了丰富的 API,可以方便地进行事件驱动编程、异步任务管理、线程调度等。
核心概念
- Observable: 数据源,发出数据流。
- Observer: 订阅 Observable,处理数据、错误和完成事件。
- Subscriber: Observer 的实现类,通常用于订阅 Observable。
- Operators: 用于转换、过滤、组合数据流的操作符,如
map、filter、flatMap等。 - 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)一个新的字符串。
- 这个字符串会流向
Observer或Subscriber进行订阅处理。
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?
如果你输入:
"B"(API 请求 1)"Be"(API 请求 2)"Bei"(API 请求 3)"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: RxJavajust()只能发送固定数据(不能动态创建)。- 同步执行,立即发送所有数据。
示例 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 组件(如ViewModel、LiveData)确实有一些重叠的功能,但它们的设计目标和适用场景不同。RxJava 并没有被淘汰,但在 Android 开发中,Jetpack 组件(尤其是ViewModel和LiveData)已经成为更推荐的方式,特别是在处理 UI 相关的数据流时。
RxJava
优点
- 强大的异步操作:RxJava 提供了丰富的操作符(如
map、flatMap、switchMap等),可以轻松处理复杂的异步数据流。 - 线程切换:通过
subscribeOn和observeOn,可以灵活地控制任务执行的线程。 - 事件流处理:适合处理连续的事件流(如搜索框输入、实时数据更新等)。
缺点
- 学习曲线陡峭:RxJava 的概念(如 Observable、Observer、操作符等)对初学者来说较难掌握。
- 容易引发内存泄漏:如果订阅没有及时取消,可能会导致内存泄漏。
- 代码复杂度高:对于简单的 UI 数据绑定,RxJava 可能会显得过于复杂。
适用场景
- 复杂的异步数据流处理(如多个网络请求的组合、实时数据更新等)。
- 需要灵活线程调度的场景。
- 事件驱动型应用(如聊天应用、实时数据监控等)。
Jetpack
Jetpack 是 Google 推出的一套 Android 开发工具库,旨在简化开发流程并提高代码质量。其中与 RxJava 功能重叠的主要是 ViewModel 和 LiveData。
**ViewModel **
- 作用:用于管理 UI 相关的数据,并在配置更改(如屏幕旋转)时保持数据的一致性。
- 生命周期感知:
ViewModel的生命周期与 Activity/Fragment 分离,避免内存泄漏。 - 数据持久化:在 Activity/Fragment 重建时,
ViewModel中的数据不会丢失。
LiveData
- 作用:用于观察数据的变化,并在数据更新时自动通知 UI。
- 生命周期感知:
LiveData只会通知处于活跃生命周期状态的观察者,避免内存泄漏。 - UI 数据绑定:与
ViewModel结合使用,可以轻松实现数据驱动 UI。
优点
- 简单易用:
ViewModel和LiveData的概念简单,学习成本低。 - 生命周期感知:自动管理生命周期,避免内存泄漏。
- 与 Android 架构深度集成:Jetpack 组件是 Android 官方推荐的开发方式,与 Android Studio 和其他工具集成良好。
缺点
- 功能相对简单:相比 RxJava,
LiveData的功能较为简单,不适合处理复杂的异步数据流。 - 缺乏操作符:
LiveData没有 RxJava 那样丰富的操作符,处理复杂逻辑时需要额外代码。
适用场景
- UI 数据绑定(如将网络请求结果绑定到 UI)。
- 简单的异步任务(如数据库查询、网络请求等)。
- 需要生命周期感知的场景。
RxJava 和 Jetpack 组件的区别
| 特性 | RxJava | Jetpack 组件(ViewModel + LiveData) |
|---|---|---|
| 异步处理 | 强大,支持复杂异步数据流 | 简单,适合 UI 数据绑定 |
| 线程调度 | 灵活,支持多线程切换 | 需要手动切换线程(如配合 Coroutine) |
| 生命周期感知 | 需要手动管理(如 CompositeDisposable) | 自动管理,避免内存泄漏 |
| 操作符 | 丰富(如 map、flatMap、switchMap 等) | 无操作符,功能简单 |
| 学习曲线 | 陡峭,概念复杂 | 简单,易于上手 |
| 适用场景 | 复杂异步数据流、事件驱动型应用 | UI 数据绑定、简单异步任务 |
RxJava 是否已经是过时的产物?
RxJava 并没有被淘汰,但在 Android 开发中,Jetpack 组件(尤其是 ViewModel 和 LiveData)已经成为更推荐的方式,特别是在处理 UI 相关的数据流时。以下是一些原因:
- Jetpack 组件的普及:
- Jetpack 是 Google 官方推荐的开发方式,与 Android 生态深度集成。
ViewModel和LiveData的设计更符合 Android 的生命周期模型。
- Kotlin 协程的崛起:
- Kotlin 协程提供了更简单、更直观的异步编程方式,逐渐取代了 RxJava 的部分功能。
- 协程与
LiveData和ViewModel结合使用,可以轻松实现复杂的异步任务。
- 开发效率:
- 对于大多数应用场景,Jetpack 组件和协程已经足够使用,且代码更简洁、更易维护。
- RxJava 的复杂性和学习曲线使得它在简单场景中显得“杀鸡用牛刀”。
对比 Jetpack+协程 RxJava 有什么优势?
| 特性 | RxJava | Jetpack + 协程 (LiveData/Flow) |
|---|---|---|
| 异步流处理 | ✅ 强大,支持事件流(Observable、Flowable) | ✅ Flow 也支持流式数据处理 |
| 背压处理 | ✅ Flowable 专门用于背压处理 | ⚠️ Flow 默认不支持背压,但可以用 buffer() |
| 组合多个异步数据源 | ✅ 强大的 merge(), combineLatest(), zip() | ⚠️ Flow 也支持 combine(),但功能略弱 |
| 多线程调度 | ✅ subscribeOn() 和 observeOn() 控制流切换 | ✅ withContext() 控制协程切换 |
| 生命周期感知 | ❌ 需要手动管理 Disposable,容易内存泄漏 | ✅ LiveData 和 Flow 自动感知生命周期 |
| 延迟计算 | ✅ defer() 可以动态创建 Observable | ✅ suspend 方法 + 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,当 A1 或 B1 变化时,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) }🚀 这段代码的 核心思想:
- 监听搜索框输入(流)
- 防止高频 API 请求(防抖)
- 避免重复查询(去重)
- 自动切换最新请求(
switchMap) - 订阅并更新 UI
这就是 RxJava + 响应式编程 的典型用法,它让 异步数据流管理变得优雅高效!
更新: 2025-03-24 10:43:25
原文: https://www.yuque.com/dongpozhouzi-mshe3/zhm85g/hn0tcyl47wwb446k