【问题标题】:Rx.Net GroupJoin two Observables with time in join conditionRx.Net GroupJoin 两个 Observables 在加入条件下的时间
【发布时间】:2018-01-03 10:08:53
【问题描述】:

给定 2 个热门的 observables t1 和 t2,我将如何 GoupJoin 以便从 t2 获取所有事件,这些事件发生在 t1 中每个事件之前 x 秒和 y 秒之后?

给定:

t1 -----A-----B-----C

t2 --1--2--3--4--5--6

如果 t1 相隔 2 秒,t2 相隔 1 秒,并且我们正在寻找每个 t1 事件两侧相隔 1 秒的 t2 事件,则结果如下。

结果:

{ A, [1,2,3] }

{ B, [3,4,5] }

{ C, [5,6] }

下面是一个真实的例子,我们需要解决上述问题。 我们有一个电子邮件流和另一个文本消息流。我们需要发出另一个流的结果,其中包含电子邮件,并且文本消息发生在电子邮件发送时间之前或之后的 1 分钟内。

【问题讨论】:

    标签: system.reactive rx.net


    【解决方案1】:

    这里的问题(正如 Shlomo 所提到的)是我们需要在t1 事件发生之前打开t2 中的窗口。不幸的是,这是不可能的,因为一旦我们到达t1 中的事件,我们就已经过了需要在t2 中打开窗口的地步。

    我们可以做的是使用Delay() 将t2 及时向前移动。如果我们将其抵消x(之前的时间),我们可以将问题重新定义为“获取t2 中发生在t1 窗口打开和t1 + x + y 关闭的事件。我们可以使用GroupJoin解决这个问题。

    var scheduler = new HistoricalScheduler();
    
    var t1 = Observable.Interval(TimeSpan.FromMilliseconds(200), scheduler)
        .Select(l => (char)('A' + l));
    var t2 = Observable.Interval(TimeSpan.FromMilliseconds(100), scheduler);
    
    var x = TimeSpan.FromMilliseconds(100);     //before time
    var y = TimeSpan.FromMilliseconds(100);     //after time
    
    var delayedT2 = t2.Delay(x, scheduler);
    
    var g = t1.GroupJoin(delayedT2 ,
        _ => Observable.Timer(x + y, scheduler),
        _ => Observable.Empty<Unit>(scheduler),
        (a, b) => new { a, b}
    );
    
    scheduler.Start();
    

    这给出了结果:

    { A, [1,2] }
    { B, [3,4] }
    { C, [5,6] }
    

    这个结果仍然不是你所期望的。这是因为在您的示例中,t2 事件发生在完全相同的瞬间t1 事件中。在这种情况下,首先处理t1 + y 事件并在可以包含t2 事件之前关闭窗口。这意味着我们有效地获得了(t1-01:00) &lt;= t1 &lt; (t1 + 01:00)。例如。 A 的窗口是 01:0000 - 02.9999... 这就是为什么不包括在 03:00 发生的 3

    只需在我们的y 时间上添加一个刻度,即可将其修复为包容性

    var y = TimeSpan.FromMilliseconds(100).Add(TimeSpan.FromTicks(1)); 
    

    【讨论】:

    • 我想到了这个。我喜欢它的优雅,但我发现它不是最理想的:这意味着每个结果都有 x 的延迟。如果 x 很大,并且/或者结果是时间敏感的,那就不好了。
    【解决方案2】:

    代码转储答案(使用 100 毫秒代替 1 秒):

    var t1 = Observable.Interval(TimeSpan.FromMilliseconds(200))
        .Select(l => (char)('A' + l))
        .Delay(TimeSpan.FromMilliseconds(200));
    var t2 = Observable.Interval(TimeSpan.FromMilliseconds(100))
        .Delay(TimeSpan.FromMilliseconds(100));
    
    var x = TimeSpan.FromMilliseconds(100);     //before time
    var y = TimeSpan.FromMilliseconds(100);     //after time
    
    var g = t1.Timestamp().Join(t2.Timestamp(),
        c => Observable.Timer(y),
        i => Observable.Timer(x + y),
        (c, i) => new {GroupItem = c, RightItem = i}
    )
        .Where(a =>
            (a.GroupItem.Timestamp > a.RightItem.Timestamp && a.GroupItem.Timestamp - a.RightItem.Timestamp <= x) //group-item came first
            || (a.GroupItem.Timestamp <= a.RightItem.Timestamp && a.RightItem.Timestamp - a.GroupItem.Timestamp <= y) // right-item came first, or exact timestamp match
        )
        .Select(a => new { GroupItem = a.GroupItem.Value, RightItem = a.RightItem.Value })
        .GroupBy(a => a.GroupItem, a => a.RightItem);
    

    解释: Join 是关于“windows”的。因此,当您定义一个连接时,您必须考虑从左可观察和右可观察的每个项目打开的时间窗口。不过,我们这里的窗口很难弄清楚:我们必须以某种方式在它发生之前在左侧可观察 X 时间打开一个窗口,然后在它发生后 Y 时间关闭它。

    而不是做不可能的事,所以我们只在左侧项目出现后将其打开 Y 时间,并让右侧项目窗口由 X + Y 时间定义。但是,这会给我们留下不应该包含的项目。所以我们在时间戳上使用Where 来过滤掉它们。

    最后我们选择匿名类型和时间戳,并将它们组合在一起。

    我不认为 GroupJoin 是去这里的方式:你最终会像我所做的那样拆散小组并重新组建它..

    【讨论】:

    • 这给出了:{ A, [1,2] } { B, [3,4] ] { C, [5] } 这不是预期的结果。
    • 编辑答案以使用精确的刻度匹配(HistoricalScheduler 案例)。
    猜你喜欢
    • 2017-09-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-01-12
    • 2012-04-21
    • 1970-01-01
    • 2014-08-03
    • 1970-01-01
    相关资源
    最近更新 更多