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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 10:06:04