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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 11:45:03