【问题标题】:Handling backpressure in Rx.NET without onBackpressureLatest在没有 onBackpressureLatest 的情况下处理 Rx.NET 中的背压
【发布时间】:2017-08-23 01:41:59
【问题描述】:

我需要在 Rx.NET 中实现以下算法:

  1. stream 获取最新项目,或者如果没有新项目,则等待新项目而不阻塞。只有最新的项目很重要,其他项目可以丢弃。
  2. 将项目输入SlowFunction并打印输出。
  3. 从第 1 步开始重复。

天真的解决方案是:

let PrintLatestData (stream: IObservable<_>) =
    stream.Select(SlowFunction).Subscribe(printfn "%A")

但是,此解决方案不起作用,因为平均而言 stream 发出项目的速度比 SlowFunction 消耗它们的速度快。由于Select 不会丢弃项目,而是尝试按照从最旧到最新的顺序处理每个项目,因此在程序运行时,发出和打印项目之间的延迟将趋向无穷大。应该只从流中获取最近的最新项目,以避免这种无限增长的背压。

我搜索了文档并在 RxJava 中找到了一个名为 onBackpressureLatest 的方法,据我了解,它可以完成我上面描述的操作。但是,该方法在 Rx.NET 中不存在。如何在 Rx.NET 中实现这一点?

【问题讨论】:

  • SlowFunction 比流慢有什么问题?
  • @FyodorSoikin 如果SlowFunction 比流慢,则新项目的发射速度超过了它们的处理和打印速度。因此,当程序运行时,正在发出的新项目与正在打印的所述项目的SlowFunction 输出之间的延迟/滞后会向无穷大增长。这是不可接受的,因为我需要实时监控数据。我只关心最新的项目。
  • SlowFunction 是否同步?还是 Observable/Async?
  • @Shlomo SlowFunction 是同步的。
  • 还可以查看此线程以获取想法stackoverflow.com/questions/11010602/…

标签: c# f# system.reactive reactive-programming


【解决方案1】:

我认为您想使用 ObserveLatestOn 之类的东西。它有效地将传入事件队列替换为单个值和一个标志。

James World 已在此博客 http://www.zerobugbuild.com/?p=192

这个概念在 GUI 应用程序中被大量使用,这些应用程序无法相信服务器会以多快的速度向其推送数据。

您还可以在 Reactive Trader https://github.com/AdaptiveConsulting/ReactiveTrader/blob/83a6b7f312b9ba9d70327f03d8d326934b379211/src/Adaptive.ReactiveTrader.Shared/Extensions/ObservableExtensions.cs#L64 中看到一个实现 以及解释 ReactiveTrader 的支持演示 https://leecampbell.com/presentations/#ReactConfLondon2014

需要明确的是,这是一种减载算法,而不是背压算法。

【讨论】:

【解决方案2】:

同步/异步建议可能会有所帮助,但是,鉴于慢函数总是比事件流慢,使其异步可能允许您并行处理(在线程池上观察),但代价是(最终) 只是用完线程或通过上下文切换增加更多延迟。这听起来不像是我的解决方案。

我建议您查看 Dave Sexton 编写的开源 Rxx 'Introspective' 运算符。这些可以改变您最近获得的缓冲/节流周期,因为队列由于消费者缓慢而备份。如果慢功能突然变得更快,它根本不会缓冲东西。如果它变慢,它将缓冲更多。 您必须检查是否有“最新来源”类型,或者只是修改现有类型以满足您的需要。例如。使用缓冲区并仅获取缓冲区中的最后一项,或进一步增强以仅在内部存储最新的。谷歌“Rxx”,你会在 Github 的某个地方找到它。

如果“慢功能”的时间相当可预测,一种更简单的方法是简单地将流限制超过此时间的量。显然,我指的不是标准的 rx '节流',而是一种允许更新而不是旧更新的方法。这里有很多解决此类问题的方法。

【讨论】:

    【解决方案3】:

    您可以sample 以您知道SlowFunction 可以处理的时间间隔进行流式传输。这是java中的一个例子:

    TestScheduler ts = new TestScheduler();
    
    Observable<Long> stream = Observable.interval(1, TimeUnit.MILLISECONDS, ts).take(500);
    stream.sample(100, TimeUnit.MILLISECONDS, ts).subscribe(System.out::println);
    
    ts.advanceTimeBy(1000, TimeUnit.MILLISECONDS);
    
    98
    198
    298
    398
    498
    499
    

    sample 不会造成背压,并且总是抓取流中的最新值,因此符合您的要求。另外sample 不会通过相同的值发送两次(从上面可以看出499 只打印一次)

    我认为这将是一个有效的C#/F# 解决方案:

    static IDisposable PrintLatestData<T>(IObservable<T> stream) {
        return stream.Sample(TimeSpan.FromMilliseconds(100))
            .Select(SlowFunction)
            .Subscribe(Console.WriteLine);
    }
    
    let PrintLatestData (stream: IObservable<_>) =
        stream.Sample(TimeSpan.FromMilliseconds(100))
            .Select(SlowFunction)
            .Subscribe(printfn "%A")
    

    【讨论】:

    • 这是一个非常简单的问题解决方案,但在我的特定用例中,SlowFunction 的运行时间变化太大,无法使用这种方法。感谢分享。
    【解决方案4】:

    前段时间我也遇到过同样的问题,但我没有找到可以完全做到这一点的内置运算符。所以我写了自己的,我称之为Latest。实现起来并不简单,但发现它在我当前的项目中非常有用。

    它的工作原理是这样的:当观察者忙于处理先前的通知时(当然是在它自己的线程上),它将最后一个最多 n 个通知(n >= 0)和OnNexts 观察者排队因为它变得空闲。所以:

    • Latest(0): 仅在观察者空闲时观察到达的项目
    • Latest(1):时刻关注最新动态
    • Latest(1000)(例如):通常处理所有项目,但如果有什么事情被卡住了,宁愿错过一些也不要得到一个OutOfMemoryException
    • Latest(int.MaxValue): 不会错过任何一项,而是在生产者和消费者之间进行负载平衡。

    因此,您的代码将是:stream.Latest(1).Select(SlowFunction).Subscribe(printfn "%A")

    签名如下:

    /// <summary>
    /// Avoids backpressure by enqueuing items when the <paramref name="source"/> produces them more rapidly than the observer can process.
    /// </summary>
    /// <param name="source">The source sequence.</param>
    /// <param name="maxQueueSize">Maximum queue size. If the queue gets full, less recent items are discarded from the queue.</param>
    /// <param name="scheduler">Optional, default: <see cref="Scheduler.Default"/>: <see cref="IScheduler"/> on which to observe notifications.</param>
    /// <exception cref="ArgumentNullException"><paramref name="source"/> is null.</exception>
    /// <exception cref="ArgumentOutOfRangeException"><paramref name="maxQueueSize"/> is negative.</exception>
    /// <remarks>
    /// A <paramref name="maxQueueSize"/> of 0 observes items only if the subscriber is ready.
    /// A <paramref name="maxQueueSize"/> of 1 guarantees to observe the last item in the sequence, if any.
    /// To observe the whole source sequence, specify <see cref="int.MaxValue"/>.
    /// </remarks>
    public static IObservable<TSource> Latest<TSource>(this IObservable<TSource> source, int maxQueueSize, IScheduler scheduler = null)
    

    实现太大了,无法在此处发布,但如果有人感兴趣,我很乐意分享。告诉我。

    【讨论】:

    • 不确定是否仍然相关,但我想看看你的实现
    • @tinudu 请分享你的实现。
    猜你喜欢
    • 2023-03-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-10-16
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多