如何实现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
相关产品推荐
相关产品推荐

