【问题标题】:Why is initial stream being triggered again after combineLatest and merge in this example?在本例中,为什么在 combineLatest 和 merge 之后再次触发初始流?
【发布时间】:2018-10-09 09:59:41
【问题描述】:

看下面的摘录:

  let requestStream = Rx.Observable
    .of(`${GITHUB_API}?since=${randomNumber()}`)
    .mergeMap(url => {
      console.log(`performing request to: ${url}`)
      return Rx.Observable.from(jQuery.getJSON(url))
    });

  let refreshStream = Rx.Observable.fromEvent(refreshButton, 'click')
    .startWith('click')
    .do(_ => users.empty())
    .combineLatest(requestStream, (_, users) => users.slice(randomNumber(users.length)));

  let randomUserStream = userRemovedStream
    .combineLatest(requestStream, (_, users) => users[randomNumber(users.length)]);

  requestStream
    .merge(refreshStream)
    .flatMap(users => users)
    .merge(randomUserStream)
    .filter(_ => users.children().length < MAX_SUGGESTIONS)
    .do(user => users.append(createItem(user)))
    .mergeMap(user => Rx.Observable.fromEvent($(`#close-${user.login}`), 'click'))
    .map(event => event.target.parentNode)
    .subscribe(user => {
      user.remove();
      userRemovedStream.next('');
    });

requestStream 返回一个包含 100 个用户的数组,但是,我当时只使用了其中的三个 (MAX_SUGGESTIONS)。 refreshStreamrandomUserStream 的存在是为了重用来自 requestStream 的其他 97 个用户。问题是,当我运行上面的代码时,我仍然在控制台上看到performing request to: ... 三次。

我注意到在最后一个流中添加 merge 方法后会发生这种情况,但是,我不确定为什么会发生这种行为。

我的理解是:当我mergerefreshStreamrandomUserStream时,每当有新项目发出时,前者点击refresh按钮,后者点击remove按钮, requestStream 上先前发出的数组将被解析并向前传递,而不是单击本身。这不应该重新触发requestStream

有人可以帮助我了解为什么会发生这种情况以及如何处理这种情况吗? - 所以我可以在第一次调用期间从 API 已经返回的用户中获取最大用户数?

【问题讨论】:

    标签: rxjs rxjs5


    【解决方案1】:

    这是因为您有效地订阅了三个requestStream。但是,您对这三者如何交互的直觉是正确的,因为您的 requestStream Observable 是冷的,它会在每次订阅时创建一个新流。

    这并不一定很明显,因为只有一个订阅是显式的,但每次你将requestStream 传递给combineLatest 时,它最终都会创建一个新订阅,然后又会启动一个新流,在这种情况下调用你的底层 API。

    如果您不希望这种情况发生,我建议您使用像 publishLast 这样的多播运算符

    所以requestStream 会变成:

    let requestStream = Rx.Observable
        .of(`${GITHUB_API}?since=${randomNumber()}`)
        .mergeMap(url => {
          console.log(`performing request to: ${url}`)
          return Rx.Observable.from(jQuery.getJSON(url))
        })
        .publishLast();
    

    在这种情况下,requestStream 现在实际上是 ConnectableObservable,因此您还需要在某个时候启动它,通常您会等到所有订阅者都连接上。

    /* Rest of you example */
    .map(event => event.target.parentNode)
    .subscribe(user => {
      user.remove();
      userRemovedStream.next('');
    });
    
    requestStream.connect();
    

    【讨论】:

    • 那是一个很好的解释,非常非常感谢!我将更多地研究冷与热 Observables 和多播。到目前为止效果很好!
    猜你喜欢
    • 1970-01-01
    • 2013-02-23
    • 2014-04-19
    • 2017-04-09
    • 2015-01-07
    • 1970-01-01
    • 2015-08-09
    • 1970-01-01
    • 2011-04-01
    相关资源
    最近更新 更多