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

Rx限制并发数Scheduler实现问题:信号量未释放及优化咨询

有限并发Scheduler实现失效原因及优化方案

失效原因

  1. Disposable管理逻辑错误
    你实现的ThreeAtATimeScheduler.Schedule方法中,返回的CompositeDisposable将Task<IDisposable>(ContinueWith的返回值)作为参数传入,但Task并非IDisposable类型,导致任务实际执行的Disposable未被纳入CompositeDisposable的管理范围。同时,释放信号量的Disposable没有和任务的生命周期绑定,任务执行完成后没有任何机制触发该Disposable的Dispose方法,信号量因此无法释放。

  2. 任务生命周期未关联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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 01:52:17