【问题标题】:Lazy observable sequence that replays value or error重放值或错误的惰性可观察序列
【发布时间】: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


【解决方案1】:

试试这个版本的你ConnectUntilCompleted方法:

public static IObservable<T> ConnectUntilCompleted<T>(this IConnectableObservable<T> @this)
{
    return Observable.Create<T>(o =>
    {
        var subscription = @this.Subscribe(o);
        var connection = @this.Connect();
        return new CompositeDisposable(subscription, connection);
    });
}

允许 Rx 正常运行。

现在我已经添加到它以帮助显示正在发生的事情:

public static IObservable<T> ConnectUntilCompleted<T>(this IConnectableObservable<T> @this)
{
    return Observable.Create<T>(o =>
    {
        var disposed = Disposable.Create(() => Console.WriteLine("Disposed!"));
        var subscription = Observable
            .Defer<T>(() => { Console.WriteLine("Subscribing!"); return @this; })
            .Subscribe(o);
        Console.WriteLine("Connecting!");
        var connection = @this.Connect();
        return new CompositeDisposable(disposed, subscription, connection);
    });
}

现在你的 observable 看起来像这样:

var o =
    Observable
        .Defer(() =>
            Observable
                .Start(() =>
                {
                    Console.WriteLine("Started.");
                    throw new InvalidOperationException();
                }))
        .PublishLast()
        .ConnectUntilCompleted();

最后的关键是实际处理订阅中的错误——所以仅仅做o.Subscribe()是不够的。

这样做:

        o.Subscribe(
            x => Console.WriteLine(x),
            e => Console.WriteLine(e.Message),
            () =>  Console.WriteLine("Done."));

        o.Subscribe(
            x => Console.WriteLine(x),
            e => Console.WriteLine(e.Message),
            () =>  Console.WriteLine("Done."));

        o.Subscribe(
            x => Console.WriteLine(x),
            e => Console.WriteLine(e.Message),
            () =>  Console.WriteLine("Done."));         

当我运行时,我得到了这个:

订阅! 连接! 订阅! 连接! 订阅! 连接! 开始了。 由于对象的当前状态,操作无效。 处置! 由于对象的当前状态,操作无效。 处置! 由于对象的当前状态,操作无效。 处置!

注意“Started”只出现一次,但报错3次。

(有时Started 在第一次订阅后出现在列表中较高的位置。)

我想这就是你想要的描述。

【讨论】:

  • 感谢@Enigmativity 的有用回复。不幸的是,它并没有完全表现出我所追求的行为。我认为ConnectUntilCompleted 是一个糟糕的名字。我想要的是任何未来的订阅者都能收到相同的结果,即使所有以前的订阅者都已经处理了他们的订阅。也许这实际上是一个坏主意(它确实感觉不合时宜),但我正在尝试对异步方法调用进行建模。所以 obs 代表该方法调用的结果,任何未来的订阅者都应该收到相同的结果。如果需要一个新的调用,就需要一个新的 obs。
  • @KentBoogaart - 所有未来的订阅者都会收到相同的结果。如果您将throw 替换为return 42,那么您将看到"Started." 只发生一次,但三个订阅都获得相同的值。这确实发生在之前的调用者处理了他们的订阅之后。
  • 我想我对以下事实感到困惑:有时Subscribe() 会立即抛出(或者当我提前调度程序时),而其他时候它不会(并且需要一个明确的错误处理程序来“请参阅“错误)。
  • @KentBoogaart - .Subscribe(...) 不会抛出。可观察的管道抛出。如果您的订阅者中有错误处理程序,那么您可以干净地处理错误。如果没有,那么管道应该抛出。按照惯例,您应该编写不会抛出异常的代码,除非无法避免。
【解决方案2】:

只是为了支持@Engimativity 的回答,我想展示你应该如何运行测试,这样你就不会得到这些“惊喜”。您的测试是不确定的,因为它们是多线程/并发的。您在不提供IScheduler 的情况下使用Observable.Start 是有问题的。如果您使用 TestScheduler 运行测试,您的测试现在将是单线程和确定性的

[Test]
public void Test()
{
    var testScheduler = new TestScheduler();
    var o = Observable
        .Defer(() => Observable
            .Start(() => { throw new InvalidOperationException(); }, testScheduler)
            .PublishLast()
            .ConnectUntilCompleted());

    var observer = testScheduler.CreateObserver<Unit>();
    o.Subscribe(observer);

    testScheduler.Start();

    CollectionAssert.IsNotEmpty(observer.Messages);
    Assert.AreEqual(NotificationKind.OnError, observer.Messages[0].Value.Kind);
}

【讨论】:

  • 我完全忘记了Start 默认不使用ImmediateScheduler。我真的需要编写那些 Roslyn 分析器,以便在您调用采用调度程序的方法并且您不提供调度程序时发出警告。谢谢李。也就是说,我试图构建的 observable 的目标是,如果它已经完成,它不应该重要 - 你应该在订阅时获得值(或错误)。
  • 是的,后续订阅者也是如此(如果原始订阅在此之前发生了序列错误)
【解决方案3】:

实现您的要求的另一种方法是:

var lazy = new Lazy<Task>(async () => { /* execute once */ }, isThreadSafe: true);
var o = Observable.FromAsync(() => lazy.Value);

第一次订阅时,lazy 将创建(并执行)任务。对于其他订阅,lazy 将返回相同的(可能已经完成或失败)任务。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-08-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多