RxSwift如何在不终止先前异步任务时仅订阅最新结果?
保留所有异步任务执行但仅订阅最新结果(RxSwift)
我正在用RxSwift做一个异步任务的模拟,代码如下,目的是把Int类型的Observable转换成String类型的:
let pubishSubject = PublishSubject<Int>() pubishSubject.asObservable() .debug("before") .flatMap ({ (value) -> Observable<String> in let task: Observable<String> = Observable.create { observer in DispatchQueue.main.asyncAfter(deadline: .now() + 1, execute: { print("Executed \(value)") observer.on(.next("Async \(value)")) observer.on(.completed) }) return Disposables.create(with: { print("Disposed \(value)") }) } return task }) .debug("after") .subscribe() pubishSubject.onNext(1) pubishSubject.onNext(2) pubishSubject.onNext(3)
当前的调试输出会打印所有三个异步任务的结果:
2018-06-01 14:08:35.748: after -> subscribed 2018-06-01 14:08:35.749: before -> subscribed 2018-06-01 14:08:35.751: before -> Event next(1) 2018-06-01 14:08:35.753: before -> Event next(2) 2018-06-01 14:08:35.753: before -> Event next(3) Executed 1 2018-06-01 14:08:36.785: after -> Event next(Async 1) Disposed 1 Executed 2 2018-06-01 14:08:36.786: after -> Event next(Async 2) Disposed 2 Executed 3 2018-06-01 14:08:36.787: after -> Event next(Async 3) Disposed 3
我的需求是:所有异步任务必须完整执行(不能用flatMapLatest,因为它会取消之前的任务),但最终我只需要订阅并接收最新的那个结果(也就是"Async 3")。我完全没头绪,有没有合适的RxSwift操作符或者思路能解决这个问题?
解决方案:跟踪任务序号过滤结果
你需要的是一种“让所有任务跑完,但只放行最后一个任务结果”的逻辑,这里可以通过给每个任务标记序号,然后过滤出序号匹配最新值的结果来实现,具体步骤如下:
- 给每个原始元素分配唯一递增序号:用
scan操作符给PublishSubject发出的每个Int值绑定一个递增的序号,生成(value: Int, index: Int)的序列,这样每个异步任务都能关联到自己的序号。 - 异步任务携带序号返回结果:在
flatMap中,把异步任务的结果和对应的序号一起返回,变成Observable<(String, Int)>类型。 - 结合最新序号过滤结果:用
combineLatest跟踪当前最新的序号,然后和每个异步结果对比,只保留结果序号等于最新序号的元素,最后提取出String结果。
修改后的完整代码:
let pubishSubject = PublishSubject<Int>() // 生成带有序号的序列,share保证多订阅共享同一个序号流 let numberedSequence = pubishSubject .scan((value: 0, index: 0)) { previous, current in (value: current, index: previous.index + 1) } .share(replay: 1) // 执行异步任务,返回结果+序号的序列 let asyncResultsWithIndex = numberedSequence .debug("before") .flatMap { item -> Observable<(String, Int)> in Observable.create { observer in DispatchQueue.main.asyncAfter(deadline: .now() + 1) { print("Executed \(item.value)") observer.onNext(("Async \(item.value)", item.index)) observer.onCompleted() } return Disposables.create { print("Disposed \(item.value)") } } } // 过滤出仅最新序号对应的结果 numberedSequence .map(\.index) // 提取最新序号 .combineLatest(asyncResultsWithIndex) { latestIndex, result in (result.0, result.1, latestIndex) } .filter { _, resultIndex, latestIndex in resultIndex == latestIndex // 只保留序号匹配的结果 } .map(\.0) // 提取最终的String结果 .debug("after") .subscribe() pubishSubject.onNext(1) pubishSubject.onNext(2) pubishSubject.onNext(3)
运行这段代码后,你会看到所有三个异步任务依然会执行(Executed 1、Executed 2、Executed 3都会打印),但after的调试输出只会显示:
20XX-XX-XX XX:XX:XX.XXX: after -> subscribed 20XX-XX-XX XX:XX:XX.XXX: after -> Event next(Async 3)
完全符合你的需求:既保留了所有异步任务的执行,又只接收最新的结果。
核心逻辑说明
scan生成的序号是递增的,确保每个元素都有唯一标识,share(replay:1)避免了多订阅导致序号重复生成的问题。flatMap中绑定序号,让每个异步结果都能追溯到它对应的原始元素,方便后续判断是否为最新任务的结果。combineLatest会持续监听最新的序号,每次有异步结果返回时,都会和当前最新序号对比,只有匹配的结果才会被传递到下游,从而过滤掉所有中间结果。
这个方案完美避开了flatMapLatest会取消之前任务的问题,同时实现了你要的“队列式”逻辑:所有任务都执行,但只返回最后一个的结果。
内容的提问来源于stack exchange,提问作者Jagger
相关产品推荐
相关产品推荐

