【问题标题】:How to fix the inconsistency of the Publish().RefCount() behavior?如何解决 Publish().RefCount() 行为不一致的问题?
【发布时间】:2020-11-23 01:02:13
【问题描述】:

最近我偶然发现了 Enigmativity 的 interesting statement 关于 PublishRefCount 运算符:

您正在使用危险的 .Publish().RefCount() 运算符对,它创建了一个在完成后无法订阅的序列。

此声明似乎反对 Lee Campbell 对这些运营商的评估。引用他的书Intro to Rx

Publish/RefCount 对对于获取冷的 observable 并将其作为热的 observable 序列共享给后续观察者非常有用。

一开始我不相信 Enigmativity 的说法是正确的,所以我试图反驳它。我的实验表明Publish().RefCount() 可以 确实不一致。第二次订阅已发布的序列可能会导致对源序列的新订阅,这取决于源序列是否在连接时完成。如果已完成,则不会重新订阅。如果未完成,则将重新订阅。这是此行为的演示:

var observable = Observable
    .Create<int>(o =>
    {
        o.OnNext(13);
        o.OnCompleted(); // Commenting this line alters the observed behavior
        return Disposable.Empty;
    })
    .Do(x => Console.WriteLine($"Producer generated: {x}"))
    .Finally(() => Console.WriteLine($"Producer finished"))
    .Publish()
    .RefCount()
    .Do(x => Console.WriteLine($"Consumer received #{x}"))
    .Finally(() => Console.WriteLine($"Consumer finished"));

observable.Subscribe().Dispose();
observable.Subscribe().Dispose();

在本例中,observable 由三部分组成。首先是生成单个值然后完成的生产部分。然后遵循发布机制(Publish+RefCount)。最后是观察生产者发出的值的消费部分。 observable 被订阅了两次。预期的行为是每个订阅都会收到一个值。但这不是发生的事情!这是输出:

Producer generated: 13
Consumer received #13
Producer finished
Consumer finished
Consumer finished

(Try it on fiddle)

如果我们注释o.OnCompleted(); 行,这是输出。这种细微的变化会导致一种预期和理想的行为:

Producer generated: 13
Consumer received #13
Producer finished
Consumer finished
Producer generated: 13
Consumer received #13
Producer finished
Consumer finished

在第一种情况下,cold 生产者(Publish().RefCount() 之前的部分)只订阅了一次。第一个消费者收到了发出的值,但第二个消费者什么也没收到(OnCompleted 通知除外)。在第二种情况下,生产者被订阅了两次。每次它产生一个值,每个消费者得到一个值。

我的问题是:我们如何解决这个问题?我们如何修改Publish 运算符或RefCount 或两者,以使它们的行为始终一致且合乎需要?以下是理想行为的规范:

  1. 发布的序列应该将所有直接来自源序列的通知传播给它的订阅者,而不是其他任何东西。
  2. 当当前订阅者数量从零增加到一时,已发布序列应订阅源序列。
  3. 只要至少有一个订阅者,发布的序列就应该与源保持连接。
  4. 当当前订阅者数量为零时,已发布序列应取消订阅源。

我要求提供上述功能的自定义 PublishRefCount 运算符,或使用内置运算符实现所需功能的方法。

顺便说一句,similar question 存在,它问为什么会发生这种情况。我的问题是关于如何解决它。


更新:回想起来,上述规范导致了一种不稳定的行为,使得竞态条件不可避免。不能保证对已发布序列的两次订阅将导致对源序列的一次订阅。源序列可能在两个订阅之间完成,导致第一个订阅者取消订阅,导致RefCount 运算符取消订阅,导致下一个订阅者重新订阅源。内置 .Publish().RefCount() 的行为可以防止这种情况发生。

道德教训是.Publish().RefCount() 序列没有损坏,但它不可重复使用。它不能可靠地用于多个连接/断开会话。如果你想要第二个会话,你应该创建一个新的.Publish().RefCount() 序列。

【问题讨论】:

  • 为什么你认为你观察到的行为是错误的?这正是我期望它的行为方式。
  • @Shlomo 在RefCountdefinition 中没有提示它保持源observable 的状态(无论是否完成),并使用此内存发送OnCompleted代表其发出通知。我的期望是它只保持当前的订阅者数量。换句话说,RefCount 做的比它应该做的要多。我宁愿它是愚蠢和可预测的,而不是有自己的想法。
  • 它没有这样做。发布是。
  • @Shlomo 那太好了。在这种情况下,只有Publish 需要修复!
  • @Enigmativity 啊,那好吧。很抱歉造成误解。

