Combine中定时任务的失败延迟重试实现方案咨询
问题描述
我需要实现按固定间隔调度API请求的逻辑:
- 请求失败时,延迟一段时间重试,达到最大重试次数后停止调用
- 请求成功时,继续按原定间隔执行
我的思路:
- 使用Timer和dataTaskPublisher两个流
- 由Timer驱动dataTaskPublisher
- 若dataTaskPublisher失败,重新调度Timer并添加延迟
- 达到最大重试次数后停止Timer(从而停止dataTaskPublisher)
目前不知道如何正确编排这两个流,附上尝试的代码(其中错误回调未触发),求示例方案参考:
var cancellables: AnyCancellable let interval = 2 let maxRetry = 2 let delay = 5 var timer = Timer.publish(every: TimeInterval(interval), on: .main, in: .common).autoconnect() var session = URLSession.shared.dataTaskPublisher(for: URL(string: "https://google.com")!).map(\.data).eraseToAnyPublisher() cancellables = timer.flatMap{ _ in return session.catch { _ -> AnyPublisher<Data, URLError> in return Publishers.Delay(upstream: session, interval: .seconds(delay), tolerance: 1, scheduler: RunLoop.main).print("retrying").retry(maxRetry).eraseToAnyPublisher() }.eraseToAnyPublisher() }.sink(receiveCompletion: { comp in if case let .failure(err) = comp { print("=============================This is error \(err)") // This doesn't trigger. } }, receiveValue: { data in print("success") })
解决方案
原代码问题分析
- 错误被吞掉:
catch捕获错误后返回了新的Publisher,导致上游错误无法传递到sink的receiveCompletion,所以错误回调永远不会触发。 - 重试逻辑混乱:
retry应该直接作用在原始请求流上,搭配延迟逻辑实现失败后重试,而不是在catch里重新包装请求。 - Timer未受控:重试耗尽后没有停止Timer,导致Timer会继续触发新的无效请求。
正确实现示例
以下是符合需求的Combine编排方案,核心是把重试逻辑封装到请求流中,同时在重试耗尽时停止Timer:
import Combine import Foundation class APIScheduler { private var cancellables = Set<AnyCancellable>() private let baseInterval: TimeInterval private let maxRetryCount: Int private let retryDelay: TimeInterval private let apiURL: URL private var timer: AnyPublisher<Void, Never>? private var isRetrying = false init(baseInterval: TimeInterval = 2, maxRetryCount: Int = 2, retryDelay: TimeInterval = 5, apiURL: URL) { self.baseInterval = baseInterval self.maxRetryCount = maxRetryCount self.retryDelay = retryDelay self.apiURL = apiURL startScheduler() } private func startScheduler() { // 创建Timer流,重试期间不触发新请求 timer = Timer.publish(every: baseInterval, on: .main, in: .common) .autoconnect() .filter { [weak self] _ in guard let self = self else { return false } return !self.isRetrying } .map { _ in } .share() .eraseToAnyPublisher() timer? .flatMap { [weak self] _ in guard let self = self else { return Empty().eraseToAnyPublisher() } return self.makeAPIRequest() } .sink(receiveCompletion: { [weak self] completion in if case .failure(let error) = completion { print("请求最终失败: \(error)") // 重试耗尽后停止调度器 self?.stopScheduler() } }, receiveValue: { data in print("请求成功,收到数据长度: \(data.count)") }) .store(in: &cancellables) } private func makeAPIRequest() -> AnyPublisher<Data, URLError> { isRetrying = true return URLSession.shared.dataTaskPublisher(for: apiURL) .map(\.data) // 失败时添加延迟再抛出错误 .delayIfFailure(interval: retryDelay, scheduler: RunLoop.main) .retry(maxRetryCount) .handleEvents(receiveCompletion: { [weak self] _ in // 请求完成后结束重试状态 self?.isRetrying = false }) .eraseToAnyPublisher() } private func stopScheduler() { timer = nil cancellables.removeAll() print("调度器已停止") } } // 自定义Operator:实现失败后延迟抛出错误 extension Publisher where Failure: Error { func delayIfFailure(interval: TimeInterval, scheduler: some Scheduler) -> AnyPublisher<Output, Failure> { self.catch { error in return Just(error) .delay(for: .seconds(interval), scheduler: scheduler) .flatMap { Fail(error: $0).eraseToAnyPublisher() } } .eraseToAnyPublisher() } } // 使用示例 guard let url = URL(string: "https://google.com") else { fatalError("无效URL") } let scheduler = APIScheduler(apiURL: url)
关键逻辑说明
- 重试状态锁:用
isRetrying标记当前是否处于重试阶段,避免Timer在重试期间触发新请求,保证同一时间只有一个请求在执行。 - 自定义延迟重试Operator:
delayIfFailure实现了请求失败后延迟再抛出错误的逻辑,配合retry就能实现失败后延迟重试的效果。 - Timer生命周期管理:当重试耗尽、请求最终失败时,调用
stopScheduler清空订阅并停止Timer,避免无效调度。 - 错误透传:没有用
catch吞掉错误,让错误能传递到sink的receiveCompletion,方便处理最终失败的场景。
内容的提问来源于stack exchange,提问作者Bishow Gurung
相关产品推荐
相关产品推荐

