有趣的问题。首先要做的是编写某种频率收集函数。可能是这样的:
public static IObservable<int> GetFrequency<T>(this IObservable<T> source, TimeSpan measuringFreq, TimeSpan lookback)
{
return source.GetFrequency(measuringFreq, lookback, Scheduler.Default);
}
public static IObservable<int> GetFrequency<T>(this IObservable<T> source, TimeSpan measuringFreq, TimeSpan lookback, IScheduler scheduler)
{
return source.Buffer(lookback, measuringFreq, scheduler)
.Select(l => l.Count);
}
如果measuringFreq 是 1 秒,lookback 是 5 秒,这意味着每一秒我们都会看到过去 5 秒内发送了多少条消息的计数。快速而肮脏的例子:
var r = new System.Random();
var nums = Observable.Generate(
0,
i => i < 100,
i => i + 1,
i => i, _ => TimeSpan.FromSeconds(r.NextDouble() * 1)
);
var freq = nums.GetFrequency(TimeSpan.FromSeconds(1), TimeSpan.FromSeconds(5));
freq.Dump(); //Linqpad
nums 是一个 observable,它应该平均每半秒生成一条消息(它随机选择 0 到 1 秒之间的持续时间)。 freq 每秒生成一个值,该值返回过去 5 秒内生成的消息 nums 的数量(平均应为 10)。在我机器上的最新运行中,我得到了这个:
11
11
12
10
12
11
9
9
10
9
8
...
一旦我们有了获得频率的方法,您就需要编写一个函数来将两个类似类型的 observables 合成在一起,并根据频率进行切换。我是这样写的:
public static IObservable<T> MaintainFrequencyImproper<T>(this IObservable<T> sourceA, IObservable<T> sourceB, TimeSpan measuringFreq, TimeSpan lookback, IScheduler scheduler, int aAdvantage = 0)
{
var aFreq = sourceA.GetFrequency(measuringFreq, lookback, scheduler);
var bFreq = sourceB.GetFrequency(measuringFreq, lookback, scheduler);
var toReturn = aFreq.Zip(bFreq, (a, b) => a + aAdvantage - b)
.Select(freqDifference => freqDifference < 0 ? sourceB : sourceA) //If advantage is 0, and a & b both popped out 5 messages in the last second, then A wins
.StartWith(sourceA)
.Switch();
return toReturn;
}
首先我们用GetFrequency 获得两个可观察的频率,然后我们将这两个压缩在一起,并比较它们。如果 B 比 A 更频繁,则使用 B。如果它们的频率相等或 A 更频繁,则使用 A。
aAdvantage 变量允许您表达对 A 比 B 更强的偏好。0(默认值)意味着源 A 赢得平局,或者当它更频繁时,但 B 获胜。 2 意味着 B 在最近一段时间内必须比 A 多产生 3 条消息才能要求使用 B。
使用适当的 Publishing 可观察对象以避免多次订阅,看起来像这样:
public static IObservable<T> MaintainFrequencyProper<T>(this IObservable<T> sourceA, IObservable<T> sourceB, TimeSpan measuringFreq, TimeSpan lookback,
IScheduler scheduler, int aAdvantage = 0)
{
return sourceA.Publish(_sourceA => sourceB.Publish(_sourceB =>
_sourceA.GetFrequency(measuringFreq, lookback, scheduler)
.Zip(_sourceB.GetFrequency(measuringFreq, lookback, scheduler), (a, b) => a + aAdvantage - b)
.Select(freqDifference => freqDifference < 0 ? _sourceB : _sourceA)
.StartWith(_sourceA)
.Switch()
))
}
我希望这会有所帮助。在如何将其融入您的代码方面,您并没有留下太多东西。如果您需要,请添加mcve。