如何在Rx.NET(C#)中实现Observable流的Debounce/DebounceDistinct扩展方法?
Rx.NET自定义Debounce扩展方法的实现与优化需求
正如标题所述,我找不到现成的Debounce()扩展方法(类似开箱即用的Throttle())。以下是基于我对Rx.NET理念的理解,一个粗略且可能存在缺陷的实现思路。我清楚这个实现大概率非线程安全,还把防抖与去重逻辑混在了一起,并不理想,需要更可靠的方案来提升它的健壮性:
/// <summary> /// 仅当在指定的debounceTime时间段内没有其他元素发射时,才从源可观察序列中发射元素。 /// </summary> /// <typeparam name="TSource">源序列中元素的类型。</typeparam> /// <param name="source">要执行防抖操作的源序列。</param> /// <param name="debounceTime">每个元素的防抖时长。</param> /// <param name="customComparer">用于比较当前元素与前一个元素的自定义比较器</param> /// <returns>经过防抖处理后的序列。</returns> /// <exception cref="ArgumentNullException"><paramref name="source"/> 为null时抛出。</exception> /// <exception cref="ArgumentOutOfRangeException"><paramref name="debounceTime"/> 小于TimeSpan.Zero时抛出。</exception> /// <remarks> /// <para> /// 该操作符通过为每个元素保留指定的<paramref name="debounceTime"/>时长来防抖源序列。如果在这个时间窗口内产生了另一个元素,前一个元素会被丢弃,并为当前元素启动一个新的计时器以尝试防抖, /// 从而重启整个流程。对于元素之间间隔从未大于或等于<paramref name="debounceTime"/>的流,最终的流不会产生任何元素。如果想要在减少流数据量的同时保证周期性产生元素,可以考虑使用Observable.Sample系列操作符。 /// </para> /// </remarks> static public IObservable<TSource> DebounceX<TSource>(this IObservable<TSource> source, TimeSpan debounceTime, IEqualityComparer<TSource> customComparer = null) { if (source == null) throw new ArgumentNullException(nameof(source)); if (debounceTime.TotalMilliseconds <= 0) throw new ArgumentOutOfRangeException(nameof(debounceTime)); var debouncedElementsStream = Observable.Create<TSource>(Subscribe); return debouncedElementsStream; IDisposable Subscribe(IObserver<TSource> observer) { // 稳定元素的子流 var timer = (Timer)null; var previousElement = default(TSource); var comparer = customComparer ?? EqualityComparer<TSource>.Default; return source.Subscribe(onNext: OnEachElementEmission_, onError: OnError_, onCompleted: OnCompleted_); void OnEachElementEmission_(TSource value) { if (timer != null && comparer.Equals(previousElement, value)) return; previousElement = value; timer?.Dispose(); timer = new Timer( state: null, period: Timeout.InfiniteTimeSpan, // 我们只希望计时器触发一次 dueTime: debounceTime, // 在防抖时间过后触发 callback: _ => OnElementHasRemainedStableForTheSpecificPeriod(value) ); } void OnError_(Exception exception) { timer?.Dispose(); timer = null; observer.OnError(exception); } void OnCompleted_() { timer?.Dispose(); timer = null; observer.OnCompleted(); } void OnElementHasRemainedStableForTheSpecificPeriod(TSource value) { timer?.Dispose(); timer = null; observer.OnNext(value); // 发射防抖后的元素 // 调用OnNext(value)本质是在流中发射那个在debounceTime指定的整个时间段内保持稳定的防抖元素 } } }
内容的提问来源于stack exchange,提问作者XDS
相关产品推荐
相关产品推荐

