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

RxSwift:如何将单例中不同类型Single<T>加入队列串行执行

解决方案:用RxSwift Concat实现串行操作队列

这个需求用RxSwift的concat操作符刚好能解决,核心思路是在单例内部维护一个共享的串行Observable队列,每次新操作都追加到队列末尾,严格保证前一个操作完成后才启动下一个。

完整实现代码

import RxSwift

class Singleton { 
    static let shared = Singleton() 
    private init() { } 
    
    private let disposeBag = DisposeBag()
    // 初始队列是一个已完成的Observable,确保第一个操作能立即执行
    private var currentQueue: Observable<Void> = Observable.just(())
    
    private func placeInQueue<T>(operation: Single<T>) -> Single<T> {
        return Single.create { [weak self] observer in
            guard let self = self else {
                observer(.error(NSError(domain: "Singleton", code: -1, userInfo: [NSLocalizedDescriptionKey: "Singleton instance deallocated"])))
                return Disposables.create()
            }
            
            // 将操作转换为Void类型的Observable,适配队列类型
            let wrappedOperation = operation
                .asObservable()
                .map { _ in () } // 忽略操作返回值,只追踪完成状态
                .catch { _ in
                    // 可选:如果想让队列在操作失败后继续执行,就保留这个catch
                    // 如果希望队列在出错时终止,直接删掉这个catch即可
                    return Observable.just(())
                }
            
            // 把新操作追加到现有队列末尾,生成新的串行队列
            let newQueue = self.currentQueue.concat(wrappedOperation)
            self.currentQueue = newQueue
            
            // 订阅原操作,把结果传递给外部调用者
            let operationDisposable = operation.subscribe(
                onSuccess: { observer(.success($0)) },
                onError: { observer(.error($0)) }
            )
            
            // 必须订阅新队列,因为Rx的Observable是冷序列,订阅才会触发执行
            let queueDisposable = newQueue.subscribe()
            return Disposables.create(operationDisposable, queueDisposable)
        }
    }
    
    func doSomethingInt() -> Single<Int> { 
        let operation = Single.just(1)
            .delay(3, scheduler: MainScheduler.instance)
        return placeInQueue(operation: operation)
    } 
    
    func doSomethingString() -> Single<String> { 
        let operation = Single.just("Wow")
            .delay(3, scheduler: MainScheduler.instance)
        return placeInQueue(operation: operation)
    } 
}

关键逻辑说明

  • 共享队列维护:currentQueue作为串行流的载体,初始值是一个已完成的Observable<Void>,确保第一个操作能立即启动。
  • 操作包装与追加:placeInQueue方法把传入的Single<T>包装成Observable<Void>(因为concat要求序列类型一致),然后用concat将新操作追加到队列末尾,再更新currentQueue为新的串行流。
  • 错误处理:catch操作是可选的——如果希望某个操作失败后,后续操作依然能执行,就保留它;如果希望队列在出错时直接终止,删掉这个catch即可。
  • 冷序列触发:因为Rx的Observable是冷序列,必须订阅才会执行,所以我们需要订阅newQueue来确保整个串行流被触发。

测试验证

比如你连续调用三次方法:

let singleton = Singleton.shared

singleton.doSomethingInt().subscribe(onSuccess: { print($0) })
singleton.doSomethingString().subscribe(onSuccess: { print($0) })
singleton.doSomethingInt().subscribe(onSuccess: { print($0) })

输出会依次在3秒、6秒、9秒后打印:1 → Wow → 1,完全符合串行执行的需求。

内容的提问来源于stack exchange,提问作者Wiingaard

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:21:07