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

Combine中定时任务的失败延迟重试实现方案咨询

问题描述

我需要实现按固定间隔调度API请求的逻辑:

  • 请求失败时,延迟一段时间重试,达到最大重试次数后停止调用
  • 请求成功时,继续按原定间隔执行

我的思路:

  1. 使用Timer和dataTaskPublisher两个流
  2. 由Timer驱动dataTaskPublisher
  3. 若dataTaskPublisher失败,重新调度Timer并添加延迟
  4. 达到最大重试次数后停止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")
        
       
    })
解决方案

原代码问题分析

  1. 错误被吞掉:catch捕获错误后返回了新的Publisher,导致上游错误无法传递到sink的receiveCompletion,所以错误回调永远不会触发。
  2. 重试逻辑混乱:retry应该直接作用在原始请求流上,搭配延迟逻辑实现失败后重试,而不是在catch里重新包装请求。
  3. 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)

关键逻辑说明

  1. 重试状态锁:用isRetrying标记当前是否处于重试阶段,避免Timer在重试期间触发新请求,保证同一时间只有一个请求在执行。
  2. 自定义延迟重试Operator:delayIfFailure实现了请求失败后延迟再抛出错误的逻辑,配合retry就能实现失败后延迟重试的效果。
  3. Timer生命周期管理:当重试耗尽、请求最终失败时,调用stopScheduler清空订阅并停止Timer,避免无效调度。
  4. 错误透传:没有用catch吞掉错误,让错误能传递到sink的receiveCompletion,方便处理最终失败的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 06:50:47