【问题标题】:How to merge observables on a regular interval?如何定期合并 observables?
【发布时间】:2015-01-29 18:50:50
【问题描述】:

我正在尝试定期合并两个传感器数据流,但在 Rx 中无法正确执行此操作。我想出的最好的是下面的示例,但是我怀疑这是对 Rx 的最佳使用。

有没有更好的方法?

我尝试过 Sample(),但传感器以不规则的时间间隔产生值,慢(> 1 秒)和快(

Observable<SensorA> sensorA = ... /* hot */
Observable<SensorB> sensorB = ... /* hot */

SensorA lastKnownSensorA;
SensorB lastKnownSensorB;

sensorA.Subscribe(s => lastKnownSensorA = s);
sensorB.Subscribe(s => lastKnownSensorB = s);

var combined = Observable.Interval(TimeSpan.FromSeconds(1))
    .Where(t => _lastKnownSensorA != null)
    .Select(t => new SensorAB(lastKnownSensorA, lastKnownSensorB)

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    我认为@JonasChapuis 的答案可能是您所追求的,但有几个问题可能存在问题:

    • CombineLatest 在所有源都至少发出一个值之前不会发出值,这可能会导致在此之前来自更快源的数据丢失。这可以通过使用StartWith 在每个传感器流上播种null object 或默认值来缓解。

    • Sample 如果在采样期间没有观察到新值,则不会发出值。我无法从问题中判断这是否可取,但如果不是,则有一个有趣的技巧可以使用“节奏”流来解决这个问题,如下所述以创建 固定 频率,而不是使用Sample 获得的最大频率。

    要解决 CombineLatest 问题,请为您的传感器流确定适当的空值 - 我通常通过类型上的静态 Null 属性使这些可用 - 这使意图非常明确。对于值类型,使用Nullable&lt;T&gt; 也是一个不错的选择:

    Observable<SensorA> sensorA = ...  .StartWith(SensorA.Null);
    Observable<SensorB> sensorB = ...  .StartWith(SensorB.Null);
    

    注意不要犯将StartWith 仅应用于CombinedLatest 的输出的常见错误......这无济于事!

    现在,如果您需要定期结果(自然可能包括最近读数的重复),请创建一个以所需时间间隔发出的“节奏”流:

    var pace = Observable.Interval(TimeSpan.FromSeconds(1));
    

    然后组合如下,从结果中省略速度值:

    var sensorReadings = Observable.CombineLatest(
        pace, sensorA, sensorB,
        (_, a, b) => new SensorAB(a,b));
    

    还值得了解MostRecent 运算符,如果您想以特定流的速度驱动输出,它可以非常有效地与Zip 结合使用。请参阅我演示该方法的这些答案:How to combine a slow moving observable with the most recent value of a fast moving observable 以及处理多个流的更有趣的调整:How do I combine three observables such that

    【讨论】:

    • 我真的很喜欢速度流的想法,但是在这种情况下,我仍然需要对结果进行采样,因为传感器可以突发数据(每秒 10-50 个),这会导致很多不必要的分配。我不认为 MostRecent+Zip 在这种情况下是实用的,因为任何一个传感器都可能爆裂或有时非常慢(虽然很酷)
    • 好的。通过要求最少的分配,你引入了一个新的要求,它改变了很多事情,这实际上是一个不同的问题,重点需要放在其他地方。我很好奇你为什么需要这个(你的系统是否会因为 GC 压力而表现出不良影响?)。这里需要更多的上下文来提供一个好的答案,特别是了解传感器数据流的性质。避免分配通常需要仔细考虑数据类型的使用和重用。
    • 在两个传感器流上应用Sample 来限制源数据当然是限制下游分配的一种明显方法——但肯定不是全部。从本质上讲,您会想要尝试设计端到端流,以便大多数分配都在源点,并通过重用来提高效率。不过,担心所有基于实际而不是想象的性能问题的事情!这种代码需要高超的技能,维护起来也很昂贵。
    • 如果没有必要,为什么要分配大量对象?我上面的例子并没有过度分配。我的问题是我的代码对我来说看起来并不像正确的 RX,但是我找不到没有负面影响的方法。或者无论如何,这会产生正确的结果。就像我说的我不能使用 Sample() 因为当流很慢时它不会触发。所以即使是 combinelatest + 起搏器也不起作用。
    • 如果您将Sample() 附加到上面的sensorAsensorB 声明中,在StartWith 之前,各个流的速度将无关紧要。
    【解决方案2】:

    如何使用CombineLatest() 运算符在每次产生一个值时合并传感器的最新值,然后使用Sample() 以确保每秒一次测量的最大频率?

    sensorA.CombineLatest(sensorB, (a, b) => new {A=a, B=b}).Sample(TimeSpan.FromSeconds(1))
    

    【讨论】:

    • 试过这样的东西。我发现了 2 个问题。如果两个传感器都在爆裂(每秒 10+ 个滴答声),我最终会从 CombineLatest 获得很多不必要的分配。其次,如果 Sample() 超过一秒钟没有收到任何内容,则不会触发。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-06-15
    • 1970-01-01
    • 1970-01-01
    • 2020-07-23
    相关资源
    最近更新 更多