Swift Combine带延迟的Repeat If自定义Publisher实现问题求助
问题背景
我需要为Swift Combine实现一个带延迟的Repeat If自定义Publisher,用于根据前一次响应设置的延迟时间轮询后端接口(最长5分钟,轮询周期3-6秒)。尝试采用递归flatMap的方式实现,但该方案不稳定:重复轮询次数随机(20至200次不等)后,首个订阅会触发finished事件,后续所有订阅也随之终止,推测该现象可能与内存情况相关。我已通过另一种方案实现了预期功能,但仍希望了解递归flatMap方案失效的原因。
原递归实现代码
enum CustomPublishers { } extension CustomPublishers { struct RepeatIf<Upstream: Publisher>: Publisher { typealias Output = Upstream.Output typealias Failure = Upstream.Failure init( upstream: Upstream, shouldRepeat: @escaping (Upstream.Output) -> Bool, withDelay: @escaping (Upstream.Output) -> Int ) { self.upstream = upstream self.shouldRepeat = shouldRepeat self.withDelay = withDelay } var upstream: Upstream var shouldRepeat: (Upstream.Output) -> Bool var withDelay: (Upstream.Output) -> Int func receive<Downstream: Subscriber>(subscriber: Downstream) where Failure == Downstream.Failure, Output == Downstream.Input { upstream .print("CustomPublishers.RepeatIf(1)>") .flatMap { output in Just((output)).setFailureType(to: Downstream.Failure.self) .delay(for: .seconds(withDelay(output)), scheduler: DispatchQueue.global()) } .flatMap { output in shouldRepeat(output) ? Self(upstream: upstream, shouldRepeat: shouldRepeat, withDelay: self.withDelay) .eraseToAnyPublisher() : Just((output)).setFailureType(to: Downstream.Failure.self) .eraseToAnyPublisher() } .catch { (error: Upstream.Failure) -> AnyPublisher<Output, Failure> in return Fail(error: error).eraseToAnyPublisher() } .print("CustomPublishers.RepeatIf(2)>") .receive(subscriber: subscriber) } } }
可行实现代码
func retryRequestWithDelay(url: URL) -> AnyPublisher<Response, AppError> { let pollPublisher = CurrentValueSubject<Int, AppError>(0) return pollPublisher.compactMap { [weak self] delay in return self?.networkingService.request(url) .delay(for: .seconds(delay), scheduler: DispatchQueue.global()) } .switchToLatest() .receive(on: DispatchQueue.main) .handleEvents(receiveOutput: { [weak self] (response: Response) in guard let self = self else { return } if response.status == .pending { self.pollPublisher.send(response.pollPeriod) } else { self.pollPublisher.send(completion: .finished) } }) .filter { (response: Response) in guard response.data.status == .pending else { return true } return false } .eraseToAnyPublisher() }
递归flatMap方案失效的核心原因
1. Upstream Publisher的复用与状态污染
递归创建RepeatIf实例时,直接复用了原始的upstream实例。如果upstream是带有内部状态的Publisher(比如封装了网络请求的Publisher),多次订阅同一个实例会导致状态混乱。即使是冷Publisher,若其内部持有资源(如URLSessionTask),重复订阅可能触发资源的意外释放或重置,导致Publisher提前发送finished事件。
2. 递归订阅链的内存累积与意外释放
每次递归都会创建新的RepeatIf实例,这些实例会持有原始的upstream、shouldRepeat和withDelay闭包。随着轮询次数增加,内存中会形成嵌套的订阅链,每个环节都持有上层引用。当内存压力增大时,ARC可能意外释放链中的某个订阅环节,导致整个订阅链收到finished信号,终止后续所有轮询。
3. 后台线程调度的线程安全问题
延迟操作使用DispatchQueue.global()作为调度器,Combine的订阅状态管理在多线程环境下容易出现线程安全问题。多个延迟操作在不同后台队列执行时,可能导致订阅状态不一致,比如某个环节的订阅状态被错误标记为已完成,进而触发整个链的终止。
4. FlatMap逻辑的订阅生命周期缺陷
第一个flatMap将输出延迟后,第二个flatMap根据shouldRepeat决定递归或返回Just。但递归创建的RepeatIf实例会重新订阅同一个upstream,如果upstream在某次订阅中因内部逻辑(如网络超时、重试机制触发)提前完成,就会导致整个递归链终止,而这种触发时机具有随机性,表现为轮询次数不稳定。
可行方案的优势
可行方案通过CurrentValueSubject作为轮询触发源,使用switchToLatest确保每次只保留最新的订阅,避免了递归带来的嵌套订阅链问题。同时通过手动控制Subject的发送事件,明确管理轮询的启动、延迟和终止,订阅生命周期更可控,内存管理也更清晰。
内容的提问来源于stack exchange,提问作者Lukas Andrlik

