System.Reactive中Subscribe订阅异常的捕获与展示方案咨询
解决方案
问题根因
Rx 标准订阅的 onError 回调默认仅能捕获源流自身推送的异常(如连接失败、数据反序列化错误等),如果你在 onNext 业务处理逻辑中主动抛出异常,默认会直接终止整个 Observable 序列,导致后续消息无法继续接收,部分封装的流组件也可能存在异常传递链路断裂的情况。
实现方案:自定义安全订阅扩展方法
我们可以封装统一的 SafeSubscribe 扩展方法,同时覆盖业务处理异常和源流自身异常两种场景,出错后自动重试保持订阅存活,无需每个订阅重复写异常处理逻辑。
扩展方法代码
using System; using System.Reactive; using System.Reactive.Linq; using Serilog; public static class RxSubscribeExtensions { /// <summary> /// 带全局异常捕获的安全订阅方法 /// </summary> /// <typeparam name="T">流元素类型</typeparam> /// <param name="source">源可观察序列</param> /// <param name="onNext">元素业务处理逻辑</param> /// <param name="customErrorHandler">自定义异常处理逻辑,不传默认使用Serilog输出错误日志</param> /// <returns>订阅句柄,可加入CompositeDisposable统一释放</returns> public static IDisposable SafeSubscribe<T>(this IObservable<T> source, Action<T> onNext, Action<Exception>? customErrorHandler = null) { // 缺省异常处理:输出完整错误堆栈到日志 var errorHandler = customErrorHandler ?? (ex => Log.Error(ex, "Rx流处理发生未捕获异常")); return source // 捕获业务处理逻辑抛出的异常,避免终止流 .Select(item => { try { onNext(item); return Unit.Default; } catch (Exception ex) { errorHandler(ex); return Unit.Default; } }) // 捕获源流自身抛出的异常 .Catch((Exception ex) => { errorHandler(ex); // 异常处理完成后返回空序列,配合Retry实现订阅不终止 return Observable.Empty<Unit>(); }) // 异常后自动重试,保持订阅存活 .Retry() .Subscribe(); } }
使用方法
把你所有的 .Subscribe 调用替换为 .SafeSubscribe 即可,示例如下:
// 无需自定义异常处理的场景,默认打错误日志 client.Streams.KlineStream .SafeSubscribe(response => { Guard.Against.Null(response, nameof(response)); Guard.Against.Null(response.Data, nameof(response.Data), "Something went wrong and the kline object is null"); Guard.Against.Null(response.Data.Data, nameof(response.Data.Data), "Something went wrong and the kline data object is null"); var kline = response.Data; var klineData = response.Data.Data; Log.Information($"Kline [{kline.Symbol}] " + $"Kline start time: {klineData.StartTime} " + $"Kline close time: {klineData.CloseTime} " + $"Interval: {klineData.Interval} " + $"First trade ID: {klineData.FirstTradeId} " + $"Last trade ID: {klineData.LastTradeId} " + $"Open price: {klineData.OpenPrice} " + $"Close price: {klineData.ClosePrice} " + $"High price: {klineData.HighPrice} " + $"Low price: {klineData.LowPrice} " + $"Base asset volume: {klineData.BaseAssetVolume} " + $"Number of trades: {klineData.NumberTrades} " + $"Is this kline closed?: {klineData.IsClosed} " + $"Quote asset volume: {klineData.QuoteAssetVolume} " + $"Taker buy base: {klineData.TakerBuyBaseAssetVolume} " + $"Taker buy quote: {klineData.TakerBuyQuoteAssetVolume} " + $"Ignore: {klineData.Ignore} " ); }) .DisposeWith(disposable); // 需要自定义异常处理的场景 client.Streams.AggregateTradesStream .SafeSubscribe(response => { throw new Exception("Asd"); Guard.Against.Null(response, nameof(response)); Guard.Against.Null(response.Data, nameof(response.Data), "Something went wrong and the aggregated trade object is null"); var trade = response.Data; Log.Information($"Aggregated trade [{trade.Symbol}] [{trade.Side}] " + $"Price: {trade.Price} Size: {trade.Quantity}"); }, ex => Console.WriteLine("Exception: {0} {1}", ex.Message, DateTime.Now)) .DisposeWith(disposable);
可选优化
如果不需要所有异常都自动重试,可以在 Catch 逻辑中添加异常类型过滤,仅对可恢复的异常执行重试,不可恢复的异常直接终止订阅即可。
内容的提问来源于stack exchange,提问作者nop
相关产品推荐
相关产品推荐

