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是我正在寻找的记录.
这是问题的最佳方法吗?我如何实现该查询?
另外,如何让我的主题以不同的间隔发送消息,以便我可以测试窗口.
谢谢!
我解决问题的方法是生成一个扩展方法,它接受一个 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.
pub.Where(x => x.Status == "STARTED")
Run Code Online (Sandbox Code Playgroud)
这一点很简单,我们只是过滤了源代码以获得"STARTED"项目.
这有点棘手.从一个项目出现的那一刻起,我们知道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- 这将有助于使用模拟时间进行测试.
现在,如果在30分钟之前让我们"解开"的合格项目到达,我们将不再对此项目感兴趣.
"unsticking"项目的合格标准是一个与Id匹配并且没有"STARTED"状态的标准 - 我认为该项目的另一个"STARTED"副本不足以解开卡住的项目!(如果这是错的,那么任何状态都有资格进行取消).如果未放置的项目到达原始未延迟的流(pub),我们将使用TakeUntil在延迟的项目有机会出现之前终止延迟的流:
.TakeUntil(pub.Where(y => y.Id == x.Id
&& y.Status != "STARTED"))));
Run Code Online (Sandbox Code Playgroud)
现在,我们有点乱,因为我们将每个项目投射到它自己的流中 - 我们有一个流流,我们需要回到单个流.为此,我们使用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)
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一个很好的博客文章在这里.