Rx:运算符,用于从Observable流中获取第一个和最近的值

Dam*_*ian 3 system.reactive

对于基于Rx的变更跟踪解决方案,我需要一个可以在可观察序列中获取第一个和最新项目的运算符.

我如何编写一个生成以下大理石图的Rx运算符(注意:括号仅用于排列项目...我不确定如何最好地在文本中表示这一点):

     xs:---[a  ]---[b  ]-----[c  ]-----[d  ]---------|
desired:---[a,a]---[a,b]-----[a,c]-----[a,d]---------| 
Run Code Online (Sandbox Code Playgroud)

yam*_*men 5

使用与@Wilka相同的命名,您可以使用以下扩展名,这有点不言自明:

public static IObservable<TResult> FirstAndLatest<T, TResult>(this IObservable<T> source, Func<T,T,TResult> func)
{
    var published = source.Publish().RefCount();
    var first = published.Take(1);        
    return first.CombineLatest(published, func);
}
Run Code Online (Sandbox Code Playgroud)

请注意,它不一定返回a Tuple,而是为您提供在结果上传递选择器函数的选项.这使它与底层的主要操作(CombineLatest)保持一致.这显然很容易改变.

用法(如果您想在结果流中使用元组):

Observable.Interval(TimeSpan.FromSeconds(0.1))
          .FirstAndLatest((a,b) => Tuple.Create(a,b))
          .Subscribe(Console.WriteLine);
Run Code Online (Sandbox Code Playgroud)

  • 您应该添加一个Publish运算符来共享订阅源的副作用.上面的FirstAndLatest实现将导致对其结果的每个订阅的两个订阅源,这可能导致大量重复计算(或者更糟糕的是,诸如启动I/O等诸如此类的副作用). (3认同)