You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 09:07:44