Rx.NET中如何正确处置Observable.Create创建的可释放资源?
Rx.NET 中正确处置Observable.Create发射的IDisposable资源
我在项目中使用Rx(.NET)时,遇到了一个关键问题:如何正确处置在Observable.Create()期间创建、并通过OnNext()发射的IDisposable资源。这里的RunData包含数据库事务和Autofac LifetimeScope,必须在订阅者完成数据处理后释放。
现有实现代码
无限序列生成逻辑
var obs = Observable.Create<RunData>(async (o) => { if (someCondition) { RunData runData = await CreateRunData(); // RunData是IDisposable,需要释放 o.OnNext(runData); } o.OnCompleted(); return Disposable.Empty; }) .Concat(Observable.Empty<RunData>().Delay(TimeSpan.FromSeconds(2))) .Repeat() // 序列完成后无限重订阅 .Publish().RefCount(); // 共享订阅
序列转换逻辑
var final = obs .Select(runData => /* 业务转换逻辑 */) // 其他操作符处理 .Select(tuple => (tuple.runData, tuple.result));
订阅处置逻辑
final.Subscribe( async (tuple) => { var (runData, result) = tuple; try { // 使用runData和result执行异步业务逻辑 } catch (Exception e) { // 错误处理 } finally { // 手动释放runData runData.Dispose(); } }, (Exception e) => { // 错误处理 });
现有实现的问题
- 泄漏风险:若序列中途抛出异常、或多订阅者场景下,
RunData可能无法被正确处置 - 职责错位:让订阅者负责资源处置不符合单一职责原则,增加维护成本
- 尝试过的方案局限性:
Observable.Using():仅在序列结束时释放资源,但我的序列是无限的(需要用Scan()等操作构建中间状态),无法适用Observable.Create()返回的Disposable回调:会在Create逻辑完成后立即触发,而非订阅者处理完数据后,引发竞态条件
正确实现方案
核心思路
让资源生命周期与单个数据项的处理周期绑定,将资源释放逻辑封装在序列内部,订阅者仅负责业务处理,无需关心资源管理。
方案1:用信号跟踪处理完成状态(适合异步处理场景)
修改序列生成逻辑,为每个RunData添加处理完成信号,确保资源在所有订阅者处理完成后释放:
var obs = Observable.Create<(RunData, IObserver<Unit>)>(async (observer, cancellationToken) => { var disposables = new CompositeDisposable(); // 构建2秒周期的任务流 var periodicStream = Observable.Interval(TimeSpan.FromSeconds(2)) .SelectMany(async _ => { if (someCondition) { var runData = await CreateRunData(); // 创建异步信号,跟踪处理完成状态 var processedSignal = new AsyncSubject<Unit>(); // 发射资源+处理信号 observer.OnNext((runData, processedSignal.AsObserver())); // 等待处理完成后释放资源 await processedSignal.FirstAsync(cancellationToken); runData.Dispose(); } return Unit.Default; }) .Subscribe(_ => { }, observer.OnError, observer.OnCompleted); disposables.Add(periodicStream); return disposables; }) // 保留信号的同时执行业务转换 .Select(item => { var (runData, signal) = item; // 你的业务转换逻辑,生成result var result = /* ... */; return (runData, result, signal); }) .Publish().RefCount();
订阅时,处理完成后触发信号:
final.Subscribe( async tuple => { var (runData, result, processedSignal) = tuple; try { // 执行异步业务逻辑 await ProcessResult(runData, result); } catch (Exception e) { // 错误处理 } finally { // 触发处理完成信号,通知序列释放资源 processedSignal.OnNext(Unit.Default); processedSignal.OnCompleted(); } }, ex => /* 错误处理 */ );
方案2:Using+Defer+Repeat(适合同步/短周期处理场景)
如果每个RunData仅需单次处理,可将Observable.Using与Defer、Repeat结合,让每个周期创建独立的资源上下文,确保资源在单次数据发射后(处理完成前)不被释放:
var obs = Observable.Defer(() => Observable.Using( // 每个周期创建新资源 async () => await CreateRunData(), runData => { if (someCondition) { // 发射资源后立即完成当前子序列,触发Using释放资源 return Observable.Return(runData); } return Observable.Empty<RunData>(); } ) ) // 等待2秒后重复 .Concat(Observable.Empty<RunData>().Delay(TimeSpan.FromSeconds(2))) .Repeat() .Publish().RefCount();
注意:该方案仅适用于订阅者处理逻辑为同步、或Rx操作符能保证处理完成后才触发序列完成的场景。若存在异步处理,仍需方案1的信号跟踪机制。
多订阅者场景补充
使用Publish().RefCount()时,多个订阅者会共享同一个RunData,此时需要跟踪所有订阅者的处理完成状态,可将AsyncSubject替换为CountDownEvent或使用Observable.WhenAll合并所有订阅者的处理信号,确保所有订阅者处理完成后再释放资源。
内容的提问来源于stack exchange,提问作者Maximilian Csuk
相关产品推荐
相关产品推荐

