如何用Combine实现支持重试的Swift多上传任务动态队列
基于Combine的多任务上传队列实现思路
- 首先设计可观测队列类,需要符合
ObservableObject协议,才能直接和SwiftUI的状态监听能力打通。内部用CurrentValueSubject持有待执行任务列表,单独维护每个任务的剩余重试次数,避免污染原始任务对象。 - 建议对原始任务做一层封装,不要直接使用裸函数,比如定义
UploadTask结构体,内置任务执行闭包、唯一标识、剩余重试次数字段,方便后续任务追踪和管理。 - 任务执行流用
flatMap(maxPublishers: .max(并发数))控制同时执行的任务上限,比如设置.max(2)即可实现最多2个任务并行上传的效果。 - 单个任务的执行流套上Combine原生的
retry(_:)操作符,传入你配置的重试次数,任务执行失败抛出错误后,Combine会自动完成重试逻辑,不需要自己手写重试调度。 - 任务执行成功或重试全部失败后,直接从待执行任务列表中移除对应任务即可,剩余任务数的属性用
@Published修饰,SwiftUI可以直接绑定该属性自动刷新UI。
最简实现示例
import Combine import SwiftUI // 上传任务定义,可根据业务需求修改返回值类型 typealias UploadTask = () async throws -> Data class SomeQueue: ObservableObject { /// 未完成任务数,直接供SwiftUI绑定 @Published var remaining: Int = 0 /// 待执行任务队列 private var tasks = CurrentValueSubject<[UploadTask], Never>([]) /// 全局最大重试次数 private var maxRetryCount: Int = 0 /// 任务最终失败回调 private var failureHandler: ((Error) -> Void)? /// 并发数限制,可根据需求修改或对外暴露配置 private let concurrencyLimit: Int = 3 private var cancellables = Set<AnyCancellable>() init(from initialTasks: [UploadTask]) { tasks.send(initialTasks) remaining = initialTasks.count } /// 配置重试次数,链式调用返回自身 func retry(_ count: Int) -> Self { maxRetryCount = count return self } /// 配置失败回调,链式调用返回自身 func onFailure(_ handler: @escaping (Error) -> Void) -> Self { failureHandler = handler return self } /// 启动队列执行,链式调用返回自身 func execute() -> Self { tasks .flatMap { tasks -> AnyPublisher<UploadTask, Never> in Publishers.Sequence(sequence: tasks) .eraseToAnyPublisher() } .flatMap(maxPublishers: .max(concurrencyLimit)) { [weak self] task -> AnyPublisher<Data, Error> in Future { promise in Task { do { let result = try await task() promise(.success(result)) } catch { promise(.failure(error)) } } } .retry(self?.maxRetryCount ?? 0) .eraseToAnyPublisher() } .sink(receiveCompletion: { [weak self] completion in guard let self else { return } if case .failure(let error) = completion { self.failureHandler?(error) // 重试全部失败后移除任务 var currentTasks = self.tasks.value currentTasks.removeFirst() self.tasks.send(currentTasks) self.remaining = currentTasks.count } }, receiveValue: { [weak self] _ in guard let self else { return } // 执行成功后移除任务 var currentTasks = self.tasks.value currentTasks.removeFirst() self.tasks.send(currentTasks) self.remaining = currentTasks.count }) .store(in: &cancellables) return self } /// 动态追加任务 func appendWorkItems(_ items: [UploadTask]) { var currentTasks = tasks.value currentTasks.append(contentsOf: items) tasks.send(currentTasks) remaining = currentTasks.count } }
可选优化点
- 上述示例默认按队列顺序移除任务,如果你需要支持任务乱序完成后正确移除,可以给每个
UploadTask增加唯一ID属性,任务完成后根据ID过滤移除对应任务即可,不要直接调用removeFirst()。 - 可以额外暴露任务进度、执行状态的Publisher,方便业务层展示整体上传进度、单个任务状态。
- 如果需要支持队列暂停/继续能力,新增一个
isSuspended的CurrentValueSubject,用filter操作符控制任务流的发射即可。
内容的提问来源于stack exchange,提问作者Simon
相关产品推荐
相关产品推荐

