【发布时间】:2016-02-11 10:35:01
【问题描述】:
我正在尝试创建具有以下特征的可观察管道:
- 很懒(在有人订阅之前什么都不做)
- 无论收到多少订阅,最多执行一次
- 重放其结果值,如果有的话
- 重放其产生的错误(如果有)
在我的一生中,我无法找出正确的语义来实现这一点。我认为这将是一个简单的例子:
Observable
.Defer(() => Observable
.Start(() => { /* do something */ })
.PublishLast()
.ConnectUntilCompleted());
ConnectUntilCompleted 就像它听起来的那样:
public static IObservable<T> ConnectUntilCompleted<T>(this IConnectableObservable<T> @this)
{
@this.Connect();
return @this;
}
这似乎在可观察对象成功终止时有效,但在出现错误时无效。任何订阅者都不会收到错误消息:
[Fact]
public void test()
{
var o = Observable
.Defer(() => Observable
.Start(() => { throw new InvalidOperationException(); })
.PublishLast()
.ConnectUntilCompleted());
// this does not throw!
o.Subscribe();
}
谁能告诉我我做错了什么?为什么Publish 不重播它收到的任何错误?
更新:它变得更加陌生:
[Fact]
public void test()
{
var o = Observable
.Defer(() => Observable
.Start(() => { throw new InvalidOperationException(); })
.PublishLast()
.ConnectUntilCompleted())
.Do(
_ => { },
ex => { /* this executes */ });
// this does not throw!
o.Subscribe();
o.Subscribe(
_ => { },
ex => { /* even though this executes */ });
}
【问题讨论】:
-
一旦您放弃订阅结果,您的订阅者就会停止接收通知。试着抓住它们一会儿。
-
你对 observables 的理解似乎有些问题。他们都是懒惰的——甚至是热的——并且在有订阅者之前什么都不做。每个 observable 都是一个定义。当您订阅时,您会为每个订阅创建一个新管道。管道彼此独立(除非您使用发布或主题来共享管道的一部分)。因此,您无法重播错误 - 一旦发生错误,已创建的管道就会关闭。您可以重复值,因为合同是
OnNext*(OnError|OnCompleted)- 所以多个值,一个错误。 -
啊,但我想我明白你想用
.PublishLast()做什么。它将为每个未来的订阅保留。我想我有一些代码可以帮助你。 -
我认为这里有竞争条件。如果我在没有附加调试器的情况下多次运行它,我可以在几次尝试后让它抛出。
-
@ErenErsönmez 是正确的。当您以多线程方式执行测试时(Obs.Start 使用 TaskPoolScheduler 的默认值),您的测试在抛出异常之前完成。如果您将 TestScheduler 放入测试中,您会看到抛出的异常。
标签: c# .net system.reactive