使用RX查询,如何获得每秒3秒的窗口具有相同状态的记录?

Lau*_*ura 7 .net c# system.reactive

我有几天看着RX,我已经阅读了很多; 我读过IntroToRx; 我还看了101个RX样品和许多其他地方,但我无法弄清楚这一点.这听起来很简单,但我无法得到我需要的东西:我需要知道哪个"ID"已经"陷入"状态'STARTED'至少30分钟.

我有一个MyInfo类,如下所示:

public class MyInfo
{
    public string ID { get; set; }
    public string Status { get; set; }
}
Run Code Online (Sandbox Code Playgroud)

我编写了一个主题来帮助我测试这样:

        var subject = new Subject<MyInfo>();
        subject.OnNext(new MyInfo() { ID = "1", Status = "STARTED" });
        subject.OnNext(new MyInfo() { ID = "2", Status = "PHASE1" });
        subject.OnNext(new MyInfo() { ID = "3", Status = "STOPPED" });
        subject.OnNext(new MyInfo() { ID = "4", Status = "STARTED" });
        subject.OnNext(new MyInfo() { ID = "1", Status = "STARTED" });
        subject.OnNext(new MyInfo() { ID = "2", Status = "PHASE1" });
        subject.OnNext(new MyInfo() { ID = "3", Status = "STOPPED" });
        subject.OnNext(new MyInfo() { ID = "4", Status = "PHASE2" });
        subject.OnNext(new MyInfo() { ID = "1", Status = "STARTED" });
        subject.OnNext(new MyInfo() { ID = "2", Status = "STOPPED" });
        subject.OnNext(new MyInfo() { ID = "3", Status = "STOPPED" });
        subject.OnNext(new MyInfo() { ID = "4", Status = "STARTED" });
        subject.OnCompleted();
Run Code Online (Sandbox Code Playgroud)

到目前为止,我的查询和订阅看起来像这样:(我在样本中使用秒数)

var q8 = from e in subject
                 group e by new { ID = e.ID, Status = e.Status } into g
                 from w in g.Buffer(timeSpan: TimeSpan.FromSeconds(3)
                                   , timeShift: TimeSpan.FromSeconds(1))
                 select new
                 {
                     ID = g.Key.ID,
                     Status = g.Key.Status,
                     count = w.Count
                 };
        var subsc = q8.Subscribe(a => Console.WriteLine("{0} {1} {2}", a.ID, a.Status, a.count));
Run Code Online (Sandbox Code Playgroud)

现在我可以得到一个输出,告诉我ID以及ID在一段时间内看到的状态.

ID   Status   Count
1    STARTED   3
2    PHASE1    2
3    STOPPED   3
4    STARTED   2
4    PHASE2    1
2    STOPPED   1

我接下来要做的是首先丢弃那些在区间内看到超过1个状态的东西(所以ID 2和4将被淘汰),剩下的那些,丢弃那些状态不是"开始"(这将消除ID 3).ID 1是我正在寻找的记录.

这是问题的最佳方法吗?我如何实现该查询?

另外,如何让我的主题以不同的间隔发送消息,以便我可以测试窗口.

谢谢!

Jam*_*rld 7

实施解决方案

我解决问题的方法是生成一个扩展方法,它接受一个 IObservable<MyInfo>输入流(可能是a Subject)和一个IScheduler,并返回一个已经卡住的项目流.它看起来像这样:

public static class ObservableExtensions
{
    public static IObservable<MyInfo> StuckInfos(this IObservable<MyInfo> source,
        IScheduler scheduler = null)
    {
        scheduler = scheduler ?? Scheduler.Default;

        return source.Publish(pub =>
            pub.Where(x => x.Status == "STARTED")
                .SelectMany(
                    x => Observable.Return(x)
                        .Delay(TimeSpan.FromMinutes(30), scheduler)
                        .TakeUntil(pub.Where(y => y.Id == x.Id
                                                  && y.Status != "STARTED"))));

    }
}
Run Code Online (Sandbox Code Playgroud)

这里有很多Rx!让我们一点一滴......

一般的想法是,我们希望MyInfo在"已启动"状态下查找具有匹配Id 30分钟的未启动项目"未应答"的实例(以下称"项目").

暂时忽略这Publish一点,我会回到那个.想象一下pub变量是source.

第1步 - 过滤"已启动"项目

pub.Where(x => x.Status == "STARTED")
Run Code Online (Sandbox Code Playgroud)

这一点很简单,我们只是过滤了源代码以获得"STARTED"项目.

第2步 - 将每个项目转换为自己的延迟流

