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

Swift的Combine框架是否有类似RxSwift的sample(on:)运算符?如何实现采样?

Combine sample(on:) 运算符问题解答

关于sample操作符的行为

sample操作符的核心逻辑为:存在两个流,当其中的采样触发流被触发时,若被采样流的最新值尚未被发送过,则将该最新值发送出去,行为示意图如下:
sample操作符行为示意图


1. Combine是否原生支持对应运算符

Combine 没有提供与RxSwift、Reactive Swift中sample(on:)功能完全一致的原生运算符。

现有原生的combineLatest、zip等多流聚合运算符都无法匹配采样逻辑:

  • combineLatest:任意一个流产生新值时都会触发输出,无法做到仅由触发流控制输出时机
  • zip:要求两个流都产生新值才会输出,无法满足仅取被采样流最新值的需求

2. Combine中自定义实现sample(on:)的方案

可以通过扩展Publisher的方式实现通用的采样运算符,代码如下:

import Combine

extension Publisher {
    /// 模拟RxSwift的sample(on:)运算符
    /// - Parameter trigger: 采样触发流,每当该流发出事件时触发一次采样
    /// - Returns: 仅在采样触发时输出上游最新未发送值的Publisher
    func sample<Trigger: Publisher>(on trigger: Trigger) -> AnyPublisher<Output, Failure> 
    where Trigger.Failure == Failure {
        let latestValue = CurrentValueSubject<Output?, Failure>(nil)
        let hasUnsentNewValue = CurrentValueSubject<Bool, Failure>(false)
        
        // 监听上游输出,更新最新值和未发送标记
        let upstreamCancellable = self
            .handleEvents(
                receiveOutput: { value in
                    latestValue.send(value)
                    hasUnsentNewValue.send(true)
                },
                receiveCompletion: { completion in
                    latestValue.send(completion: completion)
                    hasUnsentNewValue.send(completion: completion)
                }
            )
            .sink(receiveCompletion: { _ in }, receiveValue: { _ in })
        
        // 监听触发流,符合条件时输出最新值
        return trigger
            .compactMap { _ in
                guard hasUnsentNewValue.value, let value = latestValue.value else {
                    return nil
                }
                hasUnsentNewValue.send(false)
                return value
            }
            .handleEvents(receiveCancel: {
                upstreamCancellable.cancel()
            })
            .eraseToAnyPublisher()
    }
}

使用示例

var cancellables = Set<AnyCancellable>()

// 被采样流:每1秒生成一个递增整数
let sampledStream = Timer.publish(every: 1, on: .main, in: .common)
    .autoconnect()
    .scan(0, { accum, _ in accum + 1 })

// 触发流:每3秒触发一次采样
let triggerStream = Timer.publish(every: 3, on: .main, in: .common)
    .autoconnect()

// 订阅采样结果,预期每3秒输出一次被采样流的最新值:3、6、9...
sampledStream
    .sample(on: triggerStream)
    .sink(receiveValue: { print("采样结果:\($0)") })
    .store(in: &cancellables)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 20:06:03