【问题标题】:windowCount dropping valueswindowCount 丢弃值
【发布时间】:2016-10-01 10:47:41
【问题描述】:

我正在尝试使用 windowCount 将我的可观察值分组,并为每个组的每个值发送请求。
然后,连接这些组,以便在当前组的请求未完成之前不会启动下一个组的请求.
问题是某些值被跳过了。

这是我的代码。
(我没有在此处进行实际的 ajax 调用,但 Observable.timer 应该可以作为示例)。

Observable.interval(300)
     .take(12)
     .windowCount(3)
     .concatMap(obs => {
         return obs.mergeMap(
             v => Observable.timer(Math.random() * 1500).mapTo(v)
         );
     })
     .do(v => console.log(v))
     .finally(() => console.log('fin'))
     .subscribe();

我尝试通过手动创建组来替换 windowCount。而且效果很好。没有跳过任何值。

Observable.interval(900)
    .take(4)
    .map(i => Observable.interval(300).take(3).map(j => j + i * 3))
    .concatMap(obs => {
        return obs.mergeMap(
            v => Observable.timer(Math.random() * 1500).mapTo(v)
        );
    })
    .do(v => console.log(v))
    .finally(() => console.log('fin'))
    .subscribe();

我的印象是 windowCount 应该以相同的方式对发出的值进行分组。
但是,显然它做了其他事情。

我将非常感谢对其行为的任何解释。

谢谢!

【问题讨论】:

    标签: reactive-programming rxjs5


    【解决方案1】:

    缺失值是使用热可观察对象 (Observable.interval(300)) 的结果,该可观察对象继续输出您未存储以供使用的值。

    以下是您的代码的略微简化版本,它还记录了发出数字的时间。我用1 替换了Math.random(),这样输出是确定的。我还把jsbin中的代码加载了给大家试用:

    https://jsbin.com/burocu/edit?js,console

    Observable.interval(300)
        .do(x => console.log(x + ") hot observable at: " + (x * 300 + 300)))
        .take(12)
        .windowCount(3)
        .do(observe3 => {observe3.toArray()
          .subscribe(x => console.log(x + " do window count at: " + (x[2] * 300 + 300)));})
        .concatMap(obs => {
            return obs.mergeMap(
                v => Observable.timer(1 * 1500).mapTo(v)
            )
            .do(v => console.log(v + " merge map at: " + (v * 300 + 300 + 1500)));
        })
        .finally(() => console.log('fin windowCount'))
        .subscribe();
    

    它导致下面的输出。请注意,当其他运算符仍在处理时,热门的 observables 会继续前进。

    这就是给您的印象,即价值正在被丢弃。您可以看到 windowCount(3) 正在按照您的想法进行操作,而不是在何时按照您的想法进行操作。

    "0) hot observable at: 300"
    "1) hot observable at: 600"
    "2) hot observable at: 900"
    "0,1,2 do window count at: 900"
    "3) hot observable at: 1200"
    "4) hot observable at: 1500"
    "5) hot observable at: 1800"
    "3,4,5 do window count at: 1800"
    "0 merge map at: 1800"
    "6) hot observable at: 2100"
    "1 merge map at: 2100"
    "7) hot observable at: 2400"
    "2 merge map at: 2400"
    "8) hot observable at: 2700"
    "6,7,8 do window count at: 2700"
    "9) hot observable at: 3000"
    "10) hot observable at: 3300"
    "11) hot observable at: 3600"
    "9,10,11 do window count at: 3600"
    " do window count at: NaN"
    "8 merge map at: 4200"
    "fin windowCount"
    

    编辑:进一步解释...

    windowCount(3) 之后有一个对concatMap 的调用。 concatMapmapconcatAll 的组合。

    concatAll:

    加入源发出的每个 Observable(高阶 可观察的),以串行方式。它订阅每个内部 Observable 仅在前一个内部 Observable 完成后(添加了重点),并且 将它们的所有值合并到返回的 observable 中。

    因此,查看上面的输出,我们看到第一个 windowCount(3) 值 [0,1,2] 在 1800 到 2400 之间发出。

    请注意,第二个windowCount(3) 值 [3,4,5] 在 1800 处发出。concatAll 在发出 [3,4,5] 时还没有准备好订阅,因为 之前的内部 Observable 有尚未完成。所以这些值被有效地删除了。

    接下来,注意之前的内部 Observable [0,1,2] 在 2400 完成。concatAll 在 2400 订阅。

    下一个出现的值是 2700 处的值 8(订阅在 2400 开始后 300 毫秒)。然后mergeMap 在 4200 处输出值 8,因为从订阅开始点 2400 的间隔延迟为 300,然后是 1500 的计时器延迟(即 2400 + 300 + 1500 = 4200)。

    在这一点之后,序列完成,因此不再发出任何值。

    如果需要更多说明,请添加评论。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-11-06
      • 2015-09-05
      • 1970-01-01
      • 2021-07-27
      • 1970-01-01
      • 1970-01-01
      • 2018-05-02
      • 2017-06-02
      相关资源
      最近更新 更多