标签: c# system.reactive rx.net


【解决方案1】:

Lee 做了一个good job 解释IConnectableObservable,但Publish 解释得不是很好。这是一种非常简单的动物,很难解释。我假设你理解IConnectableObservable:

如果我们简单而懒惰地重新实现零参数 Publish 函数,它看起来像这样:

//  For illustrative purposes only: don't use this code
public class PublishObservable<T> : IConnectableObservable<T>
{
    private readonly IObservable<T> _source;
    private readonly Subject<T> _proxy = new Subject<T>();
    private IDisposable _connection;
    
    public PublishObservable(IObservable<T> source)
    {
        _source = source;
    }
    
    public IDisposable Connect()
    {
        if(_connection == null)
            _connection = _source.Subscribe(_proxy);
        var disposable = Disposable.Create(() =>
        {
            _connection.Dispose();
            _connection = null;
        });
        return _connection;
    }

    public IDisposable Subscribe(IObserver<T> observer)
    {
        var _subscription = _proxy.Subscribe(observer);
        return _subscription;
    }
}

public static class X
{
    public static IConnectableObservable<T> Publish<T>(this IObservable<T> source)
    {
        return new PublishObservable<T>(source);
    }
}

Publish 创建一个订阅源 observable 的代理 Subject。代理可以根据连接订阅/取消订阅源:调用Connect,代理订阅源。在连接上调用Dispose,代理从源取消订阅。重要的一点是,有一个Subject 代理与源的任何连接。不能保证您只订阅一个源,但可以保证您有一个代理和一个并发连接。您可以通过连接/断开连接进行多个订阅。

RefCount 处理调用Connect 的部分事情:这是一个简单的重新实现:

//  For illustrative purposes only: don't use this code
public class RefCountObservable<T> : IObservable<T>
{
    private readonly IConnectableObservable<T> _source;
    private IDisposable _connection;
    private int _refCount = 0;

    public RefCountObservable(IConnectableObservable<T> source)
    {
        _source = source;
    }
    public IDisposable Subscribe(IObserver<T> observer)
    {
        var subscription = _source.Subscribe(observer);
        var disposable = Disposable.Create(() =>
        {
            subscription.Dispose();
            DecrementCount();
        });
        if(++_refCount == 1)
            _connection = _source.Connect();
            
        return disposable;
    }

    private void DecrementCount()
    {
        if(--_refCount == 0)
            _connection.Dispose();
    }
}
public static class X
{
    public static IObservable<T> RefCount<T>(this IConnectableObservable<T> source)
    {
        return new RefCountObservable<T>(source);
    }

}

更多代码,但仍然非常简单:如果 refcount 上升到 1,则在 ConnectableObservable 上调用 Connect,如果下降到 0,则断开连接。

将两者放在一起,您将得到一对保证只有一个并发订阅源 observable,通过一个持久的Subject 代理。 Subject 只会在有 >0 个下游订阅时订阅源。


鉴于您的介绍,您的问题中有很多误解,所以我将一一解释:

... Publish().RefCount() 确实可能不一致。订阅一秒 发布序列的时间可能会导致新订阅 源序列与否,取决于源序列是否 连接时完成。如果完成了就不会了 重新订阅。如果未完成,则会重新订阅。

.Publish().RefCount() 将仅在一种情况下重新订阅源:当它从零订阅者变为 1 时。如果订阅者数量从 0 到 1 到 0 到 1出于任何原因然后你最终会重新订阅。源 observable 完成将导致 RefCount 发出 OnCompleted,并且其所有观察者都取消订阅。所以后续订阅RefCount 将触发重新订阅源的尝试。自然,如果源正确地观察可观察合同,它会立即发出OnCompleted,就是这样。

[查看带有 OnCompleted 的示例 observable...] observable 被订阅了两次。这 预期的行为是每个订阅都会收到一个 价值。

没有。预期的行为是代理 Subject 在发出 OnCompleted 后将重新发出 OnCompleted 到任何后续订阅尝试。由于您的源 observable 在您的第一个订阅结束时同步完成,因此第二个订阅将尝试订阅一个已经发出 OnCompletedSubject。这应该导致OnCompleted,否则 Observable 合约将被破坏。

