【发布时间】: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");
}
}
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