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
相关产品推荐
相关产品推荐