[参见示例 observable 没有 OnCompleted 作为第二种情况...] 在 第一种情况是冷生产者(之前的部分 Publish().RefCount()) 只订阅了一次。第一消费者 收到了发出的值,但第二个消费者什么也没收到 (除了 OnCompleted 通知)。在第二种情况下 生产者被订阅了两次。每次它产生一个值,并且 每个消费者都有一个值。

这是正确的。由于代理 Subject 从未完成,后续重新订阅源将导致冷 observable 重新运行。

我的问题是:我们如何解决这个问题? [..]

  1. 发布的序列应该将所有直接来自源序列的通知传播给它的订阅者,而不是什么 否则。
  2. 当当前订阅者数量从零增加到一时,已发布序列应订阅源序列。
  3. 只要至少有一个订阅者,发布的序列就应该与源保持连接。
  4. 当当前订阅者数量为零时,已发布序列应取消订阅源。

只要您没有完成/错误,目前.Publish.RefCount 目前都会发生上述所有情况。我不建议实施一个改变它的操作符,从而破坏 Observable 合同。


编辑

我认为与 Rx 混淆的第一大来源是 Hot/Cold observables。由于 Publish 可以“预热”冷的 observables,因此它会导致令人困惑的边缘情况也就不足为奇了。

首先,关于可观察合同。 Observable 合约更简洁地表述为OnNext 永远不能跟随OnCompleted/OnError,并且应该只有一个OnCompleted OnError 通知。这确实留下了尝试订阅终止的 observables 的边缘情况: 尝试订阅终止的 observable 会立即收到终止消息。这会破坏合同吗?也许吧,但据我所知,这是图书馆里唯一的合同作弊。另一种选择是订阅死空气。这对任何人都没有帮助。

这如何与热/冷可观察对象联系起来?不幸的是,令人困惑。订阅冰冷的 observable 会触发整个 observable 管道的重建。这意味着 subscribe-to-already-terminated 规则仅适用于 hot observables。冷的 observables 总是重新开始。

考虑这段代码,其中o 是一个冷可观察对象。:

var o = Observable.Interval(TimeSpan.FromMilliseconds(100))
    .Take(5);
var s1 = o.Subscribe(i => Console.WriteLine(i.ToString()));
await Task.Delay(TimeSpan.FromMilliseconds(600));
var s2 = o.Subscribe(i => Console.WriteLine(i.ToString()));

就合约而言,s1 后面的 observable 和 s2 后面的 observable 完全不同。因此,即使它们之间存在延迟,并且您最终会在 OnCompleted 之后看到 OnNext,但这不是问题,因为它们是完全不同的 observables。

热身后的Publish 版本会引起人们的注意。如果您要在上面的代码中将.Publish().RefCount() 添加到o 的末尾...

  • 在不更改任何其他内容的情况下,s2 将立即终止打印任何内容。
  • 将延迟更改为 400 左右,s2 将打印最后两个数字。
  • s1 更改为仅.Take(2)s2 将重新开始打印 0 到 4。

更糟糕的是,薛定谔的猫效应:如果你在o 上设置一个观察者来观察整个时间会发生什么,这会改变引用计数,影响功能!看着它,改变行为。调试噩梦。

这是尝试“热身”冷的 observables 的危险。它只是不能很好地工作,尤其是Publish/RefCount

我的建议是:

  1. 不要尝试加热冷的 observables。
  2. 如果您需要与冷或热 observable 共享订阅,请遵守 @Enigmativity 严格使用选择器 Publish 版本的一般规则
  3. 如果必须,请在 Publish/RefCount observable 上进行虚拟订阅。这至少提供了一致的 Refcount >= 1,从而降低了量子活动效应。

