【问题标题】:How to run an array of requests with rxjs like forkJoin and combineLatest but without having to wait for ALL to complete before seeing the results?如何使用 rxjs(如 forkJoin 和 combineLatest)运行一系列请求,但无需等待 ALL 完成才能看到结果?
【发布时间】:2021-10-10 05:12:20
【问题描述】:

假设你有一个array 的 URL:

urls: string[]

你创建了一个请求集合(在这个例子中,我使用了 Angular 的 HTTPClient.get,它返回一个 Observable)

const requests = urls.map((url, index) => this.http.get<Film>(url)

现在我想同时执行此请求,但不要等待所有响应才能看到所有内容。换句话说,如果我有 films$: Observable&lt;Film[]&gt; 之类的东西,我希望 films$ 在每次响应到达时逐渐更新。

现在要模拟这个,您可以将上面的 requests 更新为类似的内容

const requests = urls.map((url, index) => this.http.get<Film>(url).pipe(delay((index + 1)* 1000))

使用上面的Observables 数组,您应该从每个请求中一一获取数据,因为它们不是同时被请求的。请注意,这只是伪造来自各个请求的数据到达的不同时间。请求本身应该同时完成。

目标是更新films$ 中的元素,每次请求发出的时间值。

所以在我误解combineLatest 的工作原理时遇到这种情况之前

let films$: Observable<Film[]> = of([]);
const requests = urls.map(url => this.http.get<Film>(url)
.pipe(
  take(1),
 // Without this handling, the respective observable does not emit a value and you need ALL of the Observables to emit a value before combineLatest gives you results.
 // rxjs EMPTY short circuits the function as documented. handle the null elements on the template with *ngIf. 
  catchError(()=> of(null))
));
// Expect a value like [{...film1}, null, {...film2}] for when one of the URL's are invalid for example.
films$ = combineLatest(requests); 

我期待上面的代码逐渐更新films$,忽略了documentation的一部分

为确保输出数组始终具有相同的长度,combineLatest 实际上会等待所有输入 Observable 至少发出一次,然后才开始发出结果

这不是我想要的。

如果有一个 rxjs operatorfunction 可以实现我正在寻找的东西,我可以通过简单地使用 async 管道而不是必须处理 null 值和失败的请求。

我也试过

this.films$ = from(urls).pipe(mergeMap(url => this.http.get<Film>(url)));

this.films$ = from(requests).pipe(mergeAll());

这是不对的,因为返回的值类型是Observable&lt;Film&gt;,而不是Observable&lt;Film[]&gt;,我可以在模板上使用*ngFor="let film of films$ | async"。相反,如果我订阅它,就好像我正在监听一个记录的套接字,实时获取更新(单个响应进来)。例如,我可以手动订阅这两者中的任何一个并将Array.push 设置为单独的属性films: Film[],但这违背了目的(在带有async 管道的模板上使用Observable)。

【问题讨论】:

    标签: javascript angular typescript rxjs


    【解决方案1】:

    scan 运算符在这里非常适合您:

    
    const makeRequest = url => this.http.get<Film>(url).pipe(
      catchError(() => EMPTY))
    );
    
    films$: Observable<Film[]> = from(urls).pipe(
      mergeMap(url => makeRequest(url)),
      scan((films, film) => films.concat(film), [])
    ); 
    

    流程:

    • from 一次发出一个网址
    • mergeMap 订阅“makeRequest”并将结果发送到流中
    • scan 将结果累积到数组中,并在每次收到新发射时发射

    为了保持顺序,我可能会使用combineLatest,因为它以与输入 observables 相同的顺序发出一个数组。我们可以用startWith(undefined)开始每个observable,然后过滤掉未定义的项目:

    
    const requests = urls.map(url => this.http.get<Film>(url).pipe(startWith(undefined));
    
    films$: Observable<Film[]> = combineLatest(requests).pipe(
      map(films => films.filter(f => !!f))
    ); 
    

    【讨论】:

    • 布鲁...给我那个?
    • 您如何确保数据的顺序遵循 URL 的顺序?我有办法通过用几行替换array.concat 部分来做到这一点。只是好奇你自己会怎么做。
    • 在这种情况下我可能只使用combineLatest。我添加了一个示例。
    猜你喜欢
    • 1970-01-01
    • 2019-05-19
    • 2021-11-15
    • 1970-01-01
    • 2016-09-06
    • 1970-01-01
    • 2015-04-17
    • 2021-03-05
    • 2019-03-24
    相关资源
    最近更新 更多