Rx限制并发数Scheduler实现问题:信号量未释放及优化咨询
有限并发Scheduler实现失效原因及优化方案
失效原因
Disposable管理逻辑错误
你实现的ThreeAtATimeScheduler.Schedule方法中,返回的CompositeDisposable将Task<IDisposable>(ContinueWith的返回值)作为参数传入,但Task并非IDisposable类型,导致任务实际执行的Disposable未被纳入CompositeDisposable的管理范围。同时,释放信号量的Disposable没有和任务的生命周期绑定,任务执行完成后没有任何机制触发该Disposable的Dispose方法,信号量因此无法释放。任务生命周期未关联Scheduler的Disposable
Observable.Start在任务执行完毕后会结束订阅,但不会主动调用Scheduler返回的Disposable的Dispose方法。你的实现中没有将任务的完成事件与信号量释放逻辑挂钩,导致信号量的许可被永久占用,后续任务无法获取许可执行。
优化方案
方案一:正确实现自定义IScheduler
确保任务执行完成或取消时自动释放信号量,同时正确管理Disposable的生命周期:
public class ThreeAtATimeScheduler : IScheduler { private readonly SemaphoreSlim _semaphore = new SemaphoreSlim(3); private readonly IScheduler _innerScheduler = ThreadPoolScheduler.Instance; public DateTimeOffset Now => _innerScheduler.Now; public IDisposable Schedule<TState>(TState state, Func<IScheduler, TState, IDisposable> action) { var cts = new CancellationTokenSource(); var compositeDisposable = new CompositeDisposable(cts); _ = Task.Run(async () => { try { await _semaphore.WaitAsync(cts.Token); if (!cts.IsCancellationRequested) { var taskDisposable = action(_innerScheduler, state); compositeDisposable.Add(taskDisposable); } } catch (OperationCanceledException) { // 取消请求无需额外处理 } finally { if (!cts.IsCancellationRequested) { _semaphore.Release(); } } }, cts.Token); compositeDisposable.Add(Disposable.Create(() => { cts.Cancel(); try { cts.Token.WaitHandle.WaitOne(); } catch { } })); return compositeDisposable; } public IDisposable Schedule<TState>(TState state, TimeSpan dueTime, Func<IScheduler, TState, IDisposable> action) { return _innerScheduler.Schedule(state, dueTime, (s, st) => Schedule(st, action)); } }
方案二:使用Rx原生操作符替代自定义Scheduler
无需手动实现Scheduler,直接利用Rx的Merge操作符指定最大并发数,代码更简洁且稳定性更高:
var subscription = Observable.Range(0, 100) .Buffer(TimeSpan.FromSeconds(1), 10) .SelectMany(x => x) .Select(x => Observable.Start(() => { Console.WriteLine($"Performing for Thread {Thread.CurrentThread.ManagedThreadId}"); })) .Merge(maxConcurrent: 3) // 限制最大并发数为3 .Subscribe();
该方案由Rx自动管理任务的并发与Disposable生命周期,任务完成后自动释放并发槽,避免手动实现的潜在错误。
内容的提问来源于stack exchange,提问作者fahadash
相关产品推荐
相关产品推荐

