【问题标题】:Why subscribing an observable in RX.NET by Latest only accepts 1 subscriber?为什么通过 Latest 在 RX.NET 中订阅 observable 只接受 1 个订阅者?
【发布时间】:2016-02-09 18:28:42
【问题描述】:

我的目标是从一个 observable 中获得两个订阅者,但我只对事件流中的最新项目感兴趣。我希望其他人被丢弃。将其视为每 1 秒更新一次且忽略任何中间值的股票价格屏幕。 这是我的代码:

    var ob = Observable.Interval(TimeSpan.FromMilliseconds(100)) // fast event source
        .Latest().ToObservable().ToEvent();

    ob.OnNext += (l =>
                      {
                          Console.WriteLine(Thread.CurrentThread.ManagedThreadId);
                          Thread.Sleep(1000); // slow processing of events
                          Console.WriteLine("Latest: " + l);
                                                    });

    ob.OnNext += (l =>
        {
            Console.WriteLine(Thread.CurrentThread.ManagedThreadId);
            Thread.Sleep(1000); // slow processing of events
            Console.WriteLine("Latest1: " + l);
            //  subject.OnNext(l);
        });

然而,由于上面的代码,尽管我附加了两个事件(即使您使用订阅表示法也没关系)只有第一个订阅被定期调用。第二个根本不运行。为什么会这样?

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    首先我认为您的要求是以下之一:

    1. 您只想获取未来值
    2. 或者您想获取最新值(如果有)和任何未来值
    3. 或者您只需要最近的值(如果有)
    4. 或者您只想对刻度进行采样并获取每秒的值
    5. 或者您的消费者速度较慢,您需要执行减载(例如在 GUI 中)

    1) 的代码

    var ob = Observable.Interval(TimeSpan.FromMilliseconds(1000)) // fast event source
        .Publish();
    ob.Connect();
    

    2) 的代码

    var ob = Observable.Interval(TimeSpan.FromMilliseconds(1000)) // fast event source
        .Replay(1);
    ob.Connect();    
    

    3) 的代码

    var ob = Observable.Interval(TimeSpan.FromMilliseconds(1000)) // fast event source
        .Replay(1);
    ob.Connect();  
    var latest = ob.Take(1); 
    

    4) 的代码可以是这样,但在您认为的窗口周围有一些微妙的行为。

    var ob = Observable.Interval(TimeSpan.FromMilliseconds(200)) // fast event source
        .Replay(1);
    //Connect the hot observable
    ob.Connect();
    
    var bufferedSource = ob.Buffer(TimeSpan.FromSeconds(1))
        .Where(buffer => buffer.Any())
        .Select(buffer => buffer.Last());  
    

    5) 的代码可以在 James World 的博客 http://www.zerobugbuild.com/?p=192 上找到,并且在伦敦的许多银行应用程序中很常见。

    【讨论】:

    • 感谢您的回答。不过,这不是我的要求。没有 1 秒的缓冲要求。我想在订阅者代码结束后立即获得最新值。我发布的代码实现了这一点!但是,仅适用于单个订阅者。我试图弄清楚为什么它不会调用第二个。我已经阅读了互联网上所有相关的帖子(包括你的网站和帖子)。除了无法调用第二个订阅者之外,上述方法对我来说看起来更好。
    • 基本上我想在订阅者代码运行时忽略中间值。以上都没有实现。
    • James 的博文确实做到了这一点。您正在使用 Thread.Sleep 阻塞线程。我相信这会阻止您的其他订阅运行。使用 Rx 后,建议您避免使用 Thread.Sleeps
    • 是的,不过我同意,我放 thread.Sleeps 只是为了模拟一些工作。
    • 你试过ObserveLatestOn自定义操作符吗?
    【解决方案2】:

    我认为你不明白 .Latest() 的作用。

    public static IEnumerable<TSource> Latest<TSource>(
        this IObservable<TSource> source
    )
    

    在每次迭代时返回最后一个采样元素并随后阻塞直到可观察源序列中的下一个元素可用的可枚举序列。

    请注意,它会阻塞等待来自 observable 的下一个元素。

    因此,当您使用.Latest()IObservable&lt;&gt; 转换为IEnumerable&lt;&gt; 时,您必须使用.ToObservable() 将其转换回IObservable&lt;&gt; 才能调用.ToEvent()。那就是它倒下的地方。

    此代码的问题在于您创建的代码会阻塞。

    如果你这样做,它会起作用:

    var ob = Observable.Interval(TimeSpan.FromMilliseconds(100)).ToEvent();
    

    没有必要打电话给.Latest(),因为您总是从可观察对象中获取最新值。您永远无法获得更早的值。它是可观察的,而不是时间机器。

    我也不明白你为什么打电话给.ToEvent()。有什么需求?

    这样做:

    var ob = Observable.Interval(TimeSpan.FromMilliseconds(100));
    
    ob.Subscribe(l =>
    {
        Console.WriteLine(Thread.CurrentThread.ManagedThreadId);
        Thread.Sleep(1000); // slow processing of events
        Console.WriteLine("Latest: " + l);
    });
    
    ob.Subscribe(l =>
    {
        Console.WriteLine(Thread.CurrentThread.ManagedThreadId);
        Thread.Sleep(1000); // slow processing of events
        Console.WriteLine("Latest1: " + l);
    });
    

    【讨论】:

    • 让我解释一下:如果我使用Latest,那么在订阅执行的睡眠期间生成的很多值都会被跳过。所以我的输出变成了 9、18、27。这就是我想要的。如果我不使用最新,那么我会得到像 1,2,3 这样的值,并且没有任何项目被跳过。这第二种情况不是我想要的。至于 ToEvent,没有必要,我把它们展示了它如何与 C# 事件语法不一致。您通过 += 运算符附加两个事件处理程序,并且仅调用其中一个。这种行为不一致。
    • @ReverseBlade - 就在最后一点 - 行为并不一致。 .Latest().ToObservable().ToEvent() 链创建了一个阻塞运算符,它一次只返回一个输出并阻塞其余时间。第二个OnNext 甚至没有机会获得价值。这并不矛盾。这就是这些运营商如何协同工作。您应该避免像这样混合可枚举和可观察对象。
    • 啊,是的。其实我想通了。不幸的是,这是我想做的。我找不到更好的方法来实现这一点。
    • @ReverseBlade - 我仍在考虑如何解决您的问题。我会看看我能想出什么。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-12-03
    • 2018-01-30
    • 1970-01-01
    • 1970-01-01
    • 2019-12-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多