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

如何为手工实现的Observable返回可取消订阅的IDisposable?

手工实现支持取消订阅的IObservable

首先得给你的思路点个赞——从学术角度探究手动实现IObservable确实能帮你更深入理解Rx的底层逻辑。你提到的初始实现(返回Disposable.Empty)的核心问题在于它是同步阻塞执行的,整个循环跑完才会退出Subscribe方法,这时候返回空Disposable自然毫无意义,因为任务已经完成了。

你的第二个实现用Task.Run把工作移到后台线程,还通过自定义SubscriptionHandler管理观察者列表,方向完全正确,但还有几个关键细节需要修正:


现有实现的核心问题

  1. 线程安全隐患:List<T>不是线程安全集合,当多个线程同时调用Subscribe/Dispose时,Add/Remove操作可能引发异常或数据不一致。
  2. 取消逻辑不完整:即使把观察者从列表中移除,后台线程的循环还会继续调用observer.OnNext,这会导致已取消订阅的观察者收到无效消息,甚至可能触发空引用异常。
  3. 异常处理缺失:如果观察者的OnNext抛出异常,后台线程会直接崩溃,没有正确触发OnError通知。

正确的实现方式

我们需要解决线程安全、任务取消、生命周期管理这三个核心点,下面给你两种可行的思路:

思路1:线程安全集合 + 取消令牌

用ConcurrentDictionary存储观察者和对应的取消令牌,同时在后台任务中检查取消状态,确保及时终止执行:

public class MyObservable : IObservable<int>
{
    // 线程安全的字典:观察者 -> 对应的取消令牌
    private readonly ConcurrentDictionary<IObserver<int>, CancellationTokenSource> _observers = new();

    public IDisposable Subscribe(IObserver<int> observer)
    {
        if (observer == null) throw new ArgumentNullException(nameof(observer));
        
        var cts = new CancellationTokenSource();
        if (_observers.TryAdd(observer, cts))
        {
            // 启动后台任务发送数据
            _ = Task.Run(async () =>
            {
                try
                {
                    for (int i = 0; i < 5; i++)
                    {
                        // 检查是否已取消订阅
                        if (cts.Token.IsCancellationRequested)
                        {
                            break;
                        }
                        
                        // 用Task.Delay替代Thread.Sleep,支持取消
                        await Task.Delay(1000, cts.Token);
                        observer.OnNext(i);
                    }
                    
                    // 未取消的情况下发送完成通知
                    if (!cts.Token.IsCancellationRequested)
                    {
                        observer.OnCompleted();
                    }
                }
                catch (OperationCanceledException)
                {
                    // 取消订阅时无需额外处理,直接退出
                }
                catch (Exception ex)
                {
                    // 异常时通知观察者(排除取消异常)
                    if (!cts.Token.IsCancellationRequested)
                    {
                        observer.OnError(ex);
                    }
                }
                finally
                {
                    // 任务结束后自动清理观察者
                    _observers.TryRemove(observer, out _);
                    cts.Dispose();
                }
            }, cts.Token);
        }

        // 返回自定义Disposable,负责取消任务并移除观察者
        return new SubscriptionHandler(_observers, observer);
    }

    private class SubscriptionHandler : IDisposable
    {
        private readonly ConcurrentDictionary<IObserver<int>, CancellationTokenSource> _observers;
        private readonly IObserver<int> _observer;
        private bool _disposed;

        public SubscriptionHandler(ConcurrentDictionary<IObserver<int>, CancellationTokenSource> observers, IObserver<int> observer)
        {
            _observers = observers;
            _observer = observer;
        }

        public void Dispose()
        {
            if (_disposed) return;
            
            // 取消对应任务并清理资源
            if (_observers.TryRemove(_observer, out var tokenSource))
            {
                tokenSource.Cancel();
                tokenSource.Dispose();
            }
            
            _disposed = true;
        }
    }
}

思路2:复用Rx内置工具简化实现

Rx本身提供了很多现成的工具类,不需要完全从零写SubscriptionHandler。比如Disposable.Create(Action)可以快速创建一个IDisposable,传入的Action会在Dispose时执行:

public class MyObservable : IObservable<int>
{
    private readonly ConcurrentDictionary<IObserver<int>, CancellationTokenSource> _observers = new();

    public IDisposable Subscribe(IObserver<int> observer)
    {
        if (observer == null) throw new ArgumentNullException(nameof(observer));
        
        var cts = new CancellationTokenSource();
        if (!_observers.TryAdd(observer, cts))
        {
            // 观察者已存在,返回空Disposable
            return Disposable.Empty;
        }

        // 用Disposable.Create快速生成取消逻辑
        var subscription = Disposable.Create(() =>
        {
            if (_observers.TryRemove(observer, out var tokenSource))
            {
                tokenSource.Cancel();
                tokenSource.Dispose();
            }
        });

        // 后台任务逻辑和思路1一致
        _ = Task.Run(async () =>
        {
            try
            {
                for (int i = 0; i < 5; i++)
                {
                    if (cts.Token.IsCancellationRequested) break;
                    await Task.Delay(1000, cts.Token);
                    observer.OnNext(i);
                }
                
                if (!cts.Token.IsCancellationRequested)
                {
                    observer.OnCompleted();
                }
            }
            catch (Exception ex) when (!(ex is OperationCanceledException))
            {
                observer.OnError(ex);
            }
            finally
            {
                // 任务结束后自动清理订阅
                subscription.Dispose();
            }
        });

        return subscription;
    }
}

关于内置方式的疑问

你问的「是否存在内置方式可为手工实现的Observable返回有效的Disposable」——答案是肯定的:

  • Disposable.Create(Action):Rx最常用的工具之一,可以快速封装取消逻辑,无需自定义IDisposable类。
  • CompositeDisposable:可以组合多个IDisposable,一次性全部清理,适合需要管理多个资源的场景。
  • SerialDisposable:可以动态替换内部的Disposable,确保旧资源被正确释放。

而Observable.Create本质上就是封装了这些最佳实践:它帮你处理订阅管理、取消逻辑、线程安全,返回的Disposable会自动清理所有相关资源。如果不是为了学习底层,直接用Observable.Create会更简洁可靠:

public IObservable<int> CreateMyObservable()
{
    return Observable.Create<int>(observer =>
    {
        var cts = new CancellationTokenSource();
        
        var task = Task.Run(async () =>
        {
            try
            {
                for (int i = 0; i < 5; i++)
                {
                    if (cts.Token.IsCancellationRequested) break;
                    await Task.Delay(1000, cts.Token);
                    observer.OnNext(i);
                }
                
                if (!cts.Token.IsCancellationRequested)
                {
                    observer.OnCompleted();
                }
            }
            catch (Exception ex) when (!(ex is OperationCanceledException))
            {
                observer.OnError(ex);
            }
        });
        
        // 返回Disposable,负责取消任务并清理资源
        return Disposable.Create(() =>
        {
            cts.Cancel();
            try { task.Wait(); } catch { /* 忽略取消异常 */ }
            cts.Dispose();
        });
    });
}

总结

  • 你的初始思路是对的,但要注意线程安全和取消逻辑的完整性——不仅要移除观察者,还要终止后台任务,避免无效的消息发送。
  • Rx提供了丰富的内置工具类,可以帮你简化手动实现,不用重复造轮子。
  • Observable.Create是Rx官方推荐的自定义Observable方式,它封装了所有最佳实践,实际开发中优先使用它。

内容的提问来源于stack exchange,提问作者NotAgain

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 06:42:12