为何Rx.NET代码抛出ObjectDisposedException?SubscribeOn原理及实现疑问
咱们先从你给出的第一段代码说起,看看为什么会抛出ObjectDisposedException,再拆解SubscribeOn的工作机制,最后解决你Windows服务里的实际问题。
一、第一段代码的异常原因
先把代码贴出来方便分析:
static void Main(string[] args) { Observable.Using(() => new EventLoopScheduler(), els => Observable .Defer(() => Observable.Return(1)) .SubscribeOn(els)) .Subscribe(); Console.ReadLine(); }
这个异常的核心是**Observable.Using的资源释放时机和SubscribeOn的调度逻辑不匹配**:
Observable.Using的规则是:当内部Observable完成时,立即释放它创建的资源(这里就是EventLoopScheduler实例els)。SubscribeOn(els)会把订阅内部Observable的整个流程放到els的专属线程上执行。- 当
Observable.Return(1)发送完OnNext(1)和OnCompleted后,内部Observable就标记为完成,Using会立刻调用els.Dispose()。 - 但问题是,
els的线程此时可能还在处理OnCompleted的回调逻辑,Dispose()会直接终止这个线程,导致线程上未完成的操作抛出ObjectDisposedException。
SubscribeOn的具体工作机制
很多人会把SubscribeOn和ObserveOn搞混,我来给你理清楚它的核心逻辑:
SubscribeOn的作用是指定「订阅原始Observable」这个操作所在的线程/调度器,同时,原始Observable发送的所有通知(OnNext/OnError/OnCompleted)都会在这个调度器的线程上执行(除非后续用ObserveOn切换)。- 具体执行步骤:
- 当你订阅
source.SubscribeOn(scheduler)返回的代理Observable时,当前线程会先创建一个基础订阅对象。 - 代理Observable会把「订阅原始source的操作」包装成任务,提交给指定的
scheduler执行。 - 所有来自source的通知都会通过这个调度器的线程转发给观察者。
- 当订阅被取消或source完成时,对应的清理操作也会在调度器的线程上执行。
- 当你订阅
- 注意:如果原始Observable是热Observable(比如已经在发送数据的Subject),
SubscribeOn只会处理订阅后的通知,之前的数据不会切换线程。
二、Windows服务中停止异常的问题解决
再看你实际项目里的代码(我补全了一些缺失的变量定义,方便分析):
class DataService : ObservableBase<Unit> { private readonly HttpClient _httpClient; public DataService(HttpClient httpClient) { _httpClient = httpClient; } protected override IDisposable SubscribeCore(IObserver<Unit> o) { return Observable.Defer(() => Observable.Start(() => _httpClient.GetAsync("...").Result)) .RepeatWithDelay(TimeSpan.FromSeconds(1)) .ObserveOn(SchedulerEx.Current)//observe back to event loop .Do(_ => { /* 一些处理逻辑 */ }) .Select(_ => Unit.Default) .Subscribe(o); } } class Controller { private IDisposable _instance; private readonly DataService _dataService; public Controller(DataService dataService) { _dataService = dataService; } void Start() { _instance = Observable.Using( () => SchedulerEx.Create(), els => _dataService.SubscribeOn(els)); } void Stop() { _instance.Dispose(); } } class SchedulerEx { [ThreadStatic] public static EventLoopScheduler Current; public EventLoopScheduler Create() { var els = new EventLoopScheduler(); els.Schedule(() => SchedulerEx.Current = els); return els; } } static void Main() { var kernel = new StandardKernel();// 假设用Ninject kernel.Bind<HttpClient>().ToSelf(); kernel.Bind<DataService>().ToSelf(); kernel.Bind<Controller>().ToSelf(); var controller = kernel.Get<Controller>(); controller.Start(); controller.Stop();// 立即停止时抛出异常 }
问题根源
当你调用Stop()时,_instance.Dispose()会同时触发两个操作:
- 取消
_dataService.SubscribeOn(els)的订阅。 Observable.Using会立即释放els(调用els.Dispose())。
但此时DataService内部的RepeatWithDelay可能已经把下一次请求的任务放到了els的调度队列中,或者当前正在执行HttpClient.Get操作。els.Dispose()会直接终止线程,导致这些正在执行/等待的任务访问已释放的调度器,从而抛出ObjectDisposedException。
另外,依赖[ThreadStatic]的SchedulerEx.Current也有风险:当els被终止后,线程被回收,后续如果有代码误访问Current,可能拿到null或已释放的对象。
解决方案
我给你几个逐步优化的方案,从根源上解决问题:
1. 同步订阅生命周期与调度器生命周期
Observable.Using的资源释放时机太激进,我们需要先确保订阅完全取消、所有任务终止后,再释放调度器。可以用CompositeDisposable来手动管理生命周期:
void Start() { var els = SchedulerEx.Create(); // 先订阅,再把订阅和调度器一起加入CompositeDisposable var subscription = _dataService.SubscribeOn(els).Subscribe(); _instance = new CompositeDisposable(subscription, Disposable.Create(() => { // 先等待调度器队列清空,再释放 var tcs = new TaskCompletionSource<bool>(); els.Schedule(() => tcs.SetResult(true)); tcs.Task.Wait(); els.Dispose(); })); }
这样调用Stop()时,会先取消订阅,等待els队列中的所有任务执行完成,再优雅释放调度器,避免线程被强制终止。
2. 用CancellationToken控制任务循环
给DataService的循环逻辑加上取消令牌,让订阅取消时能优雅终止任务,而不是靠调度器终止:
class DataService : ObservableBase<Unit> { private readonly HttpClient _httpClient; public DataService(HttpClient httpClient) { _httpClient = httpClient; } protected override IDisposable SubscribeCore(IObserver<Unit> o) { var cts = new CancellationTokenSource(); var subscription = Observable.Defer(() => Observable.Start(async () => { if (cts.Token.IsCancellationRequested) throw new OperationCanceledException(cts.Token); await _httpClient.GetAsync("...", cts.Token); }, cts.Token)) .RepeatWithDelay(TimeSpan.FromSeconds(1), cts.Token)// 假设RepeatWithDelay支持取消令牌 .ObserveOn(SchedulerEx.Current) .Do(_ => { /* 处理逻辑 */ }) .Select(_ => Unit.Default) .Subscribe(o); return new CompositeDisposable(subscription, cts); } }
这样当订阅被取消时,cts会触发取消,RepeatWithDelay会停止循环,HttpClient的请求也会被取消,不会有残留任务在调度器中。
3. 避免依赖ThreadStatic变量
[ThreadStatic]的SchedulerEx.Current很容易出问题,比如线程被回收后变量不会自动清空,或者跨线程访问时拿到错误的值。建议直接把调度器传入DataService:
class DataService : ObservableBase<Unit> { private readonly HttpClient _httpClient; private readonly EventLoopScheduler _scheduler; public DataService(HttpClient httpClient, EventLoopScheduler scheduler) { _httpClient = httpClient; _scheduler = scheduler; } protected override IDisposable SubscribeCore(IObserver<Unit> o) { var cts = new CancellationTokenSource(); var subscription = Observable.Defer(() => Observable.Start(async () => { if (cts.Token.IsCancellationRequested) throw new OperationCanceledException(cts.Token); await _httpClient.GetAsync("...", cts.Token); }, cts.Token)) .RepeatWithDelay(TimeSpan.FromSeconds(1), cts.Token) .ObserveOn(_scheduler)// 直接用传入的调度器 .Do(_ => { /* 处理逻辑 */ }) .Select(_ => Unit.Default) .Subscribe(o); return new CompositeDisposable(subscription, cts); } } // Controller的Start方法修改为: void Start() { var els = SchedulerEx.Create(); var dataService = new DataService(_httpClient, els); var subscription = dataService.SubscribeOn(els).Subscribe(); _instance = new CompositeDisposable(subscription, Disposable.Create(() => { var tcs = new TaskCompletionSource<bool>(); els.Schedule(() => tcs.SetResult(true)); tcs.Task.Wait(); els.Dispose(); })); }
这样完全避免了对ThreadStatic变量的依赖,代码更可靠。
总结
不管是小例子还是实际项目,核心问题都是调度器的生命周期没有和订阅的生命周期同步:当调度器被提前释放时,正在执行的任务会访问已释放的资源,从而抛出异常。通过优雅的资源管理(先取消订阅、等待任务终止,再释放调度器)、用取消令牌控制任务执行、避免依赖ThreadStatic变量,就能彻底解决这个问题。
内容的提问来源于stack exchange,提问作者Junru Zhu

