【发布时间】:2019-01-26 04:06:25
【问题描述】:
我想要做的声明如下所示:
// Checks input source for timeouts, based on the number of elements received
// from clock since the last one received from source.
// The two selectors are used to generate output elements.
public static IObservable<R> TimeoutDetector<T1,T2,R>(
this IObservable<T1> source,
IObservable<T2> clock,
int countForTimeout,
Func<R> timedOutSelector,
Func<T1, R> okSelector)
大理石图在 ascii 中很困难,但这里是:
source --o---o--o-o----o-------------------o---
clock ----x---x---x---x---x---x---x---x---x---
output --^---^--^-^----^-----------!-------^---
我尝试寻找现有的 Observable 函数,它们可以以我可以使用的方式组合 source 和 clock,但大多数组合函数依赖于接收“每个函数之一”(And , Zip),或者他们从“缺失”的值 (CombineLatest) 中重新返回“上一个”值,或者它们离我需要的值太远了 (Amb, GroupJoin, @987654331 @、Merge、SelectMany、Timeout)。 Sample 看起来很接近,但我不想将源吞吐量限制为时钟速率。
所以现在我被困在试图填补这里的巨大空白:
return new AnonymousObservable<R>(observer =>
{
//One observer, two observables??
});
抱歉,“您尝试过什么”部分在这里有点弱:假设我已经尝试过思考它!我不是要求完整的实施,只是:
- 是否有内置功能可以帮助我,但我错过了?
- 如何构建一个基于 lambda 的观察者,它订阅两个可观察对象?
【问题讨论】:
-
给未来的读者;在您的问题中,我假设您的大理石图是为
countForTimeout参数传递 3 的结果? -
@Lee,是的,没错。
标签: c# system.reactive