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

能否在RxSwift中实现distinctUntilChanged与throttle的混合机制?

在RxSwift中实现自定义值级别的冷却过滤操作符

完全可以实现你描述的需求——为每个独立的值设置冷却期(比如5秒),同一值在冷却期内仅发送一次,冷却期从该值首次出现时开始计时,超时后允许再次发送;不同值之间互不影响。

为什么不能直接组合现有操作符?

distinctUntilChanged()只对比前一个元素,无法处理间隔出现的重复值(比如序列中的B、D);而throttle()是针对整个序列的时间窗口,不是针对单个值的独立冷却,所以两者的组合无法满足你的需求。

实现方案:自定义操作符

我们可以通过RxSwift的scan操作符维护一个状态字典,记录每个值的最后发送时间,结合时间判断来过滤元素。以下是具体实现:

import RxSwift

extension ObservableType where Element: Equatable {
    func distinctWithCooldown(_ dueTime: RxTimeInterval, scheduler: SchedulerType) -> Observable<Element> {
        return self.scan([Element: Date]()) { lastSendTimes, element in
            let now = scheduler.now
            // 检查当前元素是否在冷却期内
            if let lastSendTime = lastSendTimes[element], now.timeIntervalSince(lastSendTime) < dueTime.timeInterval {
                // 冷却期内,不更新时间,返回原字典
                return lastSendTimes
            } else {
                // 冷却期已过或首次出现,更新该元素的最后发送时间
                var updatedTimes = lastSendTimes
                updatedTimes[element] = now
                return updatedTimes
            }
        }
        // 提取每次更新后的元素(只有当元素被允许发送时,字典才会更新)
        .withLatestFrom(self) { $1 }
        // 过滤掉重复的元素(因为scan可能多次返回相同字典,对应不允许发送的元素)
        .distinctUntilChanged()
    }
}

// 辅助计算属性,将RxTimeInterval转为时间间隔秒数
private extension RxTimeInterval {
    var timeInterval: TimeInterval {
        switch self {
        case .seconds(let s): return TimeInterval(s)
        case .milliseconds(let ms): return TimeInterval(ms) / 1000
        case .microseconds(let us): return TimeInterval(us) / 1_000_000
        case .nanoseconds(let ns): return TimeInterval(ns) / 1_000_000_000
        case .minutes(let m): return TimeInterval(m) * 60
        case .hours(let h): return TimeInterval(h) * 3600
        case .days(let d): return TimeInterval(d) * 86400
        @unknown default: return 0
        }
    }
}

代码逻辑说明

  1. 状态维护:使用scan操作符维护一个[Element: Date]字典,存储每个值最后一次被发送的时间。
  2. 冷却判断:每次收到新元素时,获取当前调度器的时间,对比该元素的最后发送时间:
    • 如果未超过冷却期,忽略该元素,返回原状态字典;
    • 如果超过冷却期或首次出现,更新该元素的最后发送时间,返回新的状态字典。
  3. 元素提取与去重:通过withLatestFrom获取当前元素,再用distinctUntilChanged()过滤掉因状态未更新而重复的元素(确保只有符合条件的元素被发送)。

使用示例

let disposeBag = DisposeBag()
let elements = ["A", "B", "C", "D", "B", "D", "A", "X", "C", "Z"].enumerated()
    .map { index, element in
        return Observable.just(element).delay(.seconds(index), scheduler: MainScheduler.instance)
    }
    .merge()

elements
    .distinctWithCooldown(.seconds(5), scheduler: MainScheduler.instance)
    .subscribe(onNext: { print("Received: \($0)") })
    .disposed(by: disposeBag)

假设每个元素间隔1秒发送,输出结果为:

Received: A
Received: B
Received: C
Received: D
Received: A(第6秒,距离首次A已过5秒)
Received: X
Received: C(第8秒,距离首次C已过5秒)
Received: Z

完全符合你描述的需求:同一值在5秒内仅发送一次,超时后允许再次传递,不同值互不干扰。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 08:25:15