【问题标题】:How to combine observables in a non-blocking way?如何以非阻塞方式组合可观察对象?
【发布时间】:2019-06-12 21:33:59
【问题描述】:

我有一组 Observable,它们分别检索不同的数据类型。我正在链接那些 Observable 以获取我想要的页面的所有数据。 事实是所有这些信息都是独立的,因此加载一个不应阻塞或干扰加载其他信息。这是我无法实现的。

这是我目前所做的一个示例:

  getAll(): Observable<any[]> {
    return this.getHome().pipe(
      tap((home: Home) => this.application = home),
      mergeMap((home: Home) =>
        combineLatest([
          this.getAssets(home.images).pipe(defaultIfEmpty([])),
          this.getStyling(home.styling).pipe(defaultIfEmpty({})),
          this.getCategories(home.categories).pipe(defaultIfEmpty([])),
        ])
      )
    );
  }

在这里,Observable 已成功链接,所以我得到了我想要的所有数据,但是 combineLatest() 在上一个请求完成后执行下一个请求,并且一旦一切完成,父级就会发出收集到的数据,这会造成延迟。 我一直在尝试使用 merge() 而不是 combineLatest() 以便收到的每个数据都会立即将其发送到父 Observable,但显然我不能这样做

那么,我怎样才能设法链接这些 Observable,以便子 Observable 中的每个提交都直接发送到父 Observable?

【问题讨论】:

  • "我怎样才能设法链接这些 Observables,以便子 Observable 中的每个提交都直接发送到父 Observable",这就是问题,你怎么知道正在发出什么类型?如果你想发射到一个流中,你可以连接一个主题。
  • @AvinKavish 我需要对主题进行调查,您对此有很好的阅读吗?
  • 是的,请稍等,我也有一个半书面的主题答案。
  • 看看,它只是一个参考实现。主题可以通过多种不同的方式组合。

标签: javascript angular rxjs


【解决方案1】:

您需要forkJoin,它会返回一个Array,其中包含每个Observable 的结果(与forkJoin 参数的顺序相同。

  getAll(): Observable<any[]> {
    return this.getHome().pipe(
      tap((home: Home) => this.application = home),
      mergeMap((home: Home) =>
        forkJoin(
          this.getAssets(home.images).pipe(defaultIfEmpty([])),
          this.getStyling(home.styling).pipe(defaultIfEmpty({})),
          this.getCategories(home.categories).pipe(defaultIfEmpty([]))
        )
      )
    );
  }

【讨论】:

  • forkJoin 不是在所有内部 Observable 完成后才发出吗?这就是我从文档中了解到的
  • 是的,它会这样做,如果合并到同一个事件流中,你怎么知道哪个请求向你发送了数据? forkJoin 允许您将结果映射到它的发起者
  • 但这里的问题是,如果this.getAssets 需要 4000 毫秒,而 this.getStylingthis.getCategories 会在 50 毫秒内响应,我的父 obs 将被阻塞 4000 毫秒,而不是在 50 毫秒后发射然后重新发射4000 毫秒后。你明白我的意思吗?
  • 我肯定明白了,那为什么不直接返回一个Observable[],然后在组件中嵌套订阅呢?
  • 好吧,我一直在寻找答案已经有一段时间了,看起来使用 Observable[] 似乎是该问题的最终答案,我会稍等片刻看看人们会带来什么,但我最终可能会得到那个
【解决方案2】:

如果它们是独立的,您可以独立处理它们并将它们发射到主题上。

  getAll(): Observable<any> {
    const subject = new Subject<any>();
    const handler = (res: any) => subject.next(res);

    this.getHome().pipe(
      tap((home: Home) => this.application = home),
      tap(handler),
      tap((home: Home) => {
        forkJoin(
          this.getAssets(home.images).pipe(defaultIfEmpty([]),tap(handler))
          this.getStyling(home.styling).pipe(defaultIfEmpty({}), tap(handler))                  
          this.getCategories(home.categories).pipe(defaultIfEmpty([]), tap(handler))
        ).subscribe(res => subject.complete())
      })
    ).subscribe();
    return subject.asObservable();
  }

使用这种方法,类型检查必须由消费者完成。

Subjects are documented here.

【讨论】:

  • 现在看起来是最好的答案,它有点老套,但会完全按照我的要求行事 感谢这位伙伴!
  • 这太难了 :) 你应该选择@Reactangular 的答案,运营商似乎是 merge()
【解决方案3】:

我一直在尝试使用 merge() 而不是 combineLatest() 以便接收到的每个数据都会立即将其发送到父 Observable,但显然我不能这样做

听起来您正在寻找准备好的部分数据。 combineLatest()forkJoin() 都只会在它们可以填充长度与可观察对象数量相等的数组时发出数据。因此,您要么等待每个 observable 中的至少一个值,要么等待它们全部完成,但它们会发出一个数组,并且需要知道在其中放入什么。

我认为您希望在准备好后立即发出每条数据。

首先使用 merge() 从所有 3 个可观察对象发出值,但对于每个可观察对象,使用 map() 为给定类型(图像、样式、类别)创建对象键/值对。

最后,使用scan() 运算符将值聚合到单个对象中,该对象将在更新时发出每个属性。

流的消费者可以检查一个属性是否为null 以查看它是否已准备好。流将在完成之前发出 3 个值。

  getAll(): Observable<{home: Home, images: any[], styling: any[], categories: any[]}> {
    return this.getHome().pipe(
      tap((home: Home) => this.application = home),
      switchMap((home: Home) =>
        merge([
          this.getAssets(home.images).pipe(
             defaultIfEmpty([]),
             map(images => ({images}))
          ),
          this.getStyling(home.styling).pipe(
             defaultIfEmpty({}),
             map(styling => ({styling}))
          ),
          this.getCategories(home.categories).pipe(
             defaultIfEmpty([]),
             map(categories => ({categories})
          ),
        ]),
        scan((acc, next) => ({...acc, ...next}), {home, images: null, styling: null, categories: null})
      )
    );
  }

如果您想在所有属性都是null 的情况下发出第一个值,可以将startWith({}) 运算符添加到merge() 运算符。作为消费者的早期价值,因此他们知道阅读其他价值已经开始。

【讨论】:

  • 切换地图不会阻止第一个可观察对象的发射吗?
  • @AvinKavish 不,当外部 observable 发出值时,它会取消订阅内部 observable,但它永远不会停止外部 observable。
  • @AvinKavish 对我来说是一般规则。如有疑问,请始终使用switchMap(),因为mergeMap() 会泄漏内存。
  • 是的,所以它将流转换为新的 observable,消费者看不到 getHome() 发出的 observable。因此被称为 switch ... map
  • @AvinKavish 是的。流完全切换到映射流,但是当上层流发出第二个值时,它会再次切换。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-06-18
  • 2017-04-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-12-11
相关资源
最近更新 更多