如何为手工实现的Observable返回可取消订阅的IDisposable?
手工实现支持取消订阅的
IObservable 首先得给你的思路点个赞——从学术角度探究手动实现IObservable确实能帮你更深入理解Rx的底层逻辑。你提到的初始实现(返回Disposable.Empty)的核心问题在于它是同步阻塞执行的,整个循环跑完才会退出Subscribe方法,这时候返回空Disposable自然毫无意义,因为任务已经完成了。
你的第二个实现用Task.Run把工作移到后台线程,还通过自定义SubscriptionHandler管理观察者列表,方向完全正确,但还有几个关键细节需要修正:
现有实现的核心问题
- 线程安全隐患:
List<T>不是线程安全集合,当多个线程同时调用Subscribe/Dispose时,Add/Remove操作可能引发异常或数据不一致。 - 取消逻辑不完整:即使把观察者从列表中移除,后台线程的循环还会继续调用
observer.OnNext,这会导致已取消订阅的观察者收到无效消息,甚至可能触发空引用异常。 - 异常处理缺失:如果观察者的
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
相关产品推荐
相关产品推荐

