【发布时间】:2020-11-23 01:02:13
【问题描述】:
最近我偶然发现了 Enigmativity 的 interesting statement 关于 Publish 和 RefCount 运算符:
您正在使用危险的 .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
如果我们注释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 或两者,以使它们的行为始终一致且合乎需要?以下是理想行为的规范:
- 发布的序列应该将所有直接来自源序列的通知传播给它的订阅者,而不是其他任何东西。
- 当当前订阅者数量从零增加到一时,已发布序列应订阅源序列。
- 只要至少有一个订阅者,发布的序列就应该与源保持连接。
- 当当前订阅者数量为零时,已发布序列应取消订阅源。
我要求提供上述功能的自定义 PublishRefCount 运算符,或使用内置运算符实现所需功能的方法。
顺便说一句,similar question 存在,它问为什么会发生这种情况。我的问题是关于如何解决它。
更新:回想起来,上述规范导致了一种不稳定的行为,使得竞态条件不可避免。不能保证对已发布序列的两次订阅将导致对源序列的一次订阅。源序列可能在两个订阅之间完成,导致第一个订阅者取消订阅,导致RefCount 运算符取消订阅,导致下一个订阅者重新订阅源。内置 .Publish().RefCount() 的行为可以防止这种情况发生。
道德教训是.Publish().RefCount() 序列没有损坏,但它不可重复使用。它不能可靠地用于多个连接/断开会话。如果你想要第二个会话,你应该创建一个新的.Publish().RefCount() 序列。
【问题讨论】:
-
为什么你认为你观察到的行为是错误的?这正是我期望它的行为方式。
-
@Shlomo 在
RefCount的definition 中没有提示它保持源observable 的状态(无论是否完成),并使用此内存发送OnCompleted代表其发出通知。我的期望是它只保持当前的订阅者数量。换句话说,RefCount做的比它应该做的要多。我宁愿它是愚蠢和可预测的,而不是有自己的想法。 -
它没有这样做。发布是。
-
@Shlomo 那太好了。在这种情况下,只有
Publish需要修复! -
@Enigmativity 啊,那好吧。很抱歉造成误解。
标签: c# system.reactive rx.net