能否在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 } } }
代码逻辑说明
- 状态维护:使用
scan操作符维护一个[Element: Date]字典,存储每个值最后一次被发送的时间。 - 冷却判断:每次收到新元素时,获取当前调度器的时间,对比该元素的最后发送时间:
- 如果未超过冷却期,忽略该元素,返回原状态字典;
- 如果超过冷却期或首次出现,更新该元素的最后发送时间,返回新的状态字典。
- 元素提取与去重:通过
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
相关产品推荐
相关产品推荐

