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

Swift Combine带延迟的Repeat If自定义Publisher实现问题求助

Swift Combine递归flatMap实现RepeatIf 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 11:27:27