【发布时间】:2017-08-23 01:41:59
【问题描述】:
我需要在 Rx.NET 中实现以下算法:
- 从
stream获取最新项目,或者如果没有新项目,则等待新项目而不阻塞。只有最新的项目很重要,其他项目可以丢弃。 - 将项目输入
SlowFunction并打印输出。 - 从第 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