【讨论】:

  • 感谢 Shlomo 的回答,很有启发性!我赞成它,但我不接受它作为答案,因为它不同意问题的基本前提,即存在需要解决的问题。我编辑了问题以澄清您强调的“违反合同”问题。
  • 但现在我又糊涂了。不是所有的 cold observables 都破坏了the contract 吗?例如,Observable.Return 怎么样?每次订阅时,此运算符都会发出 OnCompleted 通知。所以它要么违反合同,要么合同管理对 observable 的单一订阅,而不是整个生命周期。在澄清这一点之前,我将撤销我的编辑!
  • 我不喜欢这种解释,即一旦发出OnCompleted,所有后续订阅都应该立即完成。我觉得当从 0 到 1 个观察者时,它应该重新订阅源 observable。这正是我所期望的。想象一下,如果Observable.Return(42) 只返回42 给第一个连接的观察者,但所有后续的观察者都得到完整的。
  • 添加到另一篇文章中。我实际上不确定它是否有帮助。只是 API 中令人困惑的部分。
  • @Shlomo - 是的,我知道它是这样工作的,但这不是我对操作员的直观理解。我希望如果引用计数为零,那么任何新订阅都会获得新的Observable.Return(42)
【解决方案2】:

作为 Shlomo pointed out,此问题与 Publish 运算符有关。 RefCount 工作正常。所以需要修复的是PublishPublish 只不过是使用标准 Subject&lt;T&gt; 作为参数调用 Multicast 运算符。这是它的source code

public IConnectableObservable<TSource> Publish<TSource>(IObservable<TSource> source)
{
    return source.Multicast(new Subject<TSource>());
}

所以Publish 运算符继承了Subject 类的行为。这个类,有很好的理由,保持它的完成状态。因此,如果您通过调用subject.OnCompleted() 表示其完成,则该主题的任何未来订阅者将立即收到OnCompleted 通知。此功能可以很好地服务于独立的主题及其订阅者,但是当Subject 用作源序列和该序列的订阅者之间的中间传播者时,它会成为一个有问题的工件。这是因为源序列已经保持了它自己的状态,并且在主体内部复制这个状态会引入两个状态变得不同步的风险。这正是PublishRefCount 运算符组合时发生的情况。对象记得源已经完成,而源作为一个冷序列,已经失去了对前世的记忆,并愿意重新开始新的生活。

因此解决方案是为Multicast 运算符提供无状态主题。不幸的是,我找不到基于内置Subject&lt;T&gt; 的方法来编写它(继承不是一种选择,因为类是密封的)。幸运的是,从头开始实现它并不是很困难。下面的实现使用ImmutableArray作为对象观察者的存储,并使用互锁操作来确保其线程安全(很像内置的Subject&lt;T&gt;implementation)。

public class StatelessSubject<T> : ISubject<T>
{
    private IImmutableList<IObserver<T>> _observers
        = ImmutableArray<IObserver<T>>.Empty;

    public void OnNext(T value)
    {
        foreach (var observer in Volatile.Read(ref _observers))
            observer.OnNext(value);
    }
    public void OnError(Exception error)
    {
        foreach (var observer in Volatile.Read(ref _observers))
            observer.OnError(error);
    }
    public void OnCompleted()
    {
        foreach (var observer in Volatile.Read(ref _observers))
            observer.OnCompleted();
    }

    public IDisposable Subscribe(IObserver<T> observer)
    {
        ImmutableInterlocked.Update(ref _observers, x => x.Add(observer));
        return Disposable.Create(() =>
        {
            ImmutableInterlocked.Update(ref _observers, x => x.Remove(observer));
        });
    }
}

现在Publish().RefCount() 可以通过替换它来修复:

.Multicast(new StatelessSubject<SomeType>()).RefCount()

这种变化会产生理想的行为。发布的序列最初是冷的,第一次订阅时变热,最后一个订阅者退订时再次变冷。循环继续,对过去的事件没有任何记忆。

关于源序列完成的另一种正常情况,完成传播到所有订阅者,导致所有订阅者自动取消订阅,导致发布的序列变冷。最终结果是源序列和已发布序列始终保持同步。它们要么都是热的,要么都是冷的。

这里有一个StatelessPublish 操作符,让类的消费更容易一点。

/// <summary>
/// Returns a connectable observable sequence that shares a single subscription to
/// the underlying sequence, without maintaining its state.
/// </summary>
public static IConnectableObservable<TSource> StatelessPublish<TSource>(
    this IObservable<TSource> source)
{
    return source.Multicast(new StatelessSubject<TSource>());
}

使用示例:

.StatelessPublish().RefCount()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-03-20
    • 2019-03-13
    • 2015-08-22
    • 2015-03-12
    相关资源
    最近更新 更多