RX.NET高阶Observable异常处理引发RunData资源未释放问题求助
问题背景
在学习RX.NET时,为了在Catch()操作符中访问流经管道的RunData实例,采用了高阶可观测对象结合Concat()的异常处理模式,但测试发现RunData的创建数量始终比释放数量多一个,导致断言失败;关闭该高阶异常处理逻辑后,测试恢复正常。
原因分析
核心问题在于高阶处理逻辑中重复订阅了源Observable:
在异常处理分支的代码里,Select(rd => TransformRunDataToResult(_obs))这一步,内部的TransformRunDataToResult(_obs)会再次订阅共享的源_obs(即便用了Publish().RefCount())。每次外层Select拿到一个RunData时,内部都会触发一次源的订阅,额外生成一个RunData实例。
这个额外生成的实例会被内部Observable处理并释放,但外层Select拿到的初始RunData实例并没有被下游订阅消费(因为最终Concat()的是内部Observable的输出),导致该实例始终未被释放,计数器出现差值。
从测试输出也能验证这一点:第一次输出的Created是外层Select获取的实例,第二个Created是内部订阅源产生的实例,后续只处理并释放了内部产生的实例,第一个实例永远没有进入下游的Dispose逻辑。
解决方案
不要在Select内部订阅整个源Observable,而是将单个RunData元素包装成独立的Observable,再针对这个单元素Observable做异常处理,这样就不会触发源的重复订阅。
修改后的核心代码
if (useHigherOrderExceptionHandling) { safeObs = obs .Select(rd => // 将单个rd包装成Observable,而非订阅整个源 Observable.Return(rd) .Select(r => (result: true, runData: r)) .Catch((Exception e) => Observable.Return((result: false, runData: rd))) ) .Concat(); }
复用Transform方法的版本
如果要保留TransformRunDataToResult的复用性,可以调整方法接收单个元素:
// 调整转换方法,接收单个RunData而非整个Observable IObservable<(bool result, RunData runData)> TransformRunDataToResult(RunData rd) { return Observable.Return(rd) .Select(r => (result: true, runData: r)); } // 异常处理分支修改为 safeObs = obs .Select(rd => TransformRunDataToResult(rd) .Catch((Exception e) => Observable.Return((result: false, runData: rd))) ) .Concat();
这样修改后,每个RunData实例只会被创建一次,且都会进入下游的Dispose逻辑,计数器最终会归零,测试通过。
内容的提问来源于stack exchange,提问作者Maximilian Csuk

