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

Reactive Extensions:处理前项时新事件未被处理的问题排查

Rx BufferUntilInactive 事件滞留未处理问题

我使用Reactive Extensions(Rx)实现最多X个事件为一组或无新事件Y毫秒后分组处理事件的逻辑,采用BufferUntilInactive扩展方法实现。原本运行正常,但发现当队列仍在处理前一批事件时发送新事件,部分事件会滞留未处理;只有在处理完成后发送新事件,队列中包括之前滞留的所有事件才会被处理。

测试代码与复现场景

为复现该行为,编写测试代码向subject中添加数字,等待数字被处理或10秒超时;若数字处理完成则添加下一个数字,第二个数字会在前一个数字处理仍在进行时发送,从而触发问题(无延迟时因时机问题复现概率低)。

测试代码:

private DispatcherScheduler _schedulerProviderDispatcher = new DispatcherScheduler(Application.Current.Dispatcher);
private Subject<int> _numberEvents = new Subject<int>();
private readonly ConcurrentDictionary<int, object> _numbersInProgress = new ConcurrentDictionary<int, object>();
private IDisposable _disposable;

public async Task SendNumbersAsync()
{
    _disposable = _numberEvents.Synchronize()
        .BufferUntilInactive(TimeSpan.FromMilliseconds(2000), _schedulerProviderDispatcher, 100)
        .Subscribe(numbers =>
        {
            if (numbers.Count == 0)
                return;

            var values = string.Empty;

            // 处理队列中的数字并从字典移除
            for (var i = 0; i < numbers.Count; ++i)
            {
                var number = numbers[i];
                if (_numbersInProgress.TryRemove(number, out _) == false)
                    Trace.WriteLine($"Failed to remove number '{number}'");

                values += $"{number}, ";
                Trace.WriteLine($"Handled Number: {number}, Count: {i + 1}/{numbers.Count}");
            }

            Trace.WriteLine($"Handled '{numbers.Count}' numbers. Values: '{values}'");

            // 模拟处理延迟,制造前一批未处理完就收到新事件的场景
            var task = Task.Delay(1000);
            Task.WaitAll(task);

            Trace.WriteLine($"Finished handling '{numbers.Count}' numbers. Values: '{values}'");
        });

    // 向subject推送数字
    var number = 0;
    while (_numbersInProgress.Count == 0)
    {
        ++number;
        Trace.WriteLine($"Sending number: {number}");
        _numberEvents.OnNext(number);
        Trace.WriteLine($"Finished sending number: {number}");

        // 将数字加入处理中的字典
        if (_numbersInProgress.TryAdd(number, 0) == false)
            Trace.WriteLine($"Failed to add number '{number}'");

        var waitCount = 0;
        var timedOut = false;
        // 等待数字处理完成,超时10秒则停止
        while (_numbersInProgress.Count != 0)
        {
            Trace.WriteLine($"Waiting ({waitCount})");
            await Task.Delay(100);
            ++waitCount;

            if (waitCount > 100)
            {
                timedOut = true;
                break;
            }
        }

        if (timedOut)
            Trace.WriteLine($"Timeout waiting: {number}");
        else
            Trace.WriteLine($"Finished waiting: {number}");
    }

    // 若仍有未处理数字,发送一个新数字触发处理
    if (_numbersInProgress.Count != 0)
    {
        Trace.WriteLine($"Failed waiting: {number}");

        Trace.WriteLine($"Sending number: {9999}");
        _numberEvents.OnNext(9999);
        Trace.WriteLine($"Finished sending number: {9999}");
    }
}

分组扩展方法BufferUntilInactive

public static IObservable<IList<T>> BufferUntilInactive<T>(this IObservable<T> stream, TimeSpan delay, IScheduler scheduler = null, Int32? maxCount = null)
{
     var s = scheduler ?? Scheduler.Default;

     var publish = stream.Publish(p =>
     {
          var closes = p.Throttle(delay, s);

          if (maxCount != null)
          {
              var overflows = p.Where((x, index) => index + 1 >= maxCount);
              closes = closes.Amb(overflows);
          }

          return p.Window(() => closes).SelectMany(window => window.ToList());
     });

     return publish;
}

问题现象与疑问

执行测试代码时,跟踪输出显示:数字2在数字1处理完成前发送,导致等待超时;发送数字9999后,队列中滞留的数字才被处理。

我已尝试将Throttle和事件分组逻辑封装在Publish中,但未解决问题。请问BufferUntilInactive方法或其他环节是否存在遗漏?

内容的提问来源于Stack Exchange,提问作者Andi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 14:40:30