【问题标题】:Observable from RefCount() doesn't stop publishing来自 RefCount() 的 Observable 不会停止发布
【发布时间】:2019-04-19 23:17:10
【问题描述】:

我正在构建一个消息处理管道,并注意到当最后一个观察者处理订阅时,可观察者仍在通过泵送数据。

我查看了 Rx 文档,我的假设是,根据文档,一旦最后一个观察者取消订阅,RefCount() 将断开可观察对象:

RefCount 然后跟踪有多少其他观察者订阅它,并且不会断开与底层可连接 Observable直到最后一个观察者这样做。

为了说明这个问题,我在下面创建了一个非常简约的示例:

class Program
{
    static void Main(string[] args)
    {
        _ = SimulateObservableIssue();

        Console.ReadKey();
    }

    public static async Task SimulateObservableIssue()
    {
        IObservable<int> source = Observable.Create<int>(async (observer) =>
        {
            for (int i = 0; i < 10; i++)
            {
                Console.WriteLine($"Source publishing {i}");
                observer.OnNext(i);
                await Task.Delay(1000);
            }

            observer.OnCompleted();

            return Disposable.Create(() => Console.WriteLine("Observable is disposed"));
        });

        var multiSource = source.Publish().RefCount();

        var subscription = multiSource.Subscribe(x => Console.WriteLine("Observer received: " + x));

        await Task.Delay(3000);

        subscription.Dispose();

        Console.WriteLine("Subscription disposed");

    }
}

Output:

Source publishing 0
Observer received: 0
Source publishing 1
Observer received: 1
Source publishing 2
Observer received: 2
Subscription disposed
Source publishing 3
Source publishing 4
Source publishing 5
Source publishing 6
Source publishing 7
Source publishing 8
Source publishing 9
Observable is disposed

为什么在subscription.Dispose() 之后,observable 仍在尝试生成数据?

【问题讨论】:

    标签: c# .net system.reactive rx.net


    【解决方案1】:

    您的 source observable 不尊重您提到的 Observable 合同。如果你用这个替换source:

        var source = Observable.Interval(TimeSpan.FromSeconds(1))
            .Do(i => Console.WriteLine($"Source publishing {i}"), () => Console.WriteLine("Observable is disposed"))
            .Take(10);
    

    ...您会看到它按预期工作。

    至于为什么,想想 observable 有两个阶段:订阅和观察。无论订阅取消如何,在订阅期间发生的代码总是会发生。 Observable.Createcode 都是订阅码。

    我写的 observable 都是 observable 代码(就像大多数 observable 代码一样)。因此它会适当地响应订阅取消。

    【讨论】:

    • 感谢您的回答。那么创建自定义可观察对象的正确方法是什么?
    • “自定义可观察”可能意味着很多事情。很可能您不想要或不需要一个,而宁愿像我上面所做的那样从标准运算符中构建一个。你想让你的 observable 做什么?
    • 根据我的示例,我希望我的 observable 将一些数据多播给多个订阅者。我已经在实际实现中实现了这一点,但正如强调的那样,当最后一个观察者取消订阅并且仍在生成数据时,这个问题面临着效率低下的问题。所以我的问题真的是,对于尊重 Observable 合同和 RefCount 能够在内部正确处理订阅者的自定义 observable 实现,“Observable.Create”的替代方案是什么?我可以处理状态,但 RefCount 似乎是开箱即用的解决方案。
    • 是的:RefCount 很好。你在 Observable.Create 中做了哪些其他 Observable 运算符无法表示的事情?
    • 好吧,我在问题中的示例可能过于简约,无法说明我在实际实现中所做的事情。但是,假设我想要一个 observable 来生成具有不同数值的对象,我可能需要探索更多的 Observable 运算符,但我目前看不到一个直接的方法来做到这一点,至少从上面提供的示例中。你知道我是否缺少任何东西,但可以用作'Observable.Create'的替代品吗?
    【解决方案2】:

    您观察到的行为与RefCount 运算符无关。如果您直接订阅source 而不是multiSource,行为将是相同的。

    问题与您如何使用带有异步 lambda 的 Observable.Create 来创建自定义 observable 有关。 Rx 库无法终止由异步 lambda 创建的 Task。因此,尽管观察者已取消订阅,但任务仍在继续。一般来说,任务只能以合作方式终止,而用于终止不再需要的任务的standard mechanism 是通过向任务观察到的CancellationToken 发出信号。 Rx 库通过具有以下签名的 Observable.Create 重载支持此模式:

    public static IObservable<TResult> Create<TResult>(
        Func<IObserver<TResult>, CancellationToken, Task> subscribeAsync);
    

    提供的CancellationToken 由库提供和管理,您的责任是通过终止循环来兑现取消信号。一种方法是在循环内的各个点调用CancellationToken.ThrowIfCancellationRequested() 方法。在具体示例中,将CancellationToken 作为参数传递给Task.Delay 方法就足够了:

    var source = Observable.Create<int>(async (observer, cancellationToken) =>
    {
        cancellationToken.Register(() => Console.WriteLine("Token is canceled"));
        for (int i = 0; i < 10; i++)
        {
            Console.WriteLine($"Source publishing {i}");
            observer.OnNext(i);
            await Task.Delay(1000, cancellationToken);
        }
        observer.OnCompleted();
        Console.WriteLine($"Source is completed");
    });
    

    在取消的情况下,await Task.Delay 抛出的异常不会被任何人观察到,因为此时 observable 将没有观察者。没有理由尝试/捕捉它,除非你想记录它。

    您可能已经注意到 lambda 不返回 Disposable。这是因为Disposable 的角色现在由提供的CancellationToken 扮演。在后台,该库使用CancellationDisposable,它在处理时取消了内部CancellationTokenSource,其中Token 将传递给您的方法。

    【讨论】:

    • 赞成:使用CancellationToken 是此场景的正确解决方案。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多