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

为何Rx.NET代码抛出ObjectDisposedException?SubscribeOn原理及实现疑问

Rx.NET中ObjectDisposedException问题分析与解决

咱们先从你给出的第一段代码说起,看看为什么会抛出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的调度逻辑不匹配**:

  1. Observable.Using的规则是:当内部Observable完成时,立即释放它创建的资源(这里就是EventLoopScheduler实例els)。
  2. SubscribeOn(els)会把订阅内部Observable的整个流程放到els的专属线程上执行。
  3. 当Observable.Return(1)发送完OnNext(1)和OnCompleted后,内部Observable就标记为完成,Using会立刻调用els.Dispose()。
  4. 但问题是,els的线程此时可能还在处理OnCompleted的回调逻辑,Dispose()会直接终止这个线程,导致线程上未完成的操作抛出ObjectDisposedException。

SubscribeOn的具体工作机制

很多人会把SubscribeOn和ObserveOn搞混,我来给你理清楚它的核心逻辑:

  • SubscribeOn的作用是指定「订阅原始Observable」这个操作所在的线程/调度器,同时,原始Observable发送的所有通知(OnNext/OnError/OnCompleted)都会在这个调度器的线程上执行(除非后续用ObserveOn切换)。
  • 具体执行步骤:
    1. 当你订阅source.SubscribeOn(scheduler)返回的代理Observable时,当前线程会先创建一个基础订阅对象。
    2. 代理Observable会把「订阅原始source的操作」包装成任务,提交给指定的scheduler执行。
    3. 所有来自source的通知都会通过这个调度器的线程转发给观察者。
    4. 当订阅被取消或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()会同时触发两个操作:

  1. 取消_dataService.SubscribeOn(els)的订阅。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:28:18