【问题标题】:Use Reactive Extensions to filter on "significant changes" in observable stream使用响应式扩展过滤可观察流中的“重大变化”
【发布时间】:2016-09-07 19:22:45
【问题描述】:

我有一个GeoLocationProvider(它实现了IObservable<System.Device.Location.GeoCoordinate>,它每 x 毫秒输出一次当前位置。

现在我想使用 RX 读取所有这些 GPS 坐标,并且仅在位置变化重大时(例如 - 行驶距离 > 10 米)通知订阅者。

我的“主要”代码

// [...]
IObservable<GeoCoordinate> locationProvider = new GeoLocationProvider();
LocationFeed locationFeed = new LocationFeed(locationProvider);

// register any interested observers on the locationFeed.
ConsoleLocationReporter c1 = new ConsoleLocationReporter("reporter0001");
locationFeed.Subscribe(c1);

我的 LocationFeed 实现如下所示:

using System;
using System.Device.Location;
using System.Reactive.Subjects;

namespace My.Namespace.Movement
{
    public class LocationFeed : ISubject<GeoCoordinate>, IDisposable
    {
        private readonly IDisposable _subscription;
        private readonly Subject<GeoCoordinate> _subject;

        public LocationFeed(IObservable<GeoCoordinate> observableSource)
        {
            _subject = new Subject<GeoCoordinate>();
            _subscription = observableSource.Subscribe(_subject); // TODO: Add logic to filter to only significant movement changes (> 10m)
        }

        public void Dispose()
        {
            _subscription?.Dispose();
            _subject?.Dispose();
        }

        public void OnNext(GeoCoordinate value)
        {
            _subject.OnNext(value);
        }

        public void OnError(Exception error)
        {
            _subject.OnError(error);
        }

        public void OnCompleted()
        {
            _subject.OnCompleted();
        }

        public IDisposable Subscribe(IObserver<GeoCoordinate> observer)
        {
            return _subject.Subscribe(observer);
        }
    }
}

问题一: GeoCoordinate 提供了一个方法c1.DistanceTo(c2) 来计算两个坐标之间的距离。如果阈值与上次推送的阈值相比大于 x,我只想报告(发布)新的地理坐标。我如何做到这一点?

问题2:主题的使用是否正常,以及我实现ISubject的方式?我不想在我的“主”代码中添加所有连线并将其全部移到单独的类中。

【问题讨论】:

  • 如果您想获得完整的答案,我很乐意提供进一步的帮助,但您能否提供测试以显示您想要的内容以及缺少的 GeoCoordinate
  • 我同意 Lee 的观点——不要实现你自己的实现 Rx 接口的类。这样做只会带来坏事。

标签: c# system.reactive observable


【解决方案1】:

我强烈建议不要实现ISubject&lt;T&gt;(或者就此而言IObservable&lt;T&gt;IObserver&lt;T&gt;)。而是尝试组合现有的工厂和类型,然后将它们公开为“具有”关系而不是“是”关系。

如您所见,您的LocationFeed 纯粹是observableSource 参数的包装,因此似乎无法解决任何问题。我建议删除它。

关于您发布的问题,一种解决方案是使用大小为 2 且步长为 1 的缓冲区。

IObservable<GeoCoordinate> locationProvider = new GeoLocationProvider();

locationProvider
    .Buffer(2,1)
    .Where(buffer=>buffer[0].DistanceTo(buffer[1]) > 10)
    .Select(buffer=>buffer[1])
    .Subscribe(
        pos => Console.WriteLine(pos),
        ex => { },
        () => {});

或者你可以使用Scan

IObservable<GeoCoordinate> locationProvider = new GeoLocationProvider();

locationProvider
    .Scan(Tuple.Create(GeoCoordinate.Zero,GeoCoordinate.Zero), (acc, cur)=>Tuple.Create(acc.Item2, cur))
    .Where(pair=>pair.Item1.DistanceTo(pair.Item2) > 10)
    .Select(pair=>pair.Item2)
    .Subscribe(
        pos => Console.WriteLine(pos),
        ex => { },
        () => {});

我不确定您对产生的第一个值有什么要求。是否应该发布?

编辑: 这是一个经过测试的解决方案(使用Point 类型),当发生“重大”更改时,它将推送Unit。如果这不是您想要的,那么您应该可以通过修补来获得您真正想要的

void Main()
{
    var zero = new System.Drawing.Point(0,0);
    var fenceDistance = 10;

    var scheduler = new TestScheduler();
    var source = scheduler.CreateColdObservable(
        ReactiveTest.OnNext(1, new System.Drawing.Point(0,0)),
        ReactiveTest.OnNext(2, new System.Drawing.Point(0,9)),  //Not far enough
        ReactiveTest.OnNext(3, new System.Drawing.Point(0,10)), //Touches the fence
        ReactiveTest.OnNext(4, new System.Drawing.Point(0,15)), //Not far enough
        ReactiveTest.OnNext(5, new System.Drawing.Point(0,40))  //Breaches the fence        
        );

    var observer = scheduler.CreateObserver<Unit>();

    source
        .Scan(Tuple.Create(zero, zero), (acc, cur) =>
        {
            if (DistanceBetween(acc.Item1, cur) >= fenceDistance)
            {
                return Tuple.Create(cur, cur);
            }
            else
            {
                return Tuple.Create(acc.Item1, cur);
            }
        })
        .Where(pair => pair.Item1 == pair.Item2)
        .Select(pair => Unit.Default)
        .Subscribe(observer);


    scheduler.Start();

    ReactiveAssert.AreElementsEqual(new[] {
        ReactiveTest.OnNext(1, Unit.Default),
        ReactiveTest.OnNext(3, Unit.Default),
        ReactiveTest.OnNext(5, Unit.Default)
    },observer.Messages);

}

// Define other methods and classes here
public static double DistanceBetween(System.Drawing.Point a, System.Drawing.Point b)
{
    var xDelta = a.X -b.X;
    var yDelta = a.Y - b.Y;

    var distanceSqr = (xDelta * xDelta) + (yDelta * yDelta);
    return Math.Sqrt(distanceSqr);
}

【讨论】:

  • 您可能需要调整扫描解决方案以匹配“如果阈值与最后一个推送的阈值相比”的要求,即仅当新值距离“acc”足够远时更新“acc” '。否则走路的人可能永远不会产生变化。
  • 这就是我最初认为的要求,但我将其读作“仅当两个后续值相隔大于 x 时才屈服”。但是,如果不是这种情况,那么确定这里的两个 sln 都需要更改。如果操作有 MCVE,那么测试会告诉我们是否正确:-)
  • 第一个值总是需要屈服,因为它随后被用作参考点。以下所有数字将与参考值进行比较。产生的下一个数字应该是距原始参考点的“距离> 10”。刚刚产生的值随后将成为下一个数字的新参考点。
  • 关于 ISubject 实现。是的,我完全同意。正如我所指出的,我唯一的动机是将所有这些逻辑封装在一个类中,这样我就可以一遍又一遍地重用相同的实现。我担心只有“本地范围”访问权限。
  • 知道了。将找到一些时间来更新这些要求的答案
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-07-31
  • 1970-01-01
  • 2014-11-15
  • 2012-12-19
  • 1970-01-01
相关资源
最近更新 更多