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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 20:17:50