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

如何实现Combine的reset操作符?设定时长后自动发送nil至下游

Combine 自定义定时重置操作符实现需求

Combine 本身没有现成的reset操作符可以直接实现「设定时长后自动发送nil重置」的需求,但我们可以通过组合现有操作符自定义一个符合要求的版本。

自定义reset操作符代码

import Combine

extension Publisher where Output: Sendable, Failure == Never {
    func reset(after interval: TimeInterval, scheduler: some Scheduler) -> AnyPublisher<Output?, Never> {
        self
            .map { value in
                // 为每个上游元素创建一个新流:先发送当前值,延迟指定时长后发送nil
                Just(value)
                    .merge(with:
                        Just(nil)
                            .delay(for: .seconds(interval), scheduler: scheduler)
                    )
            }
            .switchLatest() // 新元素到来时切换到最新流,自动取消旧流的延迟任务
            .eraseToAnyPublisher()
    }
}

逻辑说明

  • map将每个上游的具体值转换为一个临时流:这个流会先立刻发送当前值,再通过delay在指定时长后发送nil。
  • switchLatest是核心:每当上游有新元素发送时,它会自动取消之前订阅的临时流,确保只有最后一次元素对应的延迟任务会执行——也就是说,只有当设定时长内没有新元素时,才会触发nil的发送。
  • 最终输出的Output?类型,完美匹配示例中接收Int?的需求。

测试代码及输出

let pub = PassthroughSubject<Int, Never>()
var cancellables = Set<AnyCancellable>()

pub
    .reset(after: 5, scheduler: RunLoop.main)
    .sink { v in
        print("Received \(v ?? "nil")")
    }
    .store(in: &cancellables)

pub.send(1)
pub.send(2)
Thread.sleep(forTimeInterval: 6.0) // 仅示例用,实际开发避免阻塞主线程
pub.send(3)

执行后会输出预期结果:

Received 1
Received 2
Received nil
Received 3

注意事项

  • 示例中的Thread.sleep仅用于模拟延迟场景,实际开发中应该使用调度器的延迟API(如DispatchQueue.main.asyncAfter)来避免阻塞主线程。
  • 上述操作符仅支持Failure == Never的Publisher,如果需要适配带错误类型的流,可以修改泛型约束并添加错误处理逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 00:10:47