处理带间歇性错误可临时中断Observable序列的正确方案
问题根源
你遇到的所有问题本质上是你用的GraphQL.NET客户端违反了Rx.NET核心规范:Rx规范明确规定onError是序列的终止信号,一旦触发,序列就已经结束,后续不允许再调用任何onNext/onError/onCompleted方法。这个客户端在触发onError之后还继续推送消息、重复触发onError,本质上是把临时错误当成了普通事件发,却错误地占用了终止信号的通道。
你之前尝试的所有标准Rx操作符(Do/Retry/Catch)都是为符合规范的流设计的,自然无法适配这种错误实现:
Do只会执行副作用,不会拦截终止信号,错误还是会传到下游终止序列Retry/Catch的逻辑都是"上游触发onError终止后,重新订阅/切换备用流",但这个客户端的上游流触发onError后根本没有终止,还在继续发消息,这两个操作符会直接丢弃仍在工作的原流,自然会出现订阅失效、收不到后续消息的问题。- 你最开始直接把上游的
onError透传给对外暴露的Subject,第一次报错就会把Subject置为终止状态,后续所有消息都会被忽略,这是最核心的实现错误。
正确实现方案
不要强行用标准操作符适配上游的违规实现,在你的封装层直接把上游的行为纠正为符合Rx规范的流即可,核心原则非常简单:
只有真正会导致序列永久终止的错误,才能走
onError通道;临时连接错误、重连事件本质是普通通知,和业务数据一样走onNext通道即可。
推荐实现:分离数据与错误通知流
你之前基于事件的实现本来就是用两个独立事件分别传递数据和错误,转Observable的时候完全没必要强行把两类消息塞到同一个流里,保留两个独立的Observable序列是最清爽、最符合原有使用习惯的方案,完全不需要写丑陋的"值+异常"包装类。
封装层代码示例
public class PubSubSubscription { private readonly IGraphQLClient _client; public PubSubSubscription(IGraphQLClient client) { _client = client; } public (IObservable<TResponse> Messages, IObservable<Exception> ConnectionState) CreateSubscription<TResponse>(string topic) { var dataSubject = new Subject<TResponse>(); var stateSubject = new Subject<Exception>(); // 管理上游订阅的生命周期,避免内存泄漏 var upstreamDisposable = new SerialDisposable(); var stream = _client.CreateSubscriptionStream<TResponse>(topic); upstreamDisposable.Disposable = stream.Subscribe( onNext: response => { // 原有校验、映射、日志逻辑 dataSubject.OnNext(response); }, onError: ex => { // 先判断是否为永久不可恢复错误 if (IsPermanentFailure(ex)) { // 永久错误才走onError终止序列 dataSubject.OnError(ex); stateSubject.OnError(ex); return; } // 临时连接错误、重连事件作为普通状态通知走onNext,绝不终止数据流 stateSubject.OnNext(ex); }, onCompleted: () => { dataSubject.OnCompleted(); stateSubject.OnCompleted(); } ); // 绑定生命周期:下游所有订阅都释放时,自动取消上游订阅 var dataStream = dataSubject.AsObservable().Finally(() => upstreamDisposable.Dispose()); var stateStream = stateSubject.AsObservable(); return (dataStream, stateStream); } private bool IsPermanentFailure(Exception ex) { // 按业务规则判断:比如鉴权失败、订阅主题不存在、参数非法等重试也无法恢复的错误返回true return ex switch { UnauthorizedAccessException => true, ArgumentException => true, _ => false }; } }
调用方代码示例
var (orderStream, connectionStateStream) = ordersPubSub.CreateSubscription<Order>("/orders"); // 订阅业务数据 orderStream .ObserveOn(SynchronizationContext.Current) // 切到UI线程 .Subscribe( onNext: order => { /* 更新UI展示订单数据 */ }, onError: fatalEx => { /* 只有永久致命错误才会走到这,弹框提示、关闭订阅页即可 */ } ); // 订阅连接状态(可选,不需要展示状态可以不订) connectionStateStream .ObserveOn(SynchronizationContext.Current) .Subscribe(ex => { // 临时错误处理:比如状态栏展示"连接断开,正在尝试重连..." });
备选方案:单流传递所有事件
如果你确实不想拆成两个流,也不要写带两个可空属性的弱类型包装类,用可识别联合做类型安全的消息区分即可,C#可以直接用record实现:
// 类型安全的流事件定义,不会出现值和错误同时为空/同时有值的非法状态 public abstract record SubscriptionMessage<T>; public record DataMessage<T>(T Value) : SubscriptionMessage<T>; public record ConnectionErrorMessage<T>(Exception Error) : SubscriptionMessage<T>;
封装层只需要把所有消息都作为SubscriptionMessage<T>通过onNext发送即可,下游订阅时直接用模式匹配处理:
stream.Subscribe(msg => msg switch { DataMessage<Order> d => /* 处理订单数据 */, ConnectionErrorMessage<Order> e => /* 展示重连提示 */ });
方案优势
- 所有重连、错误适配逻辑全部封装在
CreateSubscription内部,调用方不需要重复写任何Retry/Catch逻辑 - 完全符合Rx规范,对外暴露的流行为可预测,不会出现意外终止的问题
- 不需要做别扭的类型包装,数据和状态分离,类型安全
- 和你之前基于事件的实现逻辑完全一致,迁移成本极低
内容的提问来源于stack exchange,提问作者Jason
相关产品推荐
相关产品推荐

