RxSwift 框架详解

一、什么是 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+,兼容性好

学习路径

  1. 基础概念 → 理解 Observable、Observer、Operator
  2. 常用操作符 → map、filter、debounce 等
  3. Subject → PublishSubject、BehaviorSubject
  4. Scheduler → 线程管理
  5. 实战应用 → 网络请求、表单验证、搜索功能
  6. 高级特性 → 自定义 Operator、错误处理
  7. RxCocoa → UI 绑定、Driver

参考资料


最后编辑于 :
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容