【问题标题】:Reusing streams of data for some time重用数据流一段时间
【发布时间】:2017-06-02 11:22:47
【问题描述】:

假设有一个 API 接受查询并返回结果流,因为某些结果可能会发生变化。

type Query = { 
  species?: "dog" | "cat" | "rat", 
  name?: "string",
  status?: "lost" | "found"
}
type Result = { species: string, name: string, status: string }[]

假设有多个组件向此 API 传递查询,其中一些可能是相同的。不想向服务器发送不必要的请求并喜欢优化 - 为了做到这一点,可以封装 API 并拦截调用。

interface ServiceApi {
  request(query: Query): Observable<Result>
}

class WrappedServiceApi implements ServiceApi {
  constructor(private service: ServiceApi) { }

  request(query: Query): Observable<Result> {
    // intercepted
    return this.service.request(query);
  }
}

但是如何使用 RxJS 5 进行这种优化呢?

在 RxJS 周围做这件事可能看起来像这样:

class WrappedServiceApi implements ServiceApi {

  private activeQueries;
  constructor(private service: ServiceApi) { 
    this.activeQueries = new Map<string, Observable<Result>>();
  }

  request(query: Query): Observable<Result> {
    // it's easy to stringify query
    const hashed: string = hash(query);
    if (this.activeQueries.has(hashed)) {
      // reuse existing stream
      return this.activeQueries.get(hashed);
    } else {
      // create multicast stream that remembers last value
      const results = this.service.request(query).publishLast();
      // store stream for reuse
      this.activeQueries.set(hashed, results);
      // delete stream 5s after it closed
      results.toPromise().then(
        () => setTimeout(
          () => this.activeQueries.delete(hashed), 
          5000
        )
      );
      return results;
    }
  }
}

是否可以通过更具声明性的 rx 方式实现相同的效果?

【问题讨论】:

    标签: javascript functional-programming rxjs reactive-programming rxjs5


    【解决方案1】:

    我没有测试它,但我会这样做:

    request(query: Query): Observable<Result> {
        return Observable.of(hash(query))
            .flatMap(hashed => {
                const activeQueries = this.activeQueries;
                if (activeQueries.has(hashed)) {
                    return activeQueries.get(hashed);
                } else {
                    const obs = this.service.request(query)
                        .publishReplay(1, 5000)
                        .refCount()
                        .take(1);
    
                    activeQueries.set(hashed, obs);
                    return obs;
                }
            });
    }
    

    基本上唯一的区别是我将 Observables 存储在 this.activeQueries 中,而不仅仅是它们的结果,然后我使用 .publishReplay(1, 5000) 将它们缓存 5 秒。这样,当您在 5 秒后订阅同一个 Observable 时,它​​只会重新订阅其源 Observable。

    【讨论】:

    • 除非你做一个 flatMap 它返回更高的订单流,所以违反了合同。原始示例还在 activeQueries 中存储了 observables(但是 publishReplay 比 publishLast 更好)。文档链接建议使用 Rx.Obervable.defer 和流工厂,以确保实际调用 API。
    • 不过,摆脱这个 activeQueries 地图也不错。
    • 不错,应该有flatMap
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-11-16
    • 2012-08-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-12-15
    相关资源
    最近更新 更多