System.Reactive 如何实现一次性批量销毁未知数量的订阅
问题解答
1. 批量销毁订阅的原生机制
System.Reactive 原生提供了CompositeDisposable类型,专门用来批量管理一组IDisposable对象(包括Rx订阅返回的销毁句柄),调用一次Dispose方法即可销毁所有内部管理的订阅,完全支持动态新增/移除订阅、适配不确定的订阅量级,且本身实现了线程安全,是这个场景的最优解。
2. Observable.Using封装是否能实现需求?
可以实现,但不推荐在这个场景下使用。Observable.Using的设计目标是将单个Observable的生命周期与某个资源的销毁逻辑绑定:当订阅这个Using返回的Observable时创建资源,当该Observable终止(OnCompleted/OnError)或订阅被销毁时自动释放资源。如果要用来批量管理多个独立订阅,需要把所有订阅逻辑都写到Using的资源创建委托里,逻辑耦合度高,且不支持后续动态新增订阅,灵活度远低于直接使用CompositeDisposable。
3. 代码改造示例
第一步:声明全局订阅管理器
// 作为类的私有字段,用来管理所有相关订阅 private readonly CompositeDisposable _allSubscriptions = new CompositeDisposable();
第二步:将所有订阅加入管理器
Rx 提供了DisposeWith扩展方法,可以直接把订阅返回的IDisposable对象加入到CompositeDisposable中:
client.Streams.PongStream.Subscribe(x => Log.Information($"Pong received ({x.Message})")) .DisposeWith(_allSubscriptions); client.Streams.FundingStream.Subscribe(response => { var funding = response.Data; Log.Information($"Funding: [{funding.Symbol}] rate:[{funding.FundingRate}] " + $"mark price: {funding.MarkPrice} next funding: {funding.NextFundingTime} " + $"index price {funding.IndexPrice}"); }).DisposeWith(_allSubscriptions); client.Streams.AggregateTradesStream.Subscribe(response => { var trade = response.Data; Log.Information($"Trade aggreg [{trade.Symbol}] [{trade.Side}] " + $"price: {trade.Price} size: {trade.Quantity}"); }).DisposeWith(_allSubscriptions); // 其余所有订阅都按照上述格式添加DisposeWith即可,后续新增订阅也只需要加这一行就能纳入统一管理
第三步:统一销毁所有订阅
需要销毁全量订阅时,只需要调用一次Dispose:
_allSubscriptions.Dispose();
补充注意事项
CompositeDisposable调用Dispose后就进入已销毁状态,无法再添加新的订阅。如果后续还需要新增订阅,可以重新实例化一个新的CompositeDisposable赋值给对应字段即可。- 如果需要单独销毁某个订阅而不影响其他订阅,可以先存储该订阅返回的
IDisposable对象,调用_allSubscriptions.Remove(yourSubscription)即可单独移除并销毁该订阅。
内容的提问来源于stack exchange,提问作者nop
相关产品推荐
相关产品推荐

