【问题标题】:System.reactive - dynamically switch by frequencySystem.reactive - 按频率动态切换
【发布时间】:2019-07-04 10:29:18
【问题描述】:

我有两个数据来源。

让我们想象一下:

  • 系统 A 以更高的频率提供更高质量的数据,例如
    1price/1sec,但有时会出现故障且没有数据或
    频率例如 1price/20sec
  • 系统 B 提供频率较低的数据,例如1 个价格/10 秒

是否有任何优雅的方式使用 system.reactive 从系统 A 中正常检索数据,但是当它失败(Feed 中没有数据)或速度变慢时,使用系统 B 中的数据? 我想实现某种开关,当它比 B 快时将使用 A 源。我不想混合源,所以我一次只能使用 SystemA 或 SystemB。


    class PriceFeed {

        public IObservable<Price> GetPricesFeed(IObservable<PriceFromA> pricesFromA, IObservable<PriceFromB> pricesFromB)
        {


        }


        private Price Convert(PriceFromA price) { //convert }

        private Price Convert(PriceFromB price) { //convert }

    }

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    有趣的问题。首先要做的是编写某种频率收集函数。可能是这样的:

    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

    【讨论】:

    • 不错。用 CombineLatest() 代替 Zip() 怎么样?我认为另一个方面是应该比较慢的频率下降更早检测到超时
    • CombineLatest 将引入竞争条件问题以组合两个频率可观察量。由于它们本质上是在同一个计时器上,所以我不明白这一点。
    • 关于超时,如果您将测量频率设置为小于超时持续时间,它将比超时检测系统更快,而不是更慢。
    • 更新:我意识到GetFrequency 可以使用Buffer 而不是更简单的Scan。我相应地更改了代码。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-20
    • 2020-04-19
    相关资源
    最近更新 更多