这有点棘手.从一个项目出现的那一刻起,我们知道30分钟因此我们想回答这个问题,"是否有另一个项目开始解开它?".为了帮助我们这样做,我们将创建一个新流,它将在30分钟后自动发出信息.我们的计划是,如果符合条件的项目出现问题,我们会缩短此流.假设x是信息,那么我们这样做:

Observable.Return(x)
          .Delay(TimeSpan.FromMinutes(30), scheduler)
Run Code Online (Sandbox Code Playgroud)

Observable.Return将一个项目转换为一个可观察的流,立即OnNext为该项目发布,然后OnComplete是s.尽管看起来有点毫无价值,但它实际上是一个非常有用的构建块(对于一些高级读取,有时也称为单元函数,IObservable monad结构的关键部分),因为它使我们能够从任何项目创建一个新的Observable流.一旦我们有这个流,我们Delay就这样,项目将在30分钟后出现.

请注意我们如何在调用时指定调度程序Delay- 这将有助于使用模拟时间进行测试.

第3步 - 如果某个项目变为"unstuck",则放弃延迟的流

现在,如果在30分钟之前让我们"解开"的合格项目到达,我们将不再对此项目感兴趣.

"unsticking"项目的合格标准是一个与Id匹配并且没有"STARTED"状态的标准 - 我认为该项目的另一个"STARTED"副本不足以解开卡住的项目!(如果这是错的,那么任何状态都有资格进行取消).如果未放置的项目到达原始未延迟的流(pub),我们将使用TakeUntil在延迟的项目有机会出现之前终止延迟的流:

.TakeUntil(pub.Where(y => y.Id == x.Id
                          && y.Status != "STARTED"))));
Run Code Online (Sandbox Code Playgroud)

第4步 - 整理所有这些项目流

现在,我们有点乱,因为我们将每个项目投射到它自己的流中 - 我们有一个流流,我们需要回到单个流.为此,我们使用SelectMany,(高级阅读:相当于monad绑定).在SelectMany我们完成两项工作在这里-它让我们的项目映射到一个流,并压平所产生的流的流回到一个流都一气呵成.我们将使用的映射函数是我们刚刚构建的函数 - 所以到目前为止我们将它们放在一起我们有:

pub.Where(x => x.Status == "STARTED")
                .SelectMany(
                 x => Observable.Return(x)
                                .Delay(TimeSpan.FromMinutes(30), scheduler)
                                .TakeUntil(pub.Where(y => y.Id == x.Id
                                                    && y.Status != "STARTED"))));
Run Code Online (Sandbox Code Playgroud)

第5步 - 确保我们不要使用来处理源流 Publish

我们看起来很好,但上面左边有一个微妙的问题需要解决.您会注意到我们订阅了多个source(pub)流 - 在初始Where过滤器AND中TakeUntil.

这样做的问题是,不止一次订阅同一个流可能会产生意想不到的后果.有些流是"冷"的 - 每个用户都会启动它自己的事件链.在时间是关键因素的查询中,这可能特别棘手.可能还有其他问题 - 但我不想在这里偏离太远.基本上,我们需要非常小心,我们只订阅一次源流.该Publish()方法可以为我们做到这一点 - 它将订阅一次源,然后将源多播到许多订阅者.

因此pub,lambda中出现的是我们可以安全地多次订阅的源的"安全"副本.

如何解决测试问题

您需要控制时间 - 最好的办法是在nuget包中使用专门构建的Rx测试工具rx-testing.有了这个,您可以使用a TestScheduler来控制时间和安排测试事件.

这是一个简单的测试,可以检测到一个简单的卡住项目.

public class StuckDetectorTests : ReactiveTest
{
    [Test]
    public void FindSingleStuckItem()
    {
        var testScheduler = new TestScheduler();

        var xs = testScheduler.CreateColdObservable(
            OnNext(TimeSpan.FromMinutes(5).Ticks, MyInfo.Started("1")));

        var results = testScheduler.CreateObserver<MyInfo>();

        xs.StuckInfos(testScheduler).Subscribe(results);

        testScheduler.Start();

        results.Messages.AssertEqual(
            OnNext(TimeSpan.FromMinutes(35).Ticks, MyInfo.Started("1")));
    }
}
Run Code Online (Sandbox Code Playgroud)

确保从中派生测试类ReactiveTest以利用OnXXX辅助方法.

我还创建了一些有用的工厂方法,MyInfo并实现了相等的重载,使测试更容易.

完整的代码非常冗长 - 我发布了一个在这里有更多测试的要点:https://gist.github.com/james-world/62dca2fe2f91531a0401

另外,我们还在测试的Rx一个很好的博客文章在这里.