RxSwift异步单元测试中Observable仅订阅一次问题求助
问题分析与解决方案
首先,咱们来拆解一下你遇到的问题根源:
核心原因
- flatMapLatest的行为特性:
flatMapLatest的设计逻辑就是当上游Observable发送新元素时,会立即取消前一个由它创建的Observable的订阅,转而订阅新的Observable。在你的测试中,input在100时间点发送1,紧接着200时间点发送10——这时候flatMapLatest会立刻取消第一个对应1的Observable订阅。 - 真实异步与测试调度器的冲突:你在
Observable.create里用了DispatchQueue.main.asyncAfter,这是真实线程的异步操作,而TestScheduler是虚拟时间调度器。当scheduler.start()执行后,input的两个事件会被快速触发,虽然第一个asyncAfter的闭包最终还是执行了(所以打印了next 1),但此时对应的Observable已经被disposed,observer.onNext(value)无法传递到下游的订阅者,自然看不到subscribe 1的输出。
解决方案
下面提供两种针对性的解决办法,你可以根据业务需求选择:
方案一:使用测试调度器模拟异步(推荐用于单元测试)
在RxSwift单元测试中,应该用测试调度器来模拟异步操作,这样能精准控制时间,避免真实线程的干扰。我们可以通过依赖注入把测试调度器传入ViewModel:
class TestViewModel: NSObject { let result : Observable<Int> // 新增scheduler参数,支持依赖注入 init(input: Observable<Int>, scheduler: SchedulerType) { result = input.flatMapLatest({ value -> Observable<Int> in return Observable.create({ observer -> Disposable in // 使用测试调度器的scheduleRelative模拟延迟 return scheduler.scheduleRelative((), dueTime: .seconds(1)) { _ in print("next \(value)") observer.onNext(value) return Disposables.create() } }) }) } } func testCount() { let expectation = XCTestExpectation(description: "async") let input = scheduler.createHotObservable([.next(100, 1), .next(200, 10)]) // 传入测试调度器scheduler let viewModel = TestViewModel(input: input.asObservable(), scheduler: scheduler) var receivedValues: [Int] = [] viewModel.result.subscribe(onNext: { value in print("subscribe \(value)") receivedValues.append(value) // 收到两个值时标记期望完成 if receivedValues.count == 2 { expectation.fulfill() } }).disposed(by: disposeBag) scheduler.start() wait(for: [expectation], timeout: timeout) // 验证结果是否符合预期 XCTAssertEqual(receivedValues, [1, 10]) }
方案二:替换flatMapLatest为flatMap
如果你业务上需要保留所有上游事件对应的异步操作(不取消之前的订阅),可以把flatMapLatest换成flatMap——它会同时订阅所有由上游元素创建的Observable,不会取消之前的订阅,这样两个onNext都会被传递到下游:
class TestViewModel: NSObject { let result : Observable<Int> init(input: Observable<Int>) { // 把flatMapLatest替换为flatMap result = input.flatMap({ value -> Observable<Int> in return Observable.create({ observer -> Disposable in DispatchQueue.main.asyncAfter(deadline: DispatchTime.now() + 1, execute: { print("next \(value)") observer.onNext(value) }) return Disposables.create() }) }) } }
注意:这种方式在input事件频繁时会同时存在多个异步操作,需要根据实际业务场景评估是否合适。
另外补充一点:你的测试代码里expectation没有被触发,需要在收到第二个值时调用expectation.fulfill(),否则测试会超时哦。
内容的提问来源于stack exchange,提问作者Karim
相关产品推荐
相关产品推荐

