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

如何区分Observable因错误、完成或主动取消订阅而触发释放?

问题描述

现有如下订阅代码:

public IObservable<Item> GetItems()
{
  return from batch in CreateItemsObservable()
         .OnDispose(() => Console.WriteLine("disposed"))
         select batch;
}

其中OnDispose是自定义扩展方法(注:原代码中CompositeObservable应为Rx标准类型CompositeDisposable):

public static IObservable<T> OnDispose<T>(this IObservable<T> observable, Action onDispose)
{
  return Observable.Create<T>(
    observer => new CompositeDisposable(observable.Subscribe(
      observer.OnNext,
      observer.OnError,
      observer.OnCompleted), Disposable.Create(onDispose)));
}

希望区分释放操作的触发原因:是流正常完成、主动取消订阅,还是发生错误导致的释放。考虑过用Materialize操作符传递通知到释放动作,但不确定如何应用,想知道是否有更优实现方式。

解决方案

完全可以修改OnDispose方法来区分释放原因,以下是两种可行方案,优先推荐状态跟踪的实现:

方案1:自定义释放原因枚举+状态跟踪

首先定义枚举明确释放原因:

public enum DisposeReason
{
    // 用户主动取消订阅
    Unsubscribed,
    // 流正常执行完成
    Completed,
    // 流抛出错误终止
    Errored
}

修改OnDispose方法,通过跟踪流的最终状态传递原因:

public static IObservable<T> OnDispose<T>(this IObservable<T> observable, Action<DisposeReason> onDispose)
{
    return Observable.Create<T>(observer =>
    {
        var disposeReason = DisposeReason.Unsubscribed;
        var subscription = observable.Subscribe(
            value => observer.OnNext(value),
            error =>
            {
                disposeReason = DisposeReason.Errored;
                observer.OnError(error);
            },
            () =>
            {
                disposeReason = DisposeReason.Completed;
                observer.OnCompleted();
            });

        return Disposable.Create(() =>
        {
            subscription.Dispose();
            onDispose(disposeReason);
        });
    });
}

使用示例

public IObservable<Item> GetItems()
{
    return CreateItemsObservable()
        .OnDispose(reason => 
        {
            Console.WriteLine($"释放原因:{reason}");
        });
}

逻辑说明

  • 默认disposeReason为Unsubscribed:如果用户主动取消订阅,流的OnError/OnCompleted不会触发,保持初始值。
  • 流触发OnError时,先更新原因再转发错误;触发OnCompleted时同理更新原因。
  • 最终在释放回调中,传入跟踪到的真实原因。

方案2:使用Materialize操作符

若需基于Materialize实现,可通过包装通知判断终止类型:

public static IObservable<T> OnDispose<T>(this IObservable<T> observable, Action<DisposeReason> onDispose)
{
    return observable.Materialize()
        .Do(notification =>
        {
            if (notification.Kind == NotificationKind.OnCompleted)
                onDispose(DisposeReason.Completed);
            else if (notification.Kind == NotificationKind.OnError)
                onDispose(DisposeReason.Errored);
        })
        .Dematerialize()
        .Finally(() => onDispose(DisposeReason.Unsubscribed));
}

注意事项

  • Materialize将所有流通知包装为Notification<T>,通过Kind属性判断终止类型。
  • Finally处理主动取消订阅的场景:无论流是否正常终止,主动取消时Finally都会触发。
  • 此方案需额外的包装/解包装操作,性能略低于方案1,适合已有Materialize使用场景的情况。

方案对比

  • 方案1:实现简洁、性能优异、逻辑直观,是首选方案。
  • 方案2:依赖Materialize特性,适合特定场景,但开销更高。

内容的提问来源于stack exchange,提问作者Flack

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 06:21:02