You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何正确对Observable列表使用forkJoin并处理订阅与完成回调?

并行Observable处理:实时更新结果+完成回调实现

问题背景

我有大量需要并行执行的有序Observable,每个Observable返回结果后,需将结果追加到BehaviorSubject中传递,但必须在所有Observable执行完成后调用特定函数。

这些Observable各自从API下载图片及相关元数据,要求尽可能快地执行请求,且每个结果一返回就立即处理,所有请求完成时触发完成回调——这意味着Observable需要并行执行。

原实现(无完成回调)

const requests: Observable[] = getRequests();

requests.forEach(obs => obs.subscribe(res => {
    const currentImages = this.$images.value;
    currentImages.push(res);
    this.$images.next(currentImages);
}));

初次尝试的方案

为实现完成时的回调,我编写了以下代码,虽然功能可行,但将forkJoin与请求订阅拆分的写法显得冗余怪异,希望找到更优雅的实现方式:

const requests: Observable[] = getRequests();
const finishedTracker = new Subject<void>();

requests.forEach(obs => obs.subscribe(res => {
    const currentImages = this.$images.value;
    currentImages.push(res);
    this.$images.next(currentImages);
}));

forkJoin(requests).subscribe(() => {
    finishedTracker.next();
    finishedTracker.complete();
    console.log('requests done');
});

修正重复请求的尝试

后来意识到两次订阅会导致请求被执行两次,于是调整了实现方式:

from(requests).pipe(
    mergeMap(o => {
        o.subscribe(res => {
            const currentImages = this.$images.value;
            currentImages.push(res);
            this.$images.next(currentImages);
        })
        return o;
    }, 10)
).subscribe(() => {
    finishedTracker.next();
    console.log('requests done');
})

注:我没有直接使用forkJoin的结果,是因为它会等待所有请求完成后才返回全部结果,而我需要每个请求完成后立即将结果传递给BehaviorSubject——毕竟请求数量经常多达数百个。

最终采用的解决方案

from(requests).pipe(
    mergeMap(request => request, 10),
    scan<ImageResponse, ImageResponse[]>((all, current, index) => {
        all = all.concat(current);
        this.$images.next(all);
        return all;
    }, [])
).subscribe({
    complete: () => {
        finishedTracker.next();
        console.log('requests done');
    }
});

内容的提问来源于stack exchange,提问作者tgm

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.31 14:06:54