一、什么是 RxSwift?
RxSwift 是 ReactiveX(响应式编程库)的 Swift 实现,它将观察者模式、迭代器模式和函数式编程结合在一起,提供了一种优雅的方式来处理异步事件流。
核心思想
RxSwift 的核心思想是一切皆序列:
Observable(可观察序列)→ Operator(操作符)→ Observer(观察者)
数据源 转换/过滤 接收/处理
为什么使用 RxSwift?
1. 统一异步处理
// ❌ 传统方式:多种异步处理方式
// 通知
NotificationCenter.default.addObserver(self, selector: #selector(handleNotification), name: .userDidLogin, object: nil)
// 定时器
Timer.scheduledTimer(timeInterval: 1.0, target: self, selector: #selector(updateTimer), userInfo: nil, repeats: true)
// 网络请求
URLSession.shared.dataTask(with: url) { data, response, error in
// 处理响应
}.resume()
// ✅ RxSwift:统一处理所有异步事件
// 通知
NotificationCenter.default.rx.notification(.userDidLogin)
.subscribe(onNext: { notification in
// 处理通知
})
// 定时器
Observable<Int>.interval(.seconds(1), scheduler: MainScheduler.instance)
.subscribe(onNext: { time in
// 处理定时器
})
// 网络请求
URLSession.shared.rx.data(request: URLRequest(url: url))
.subscribe(onNext: { data in
// 处理数据
}, onError: { error in
// 处理错误
})
2. 声明式编程
// ❌ 命令式编程:一步步告诉计算机怎么做
var users: [User] = []
var filteredUsers: [User] = []
var sortedUsers: [User] = []
func loadUsers() {
fetchUsers { users in
self.users = users
// 过滤
self.filteredUsers = users.filter { $0.isActive }
// 排序
self.sortedUsers = self.filteredUsers.sorted { $0.name < $1.name }
// 更新 UI
DispatchQueue.main.async {
self.tableView.reloadData()
}
}
}
// ✅ 声明式编程:描述想要什么结果
func loadUsers() {
fetchUsersObservable()
.observeOn(MainScheduler.instance) // 切换到主线程
.filter { $0.isActive } // 过滤活跃用户
.sorted { $0.name < $1.name } // 按名字排序
.subscribe(onNext: { users in
self.users = users
self.tableView.reloadData()
})
}
3. 链式操作
// ✅ RxSwift:优雅的链式操作
searchTextField.rx.text
.orEmpty // 解包 Optional
.debounce(.milliseconds(300), scheduler: MainScheduler.instance) // 防抖
.distinctUntilChanged() // 去重
.map { $0.lowercased() } // 转换为小写
.filter { !$0.isEmpty } // 过滤空字符串
.subscribe(onNext: { searchText in
self.viewModel.search(query: searchText)
})
.disposed(by: disposeBag)
二、核心概念
RxSwift 框架由几个核心概念组成:
1. Observable(可观察序列)
职责:发送数据流给观察者
Observable ──→ 发送 next 事件
├──→ 发送 completed 事件
└──→ 发送 error 事件
2. Observer(观察者)
职责:接收并处理数据流
Observer ──→ 接收 next 事件
├──→ 接收 completed 事件
└──→ 接收 error 事件
3. Operator(操作符)
职责:转换、过滤、组合数据流
Observable ──→ Operator ──→ Observer
转换/过滤
4. Scheduler(调度器)
职责:控制线程切换
Observable ──→ subscribeOn(后台线程) ──→ observeOn(主线程) ──→ Observer
5. Disposable(可释放)
职责:管理订阅生命周期
let disposable = observable.subscribe(...)
disposable.dispose() // 手动取消订阅
完整流程
┌──────────────┐
│ Observable │ 发送数据
│ (可观察序列) │
└──────┬───────┘
│
↓
┌──────────────┐
│ Operator │ 转换数据
│ (操作符) │
└──────┬───────┘
│
↓
┌──────────────┐
│ Observer │ 接收数据
│ (观察者) │
└──────────────┘
三、Observable(可观察序列)
什么是 Observable?
Observable 是数据流的源头,负责发送事件(next、completed、error)给观察者。
事件类型
enum Event<Element> {
case next(Element) // 发送新元素
case error(Error) // 发送错误
case completed // 发送完成
}
创建 Observable
1. just
用途:发送单个值,然后立即完成
import RxSwift
let observable = Observable.just("Hello, RxSwift!")
observable
.subscribe(onNext: { value in
print("收到值:\(value)")
}, onCompleted: {
print("完成")
})
// 输出:
// 收到值:Hello, RxSwift!
// 完成
2. of
用途:发送多个值,然后完成
import RxSwift
let observable = Observable.of(1, 2, 3, 4, 5)
observable
.subscribe(onNext: { value in
print("收到值:\(value)")
}, onCompleted: {
print("完成")
})
// 输出:
// 收到值:1
// 收到值:2
// 收到值:3
// 收到值:4
// 收到值:5
// 完成
3. from
用途:从数组创建 Observable
import RxSwift
let observable = Observable.from([1, 2, 3, 4, 5])
observable
.subscribe(onNext: { value in
print("收到值:\(value)")
})
// 输出:
// 收到值:1
// 收到值:2
// 收到值:3
// 收到值:4
// 收到值:5
4. empty
用途:不发送任何值,立即完成
import RxSwift
let observable = Observable<Int>.empty()
observable
.subscribe(onNext: { value in
print("收到值:\(value)") // 不会被调用
}, onCompleted: {
print("完成") // 输出:完成
})
5. never
用途:不发送任何值,永不完成
import RxSwift
let observable = Observable<Int>.never()
observable
.subscribe(onNext: { value in
print("收到值:\(value)") // 永远不会被调用
}, onCompleted: {
print("完成") // 永远不会被调用
})
6. error
用途:立即发送错误
import RxSwift
enum NetworkError: Error {
case noConnection
}
let observable = Observable<String>.error(NetworkError.noConnection)
observable
.subscribe(onNext: { value in
print("收到值:\(value)") // 不会被调用
}, onError: { error in
print("错误:\(error)") // 输出:错误:noConnection
})
7. range
用途:发送指定范围内的整数
import RxSwift
let observable = Observable.range(start: 1, count: 5)
observable
.subscribe(onNext: { value in
print("收到值:\(value)")
})
// 输出:
// 收到值:1
// 收到值:2
// 收到值:3
// 收到值:4
// 收到值:5
8. interval
用途:每隔一段时间发送一个整数
import RxSwift
let observable = Observable<Int>.interval(.seconds(1), scheduler: MainScheduler.instance)
let subscription = observable
.subscribe(onNext: { time in
print("时间:\(time)")
})
// 每秒输出:
// 时间:0
// 时间:1
// 时间:2
// ...
// 5 秒后取消订阅
DispatchQueue.main.asyncAfter(deadline: .now() + 5) {
subscription.dispose()
}
9. timer
用途:延迟一段时间后发送一个值
import RxSwift
let observable = Observable<Int>.timer(.seconds(2), scheduler: MainScheduler.instance)
observable
.subscribe(onNext: { value in
print("收到值:\(value)")
}, onCompleted: {
print("完成")
})
// 2 秒后输出:
// 收到值:0
// 完成
10. create
用途:自定义 Observable 创建逻辑
import RxSwift
let observable = Observable<String>.create { observer in
// 发送值
observer.onNext("Hello")
observer.onNext("World")
// 发送完成
observer.onCompleted()
// 返回 Disposable
return Disposables.create()
}
observable
.subscribe(onNext: { value in
print("收到值:\(value)")
})
// 输出:
// 收到值:Hello
// 收到值:World
四、Observer(观察者)
什么是 Observer?
Observer 是数据流的终点,负责接收并处理 Observable 发送的事件。
订阅方式
1. subscribe
用途:最常用的订阅方式,通过闭包处理事件
import RxSwift
let observable = Observable.of(1, 2, 3, 4, 5)
let disposable = observable
.subscribe(onNext: { value in
print("收到值:\(value)")
}, onError: { error in
print("错误:\(error)")
}, onCompleted: {
print("完成")
}, onDisposed: {
print("释放")
})
// 输出:
// 收到值:1
// 收到值:2
// 收到值:3
// 收到值:4
// 收到值:5
// 完成
// 释放
2. subscribe(onNext:)
用途:只处理 next 事件
import RxSwift
let observable = Observable.of(1, 2, 3)
let disposable = observable
.subscribe(onNext: { value in
print("收到值:\(value)")
})
3. bind
用途:将 Observable 绑定到观察者(RxCocoa 特有)
import RxSwift
import RxCocoa
let label = UILabel()
let observable = Observable.just("Hello, RxSwift!")
observable
.bind(to: label.rx.text)
.disposed(by: disposeBag)
Disposable(可释放)
所有订阅都会返回一个 Disposable 对象,用于管理订阅生命周期。
import RxSwift
class MyViewController: UIViewController {
let disposeBag = DisposeBag()
override func viewDidLoad() {
super.viewDidLoad()
let observable = Observable.of(1, 2, 3)
observable
.subscribe(onNext: { value in
print("收到值:\(value)")
})
.disposed(by: disposeBag) // 自动管理订阅
}
// 当 ViewController 被释放时,disposeBag 也会被释放
// 所有订阅自动取消
}
五、Operator(操作符)
什么是 Operator?
Operator 是用于转换、过滤、组合 Observable 的方法。它们返回新的 Observable,形成操作链。
操作符分类
1. 转换操作符
map
用途:转换数据类型
import RxSwift
let observable = Observable.of(1, 2, 3, 4, 5)
observable
.map { $0 * 2 } // 乘以 2
.subscribe(onNext: { value in
print("转换后的值:\(value)")
})
// 输出:
// 转换后的值:2
// 转换后的值:4
// 转换后的值:6
// 转换后的值:8
// 转换后的值:10
flatMap
用途:将每个值转换为新的 Observable,然后合并
import RxSwift
struct User {
let id: Int
let name: String
}
func fetchUserDetails(id: Int) -> Observable<String> {
return Observable.just("用户详情:\(id)")
}
let userIds = Observable.of(1, 2, 3)
userIds
.flatMap { id in
fetchUserDetails(id: id)
}
.subscribe(onNext: { details in
print(details)
})
// 输出:
// 用户详情:1
// 用户详情:2
// 用户详情:3
scan
用途:累积值,类似 reduce
import RxSwift
let observable = Observable.of(1, 2, 3, 4, 5)
observable
.scan(0, accumulator: { accumulator, value in
return accumulator + value
})
.subscribe(onNext: { value in
print("累积值:\(value)")
})
// 输出:
// 累积值:1
// 累积值:3
// 累积值:6
// 累积值:10
// 累积值:15
2. 过滤操作符
filter
用途:过滤不符合条件的值
import RxSwift
let observable = Observable.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
observable
.filter { $0 % 2 == 0 } // 只保留偶数
.subscribe(onNext: { value in
print("偶数:\(value)")
})
// 输出:
// 偶数:2
// 偶数:4
// 偶数:6
// 偶数:8
// 偶数:10
distinctUntilChanged
用途:去除连续重复的值
import RxSwift
let observable = Observable.of("苹果", "苹果", "香蕉", "香蕉", "香蕉", "橙子", "苹果")
observable
.distinctUntilChanged()
.subscribe(onNext: { word in
print("单词:\(word)")
})
// 输出:
// 单词:苹果
// 单词:香蕉
// 单词:橙子
// 单词:苹果
take
用途:只取前 n 个值
import RxSwift
let observable = Observable.of(1, 2, 3, 4, 5)
observable
.take(3) // 只取前 3 个
.subscribe(onNext: { value in
print("值:\(value)")
})
// 输出:
// 值:1
// 值:2
// 值:3
skip
用途:跳过前 n 个值
import RxSwift
let observable = Observable.of(1, 2, 3, 4, 5)
observable
.skip(2) // 跳过前 2 个
.subscribe(onNext: { value in
print("值:\(value)")
})
// 输出:
// 值:3
// 值:4
// 值:5
elementAt
用途:只取指定位置的值
import RxSwift
let observable = Observable.of(1, 2, 3, 4, 5)
observable
.elementAt(2) // 取第 3 个值(索引从 0 开始)
.subscribe(onNext: { value in
print("值:\(value)")
})
// 输出:
// 值:3
3. 时间操作符
debounce
用途:防抖,只在停止发送后的一段时间后才发送最后一个值
import RxSwift
let searchText = PublishSubject<String>()
searchText
.debounce(.milliseconds(500), scheduler: MainScheduler.instance)
.subscribe(onNext: { text in
print("搜索:\(text)")
})
// 快速输入
searchText.onNext("S")
searchText.onNext("Sw")
searchText.onNext("Swi")
searchText.onNext("Swif")
searchText.onNext("Swift")
// 停止输入 500ms 后才输出:
// 搜索:Swift
throttle
用途:节流,在指定时间间隔内只发送第一个或最后一个值
import RxSwift
let buttonTaps = PublishSubject<Void>()
buttonTaps
.throttle(.seconds(1), scheduler: MainScheduler.instance, latest: false)
.subscribe(onNext: { _ in
print("按钮被点击")
})
// 快速点击多次
buttonTaps.onNext(())
buttonTaps.onNext(())
buttonTaps.onNext(())
// 1 秒内只输出一次:
// 按钮被点击
delay
用途:延迟发送值
import RxSwift
let observable = Observable.just("Hello")
observable
.delay(.seconds(2), scheduler: MainScheduler.instance)
.subscribe(onNext: { value in
print("收到值:\(value)")
})
// 2 秒后输出:
// 收到值:Hello
timeout
用途:超时处理
import RxSwift
let observable = Observable<Int>.never()
observable
.timeout(.seconds(5), scheduler: MainScheduler.instance)
.subscribe(onNext: { value in
print("收到值:\(value)")
}, onError: { error in
print("超时错误:\(error)")
})
// 5 秒后输出:
// 超时错误:RxError.timeout
4. 组合操作符
merge
用途:合并多个 Observable,按时间顺序发送值
import RxSwift
let subject1 = PublishSubject<Int>()
let subject2 = PublishSubject<Int>()
Observable.merge(subject1, subject2)
.subscribe(onNext: { value in
print("值:\(value)")
})
subject1.onNext(1)
subject2.onNext(4)
subject1.onNext(2)
subject2.onNext(5)
// 输出:
// 值:1
// 值:4
// 值:2
// 值:5
zip
用途:将多个 Observable 的值配对发送
import RxSwift
let subject1 = PublishSubject<String>()
let subject2 = PublishSubject<Int>()
Observable.zip(subject1, subject2)
.subscribe(onNext: { fruit, count in
print("\(fruit): \(count)个")
})
subject1.onNext("苹果")
subject1.onNext("香蕉")
subject2.onNext(5)
subject2.onNext(3)
// 输出:
// 苹果: 5个
// 香蕉: 3个
combineLatest
用途:当任一 Observable 发送值时,发送所有 Observable 的最新值
import RxSwift
let username = BehaviorSubject(value: "")
let password = BehaviorSubject(value: "")
Observable.combineLatest(username, password)
.subscribe(onNext: { user, pass in
print("用户名:\(user),密码:\(pass)")
})
username.onNext("张三")
password.onNext("123456")
username.onNext("李四")
// 输出:
// 用户名:,密码:
// 用户名:张三,密码:
// 用户名:张三,密码:123456
// 用户名:李四,密码:123456
withLatestFrom
用途:当一个 Observable 发送值时,取另一个 Observable 的最新值
import RxSwift
let buttonTaps = PublishSubject<Void>()
let text = BehaviorSubject(value: "")
buttonTaps
.withLatestFrom(text)
.subscribe(onNext: { text in
print("搜索:\(text)")
})
text.onNext("Swift")
buttonTaps.onNext(())
text.onNext("RxSwift")
buttonTaps.onNext(())
// 输出:
// 搜索:Swift
// 搜索:RxSwift
5. 错误处理操作符
catchError
用途:捕获错误并返回新的 Observable
import RxSwift
enum NetworkError: Error {
case noConnection
}
let observable = Observable<Int>.error(NetworkError.noConnection)
observable
.catchError { error in
print("捕获错误:\(error)")
return Observable.just(0) // 返回默认值
}
.subscribe(onNext: { value in
print("值:\(value)")
})
// 输出:
// 捕获错误:noConnection
// 值:0
retry
用途:失败后重试
import RxSwift
var attemptCount = 0
func fetchData() -> Observable<String> {
attemptCount += 1
if attemptCount < 3 {
return Observable.error(NSError(domain: "Network", code: -1, userInfo: nil))
} else {
return Observable.just("成功")
}
}
fetchData()
.retry(3) // 最多重试 3 次
.subscribe(onNext: { value in
print("值:\(value)")
}, onError: { error in
print("错误:\(error)")
})
// 输出:
// 值:成功
catchErrorJustReturn
用途:捕获错误并返回单个值
import RxSwift
let observable = Observable<Int>.error(NetworkError.noConnection)
observable
.catchErrorJustReturn(0) // 捕获错误并返回 0
.subscribe(onNext: { value in
print("值:\(value)")
})
// 输出:
// 值:0
六、Subject
什么是 Subject?
Subject 是一种特殊的 Observable,可以手动发送值。它既是 Observable,也是 Observer。
常见的 Subject
1. PublishSubject
用途:不保存当前值,只发送新值给订阅者
import RxSwift
let subject = PublishSubject<String>()
// 第一次订阅
subject
.subscribe(onNext: { value in
print("订阅者 1:\(value)")
})
// 发送值
subject.onNext("Hello")
subject.onNext("World")
// 第二次订阅
subject
.subscribe(onNext: { value in
print("订阅者 2:\(value)")
})
subject.onNext("RxSwift")
// 输出:
// 订阅者 1:Hello
// 订阅者 1:World
// 订阅者 1:RxSwift
// 订阅者 2:RxSwift
2. BehaviorSubject
用途:保存当前值,新订阅者会立即收到当前值
import RxSwift
let subject = BehaviorSubject(value: "初始值")
// 第一次订阅
subject
.subscribe(onNext: { value in
print("订阅者 1:\(value)")
})
// 发送值
subject.onNext("新值")
// 第二次订阅
subject
.subscribe(onNext: { value in
print("订阅者 2:\(value)")
})
// 输出:
// 订阅者 1:初始值 ← 立即收到当前值
// 订阅者 1:新值
// 订阅者 2:新值 ← 立即收到当前值
3. ReplaySubject
用途:保存指定数量的历史值,新订阅者会收到历史值
import RxSwift
let subject = ReplaySubject<String>.create(bufferSize: 2)
// 发送值
subject.onNext("值1")
subject.onNext("值2")
subject.onNext("值3")
// 订阅
subject
.subscribe(onNext: { value in
print("订阅者:\(value)")
})
// 输出:
// 订阅者:值2 ← 收到最近 2 个值
// 订阅者:值3
4. AsyncSubject
用途:只发送最后一个值,且只在完成时发送
import RxSwift
let subject = AsyncSubject<String>()
subject
.subscribe(onNext: { value in
print("值:\(value)")
}, onCompleted: {
print("完成")
})
subject.onNext("值1")
subject.onNext("值2")
subject.onNext("值3")
subject.onCompleted() // 只在完成时发送最后一个值
// 输出:
// 值:值3
// 完成
Subject 的使用场景
1. 用户输入
import RxSwift
class SearchViewModel {
let searchText = BehaviorSubject(value: "")
private let disposeBag = DisposeBag()
init() {
searchText
.debounce(.milliseconds(300), scheduler: MainScheduler.instance)
.distinctUntilChanged()
.subscribe(onNext: { [weak self] text in
self?.performSearch(query: text)
})
.disposed(by: disposeBag)
}
private func performSearch(query: String) {
print("搜索:\(query)")
}
}
2. 事件通知
import RxSwift
class NotificationManager {
static let shared = NotificationManager()
let userDidLogin = PublishSubject<User>()
let userDidLogout = PublishSubject<Void>()
private init() {}
}
// 发送通知
NotificationManager.shared.userDidLogin.onNext(user)
// 订阅通知
NotificationManager.shared.userDidLogin
.subscribe(onNext: { user in
print("用户登录:\(user.name)")
})
.disposed(by: disposeBag)
七、Scheduler(调度器)
什么是 Scheduler?
Scheduler 是 RxSwift 的线程管理系统,用于控制 Observable 和 Observer 在哪个线程上执行。
常见的 Scheduler
1. MainScheduler
用途:主线程,用于 UI 更新
import RxSwift
observable
.observeOn(MainScheduler.instance) // 切换到主线程
.subscribe(onNext: { value in
self.label.text = value // 安全更新 UI
})
2. SerialDispatchQueueScheduler
用途:串行队列,按顺序执行
import RxSwift
let scheduler = SerialDispatchQueueScheduler(qos: .background)
observable
.subscribeOn(scheduler) // 在后台线程订阅
.subscribe(onNext: { value in
print("在后台线程处理:\(value)")
})
3. ConcurrentDispatchQueueScheduler
用途:并发队列,并发执行
import RxSwift
let scheduler = ConcurrentDispatchQueueScheduler(qos: .background)
observable
.subscribeOn(scheduler) // 在后台线程订阅
.subscribe(onNext: { value in
print("在后台线程处理:\(value)")
})
4. OperationQueueScheduler
用途:基于 OperationQueue 的调度器
import RxSwift
let scheduler = OperationQueueScheduler(operationQueue: OperationQueue())
observable
.subscribeOn(scheduler)
.subscribe(onNext: { value in
print("在 OperationQueue 中处理:\(value)")
})
subscribeOn vs observeOn
subscribeOn
用途:指定 Observable 在哪个线程上创建和执行
import RxSwift
Observable.create { observer in
print("创建 Observable:\(Thread.current)") // 后台线程
observer.onNext("Hello")
observer.onCompleted()
return Disposables.create()
}
.subscribeOn(ConcurrentDispatchQueueScheduler(qos: .background)) // 在后台线程创建
.observeOn(MainScheduler.instance) // 在主线程观察
.subscribe(onNext: { value in
print("接收值:\(Thread.current)") // 主线程
})
observeOn
用途:指定 Observer 在哪个线程上接收事件
import RxSwift
Observable.of(1, 2, 3)
.observeOn(MainScheduler.instance) // 在主线程观察
.subscribe(onNext: { value in
print("接收值:\(Thread.current)") // 主线程
})
线程切换示例
import RxSwift
Observable.create { observer in
// 后台线程:网络请求
print("1. 网络请求:\(Thread.current)")
let data = self.fetchDataFromNetwork()
observer.onNext(data)
observer.onCompleted()
return Disposables.create()
}
.subscribeOn(ConcurrentDispatchQueueScheduler(qos: .background)) // 后台线程执行
.observeOn(MainScheduler.instance) // 主线程观察
.subscribe(onNext: { data in
// 主线程:更新 UI
print("2. 更新 UI:\(Thread.current)")
self.updateUI(with: data)
})
// 输出:
// 1. 网络请求:后台线程
// 2. 更新 UI:主线程
八、实战示例
示例 1:搜索功能
import RxSwift
import RxCocoa
import UIKit
class SearchViewController: UIViewController {
private let searchBar = UISearchBar()
private let tableView = UITableView()
private let viewModel = SearchViewModel()
private let disposeBag = DisposeBag()
override func viewDidLoad() {
super.viewDidLoad()
setupUI()
setupBindings()
}
private func setupBindings() {
// 绑定搜索框输入到 ViewModel
searchBar.rx.text
.orEmpty
.bind(to: viewModel.searchText)
.disposed(by: disposeBag)
// 绑定搜索结果到 TableView
viewModel.searchResults
.bind(to: tableView.rx.items(cellIdentifier: "Cell")) { index, user, cell in
cell.textLabel?.text = user.name
}
.disposed(by: disposeBag)
}
}
class SearchViewModel {
let searchText = BehaviorSubject(value: "")
let searchResults: Observable<[User]>
private let disposeBag = DisposeBag()
init() {
searchResults = searchText
.debounce(.milliseconds(300), scheduler: MainScheduler.instance)
.distinctUntilChanged()
.filter { !$0.isEmpty }
.flatMapLatest { text in
self.performSearch(query: text)
}
.catchErrorJustReturn([])
}
private func performSearch(query: String) -> Observable<[User]> {
return Observable.create { observer in
UserService.shared.searchUsers(query: query) { result in
switch result {
case .success(let users):
observer.onNext(users)
observer.onCompleted()
case .failure:
observer.onNext([])
observer.onCompleted()
}
}
return Disposables.create()
}
}
}
示例 2:表单验证
import RxSwift
import RxCocoa
class RegistrationViewModel {
let username = BehaviorSubject(value: "")
let email = BehaviorSubject(value: "")
let password = BehaviorSubject(value: "")
let confirmPassword = BehaviorSubject(value: "")
let isFormValid: Observable<Bool>
let errorMessage: Observable<String?>
private let disposeBag = DisposeBag()
init() {
// 用户名验证
let isUsernameValid = username
.map { $0.count >= 3 }
// 邮箱验证
let isEmailValid = email
.map { $0.contains("@") && $0.contains(".") }
// 密码验证
let isPasswordValid = password
.map { $0.count >= 6 }
// 确认密码验证
let isConfirmPasswordValid = Observable.combineLatest(password, confirmPassword)
.map { $0 == $1 }
// 组合所有验证
isFormValid = Observable.combineLatest(
isUsernameValid,
isEmailValid,
isPasswordValid,
isConfirmPasswordValid
) { $0 && $1 && $2 && $3 }
// 错误消息
errorMessage = Observable.combineLatest(
isUsernameValid,
isEmailValid,
isPasswordValid,
isConfirmPasswordValid
) { usernameValid, emailValid, passwordValid, confirmValid in
if !usernameValid { return "用户名至少 3 个字符" }
if !emailValid { return "邮箱格式不正确" }
if !passwordValid { return "密码至少 6 个字符" }
if !confirmValid { return "两次密码不一致" }
return nil
}
}
func register() {
guard let isValid = try? isFormValid.toBlocking().first(), isValid else { return }
// 执行注册逻辑
print("注册成功!")
}
}
示例 3:网络请求
import RxSwift
import RxCocoa
class UserListViewModel {
let users = BehaviorSubject(value: [User]())
let isLoading = BehaviorSubject(value: false)
let errorMessage = BehaviorSubject(value: nil as String?)
private let disposeBag = DisposeBag()
func loadUsers() {
isLoading.onNext(true)
errorMessage.onNext(nil)
// 使用 RxSwift 进行网络请求
guard let url = URL(string: "https://jsonplaceholder.typicode.com/users") else {
return
}
let request = URLRequest(url: url)
URLSession.shared.rx.data(request: request)
.map { data in
try JSONDecoder().decode([User].self, from: data)
}
.observeOn(MainScheduler.instance)
.subscribe(onNext: { [weak self] users in
self?.users.onNext(users)
self?.isLoading.onNext(false)
}, onError: { [weak self] error in
self?.errorMessage.onNext(error.localizedDescription)
self?.isLoading.onNext(false)
})
.disposed(by: disposeBag)
}
}
示例 4:定时器
import RxSwift
class TimerViewModel {
let counter = BehaviorSubject(value: 0)
let isRunning = BehaviorSubject(value: false)
private var timerDisposable: Disposable?
func startTimer() {
isRunning.onNext(true)
timerDisposable = Observable<Int>.interval(.seconds(1), scheduler: MainScheduler.instance)
.subscribe(onNext: { [weak self] _ in
if let currentCount = try? self?.counter.value() {
self?.counter.onNext(currentCount + 1)
}
})
}
func stopTimer() {
isRunning.onNext(false)
timerDisposable?.dispose()
timerDisposable = nil
}
func resetTimer() {
stopTimer()
counter.onNext(0)
}
}
示例 5:按钮点击防抖
import RxSwift
import RxCocoa
class ViewController: UIViewController {
private let submitButton = UIButton()
private let disposeBag = DisposeBag()
override func viewDidLoad() {
super.viewDidLoad()
// 按钮点击防抖:1 秒内只响应一次
submitButton.rx.tap
.throttle(.seconds(1), scheduler: MainScheduler.instance)
.subscribe(onNext: { [weak self] in
self?.submitForm()
})
.disposed(by: disposeBag)
}
private func submitForm() {
print("提交表单")
}
}
九、最佳实践
1. 内存管理
// ✅ 正确:使用 [weak self] 避免循环引用
observable
.subscribe(onNext: { [weak self] value in
self?.updateUI(with: value)
})
.disposed(by: disposeBag)
// ❌ 错误:循环引用
observable
.subscribe(onNext: { value in
self.updateUI(with: value) // 循环引用!
})
.disposed(by: disposeBag)
2. 线程管理
// ✅ 正确:明确指定线程
URLSession.shared.rx.data(request: request)
.observeOn(MainScheduler.instance) // 切换到主线程
.subscribe(onNext: { data in
self.tableView.reloadData() // 安全更新 UI
})
// ❌ 错误:在后台线程更新 UI
URLSession.shared.rx.data(request: request)
.subscribe(onNext: { data in
self.tableView.reloadData() // 可能崩溃!
})
3. 错误处理
// ✅ 正确:统一错误处理
enum AppError: Error {
case networkError
case decodingError
case unknown
}
func fetchUsers() -> Observable<[User]> {
return URLSession.shared.rx.data(request: request)
.map { data in
do {
return try JSONDecoder().decode([User].self, from: data)
} catch {
throw AppError.decodingError
}
}
.catchError { error in
if let appError = error as? AppError {
return Observable.error(appError)
}
return Observable.error(AppError.unknown)
}
}
// 使用
fetchUsers()
.subscribe(onNext: { users in
self.users = users
}, onError: { error in
if let appError = error as? AppError {
// 统一处理错误
self.showErrorAlert(error: appError)
}
})
4. 取消订阅
class MyViewController: UIViewController {
private let disposeBag = DisposeBag()
override func viewDidLoad() {
super.viewDidLoad()
// 订阅会自动存储在 disposeBag 中
observable
.subscribe(onNext: { value in
// 处理数据
})
.disposed(by: disposeBag)
}
// 当 ViewController 被释放时,disposeBag 也会被释放
// 所有订阅自动取消
}
5. 使用 Driver
用途:Driver 是 RxCocoa 提供的特殊 Observable,保证在主线程上观察,不会发送错误
import RxCocoa
// ✅ 推荐:使用 Driver 处理 UI 绑定
let searchResults = viewModel.searchResults
.asDriver(onErrorJustReturn: []) // 错误时返回空数组
searchResults
.drive(tableView.rx.items(cellIdentifier: "Cell")) { index, user, cell in
cell.textLabel?.text = user.name
}
.disposed(by: disposeBag)
// ❌ 不推荐:直接使用 Observable
viewModel.searchResults
.bind(to: tableView.rx.items(cellIdentifier: "Cell")) { index, user, cell in
cell.textLabel?.text = user.name
}
.disposed(by: disposeBag)
十、RxSwift vs Combine
对比表
| 特性 | RxSwift | Combine |
|---|---|---|
| 维护者 | 社区 | Apple |
| 支持平台 | iOS 8+, macOS 10.10+ | iOS 13+, macOS 10.15+ |
| 学习曲线 | 较陡 | 中等 |
| 与 SwiftUI 配合 | 需要适配 | 完美 |
| 社区资源 | 丰富的第三方库 | 官方文档 |
| 操作符数量 | 更多 | 较少但够用 |
| 调试工具 | RxBlocking、RxTest | 内置 |
概念对应
| RxSwift | Combine |
|---|---|
| Observable | Publisher |
| Observer | Subscriber |
| Operator | Operator |
| Scheduler | Scheduler |
| Disposable | AnyCancellable |
| Subject | Subject |
| PublishSubject | PassthroughSubject |
| BehaviorSubject | CurrentValueSubject |
| Driver | Publisher + receive(on:) |
代码对比
RxSwift
import RxSwift
let observable = Observable.of(1, 2, 3, 4, 5)
observable
.map { $0 * 2 }
.filter { $0 > 5 }
.subscribe(onNext: { value in
print("值:\(value)")
})
.disposed(by: disposeBag)
Combine
import Combine
let publisher = [1, 2, 3, 4, 5].publisher
publisher
.map { $0 * 2 }
.filter { $0 > 5 }
.sink { value in
print("值:\(value)")
}
.store(in: &cancellables)
十一、常见问题
Q1: RxSwift 和 Combine 应该选择哪个?
A: 根据项目需求选择:
选择 RxSwift:
- 需要支持 iOS 13 以下版本
- 需要丰富的操作符
- 团队已经熟悉 RxSwift
- 需要使用 RxCommunity 的第三方库
选择 Combine:
- 只支持 iOS 13+
- 使用 SwiftUI
- 希望使用 Apple 原生框架
- 项目简单,不需要太多操作符
Q2: 如何调试 RxSwift 数据流?
A: 使用 debug 操作符:
observable
.debug("数据流") // 打印所有事件
.subscribe(onNext: { value in
// 处理数据
})
// 输出:
// 数据流 -> subscribed
// 数据流 -> Event next(1)
// 数据流 -> Event next(2)
// 数据流 -> Event completed
// 数据流 -> isDisposed
Q3: BehaviorSubject 和 PublishSubject 有什么区别?
A: 主要区别:
// PublishSubject:不保存当前值,只发送新值
let publishSubject = PublishSubject<String>()
publishSubject.onNext("值1")
publishSubject
.subscribe(onNext: { value in
print("订阅者:\(value)") // 不会收到"值1"
})
publishSubject.onNext("值2") // 输出:订阅者:值2
// BehaviorSubject:保存当前值,新订阅者立即收到当前值
let behaviorSubject = BehaviorSubject(value: "初始值")
behaviorSubject
.subscribe(onNext: { value in
print("订阅者:\(value)") // 输出:订阅者:初始值
})
Q4: 如何在 SwiftUI 中使用 RxSwift?
A: 使用 asObservable() 或 asDriver():
import SwiftUI
import RxSwift
import RxCocoa
class ViewModel: ObservableObject {
@Published var username: String = ""
private let disposeBag = DisposeBag()
init() {
// 将 @Published 属性转换为 Observable
$username.asObservable()
.debounce(.milliseconds(300), scheduler: MainScheduler.instance)
.distinctUntilChanged()
.subscribe(onNext: { username in
print("用户名:\(username)")
})
.disposed(by: disposeBag)
}
}
Q5: 如何处理多个订阅?
A: 使用 share() 操作符共享订阅:
// ❌ 错误:每次订阅都会触发网络请求
let observable = URLSession.shared.rx.data(request: request)
observable.subscribe(onNext: { _ in }) // 第一次网络请求
observable.subscribe(onNext: { _ in }) // 第二次网络请求
// ✅ 正确:共享订阅
let observable = URLSession.shared.rx.data(request: request)
.share()
observable.subscribe(onNext: { _ in }) // 只有一次网络请求
observable.subscribe(onNext: { _ in }) // 共享结果
Q6: flatMap 和 flatMapLatest 有什么区别?
A: 主要区别:
// flatMap:保留所有内部 Observable
let searchResults = searchText
.flatMap { text in
self.searchAPI(query: text) // 所有请求都会继续
}
// flatMapLatest:只保留最新的内部 Observable
let searchResults = searchText
.flatMapLatest { text in
self.searchAPI(query: text) // 旧请求会被取消
}
// 示例:
// 输入 "Sw" -> 发起请求 1
// 输入 "Swi" -> 发起请求 2
// flatMap:请求 1 和请求 2 都会完成
// flatMapLatest:请求 1 被取消,只有请求 2 完成
总结
RxSwift 的优势
✅ 统一异步处理:所有异步事件用统一的方式处理
✅ 声明式编程:代码简洁、易读
✅ 丰富的操作符:强大的数据转换和过滤能力
✅ 成熟的生态系统:丰富的第三方库和社区支持
✅ 跨平台:支持 iOS 8+,兼容性好
学习路径
- 基础概念 → 理解 Observable、Observer、Operator
- 常用操作符 → map、filter、debounce 等
- Subject → PublishSubject、BehaviorSubject
- Scheduler → 线程管理
- 实战应用 → 网络请求、表单验证、搜索功能
- 高级特性 → 自定义 Operator、错误处理
- RxCocoa → UI 绑定、Driver