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

RxJS中实现带共享结果的请求队列及缓存机制问题

问题描述

我有一个服务,其中的Observable会从服务器查询对象数组并通过shareReplay共享,直到数据更新。当其他服务/组件需要特定对象时,可传入对象键值请求该服务,服务会返回包含已有匹配对象及可能需从服务器直接查询的缺失对象的Observable。

需求

  • 缓存额外查询的对象,供其他服务/组件复用,避免重复请求;
  • 处理同时到来的请求,通过队列避免重复查询同一项目;
  • 当初始Observable更新时,最终需得到包含更新数据及所有已查询额外对象的去重数组(按键值)。

我自行实现了一段代码,但担心多次请求后missingValuesObservable会不断变长,希望得到优化方案。

现有代码

options: Observable<any[]>;
missingValuesObservable: Observable<any[]>;

public fetchMissingValues(values: any[]): Observable<any[]> {
    values = values.unique();

    // Shared replay request (if any)
    let requestObservable: Observable<any[]> | undefined | null = undefined;

    // Update the observable chain to include this request
    this.missingValuesObservable = (this.missingValuesObservable ?
        this.missingValuesObservable :
        this.options
    ).pipe(
        concatMap(options => {
            // Create the request only the first time, then just share-replay it
            if (requestObservable === undefined) {
                let missingValues = values.filter(v => !options.find(o => ...));
                if (missingValues.length > 0) {
                    // Load missing values from server
                    requestObservable = this.apiService.request(...).pipe(
                        map(r => r.success && r.data?.items?.length ? r.data.items : []),
                        shareReplay(1)
                    );
                } else
                    requestObservable = null;
            }

            if (requestObservable)
                return requestObservable.pipe(
                    map(requestOpts => options.concat(requestOpts).unique(o => ...))
                );
            else
                return of(options);
        })
    );

    // Return only the first result, and only the requested values
    return this.missingValuesObservable.pipe(
        first(),
        map(options => values.map(v => options.find(o => ...)))
    );
}
优化方案

核心问题在于每次调用fetchMissingValues都会把新的请求逻辑追加到missingValuesObservable的管道链上,导致链越来越长,内存占用和复杂度上升。我们可以通过维护独立的缓存状态和集中处理请求队列来解决这个问题,具体实现如下:

优化思路

  • 独立缓存存储:用Map存储已查询到的额外对象,避免通过Observable链累积状态;
  • 请求队列去重:维护正在进行的请求Map,相同键值的请求复用同一个Observable,避免重复发起;
  • 自动合并数据源:监听初始options的更新,每次更新后自动合并缓存中的对象并去重,保证数据一致性。

优化后的代码

import { Observable, of, combineLatest, first, shareReplay, map, tap, switchMap } from 'rxjs';

@Injectable({ providedIn: 'root' })
export class YourService {
    // 初始数据源,保持shareReplay共享
    private readonly options$: Observable<any[]>;
    // 缓存已查询的额外对象,按键值存储
    private readonly cachedItems = new Map<string | number, any>();
    // 正在进行的请求,避免重复发起相同请求
    private readonly pendingRequests = new Map<string | number, Observable<any>>();

    constructor(private apiService: ApiService) {
        // 初始化初始数据源(替换为你实际的初始请求逻辑)
        this.options$ = this.apiService.getInitialOptions().pipe(
            shareReplay(1)
        );
    }

    public fetchMissingValues(values: (string | number)[]): Observable<any[]> {
        // 去重请求的键值
        const uniqueValues = [...new Set(values)];

        // 合并初始数据与缓存数据,生成当前完整数据集
        const fullData$ = this.options$.pipe(
            map(initialOptions => {
                const merged = [...initialOptions];
                // 合并缓存数据并去重
                this.cachedItems.forEach(item => {
                    if (!merged.find(o => this.getKey(o) === this.getKey(item))) {
                        merged.push(item);
                    }
                });
                return merged;
            }),
            shareReplay(1)
        );

        // 找出当前数据集中缺失的键值
        const missingKeys$ = fullData$.pipe(
            map(fullData => uniqueValues.filter(key => !fullData.find(o => this.getKey(o) === key)))
        );

        // 处理缺失键值的请求,自动去重并缓存结果
        const fetchMissing$ = missingKeys$.pipe(
            map(missingKeys => {
                return missingKeys.map(key => {
                    // 复用正在进行的请求
                    if (this.pendingRequests.has(key)) {
                        return this.pendingRequests.get(key)!;
                    }
                    // 发起新请求并缓存状态
                    const request$ = this.apiService.requestSingleItem(key).pipe(
                        tap(item => {
                            if (item) {
                                this.cachedItems.set(this.getKey(item), item);
                            }
                            this.pendingRequests.delete(key);
                        }),
                        shareReplay(1)
                    );
                    this.pendingRequests.set(key, request$);
                    return request$;
                });
            }),
            // 等待所有缺失项请求完成
            switchMap(requests => requests.length ? combineLatest(requests) : of([]))
        );

        // 返回最终结果:等待缺失项请求完成后,从完整数据集中提取目标对象
        return fetchMissing$.pipe(
            switchMap(() => fullData$),
            first(),
            map(fullData => uniqueValues.map(key => fullData.find(o => this.getKey(o) === key)))
        );
    }

    // 辅助方法:根据实际业务获取对象的键值(比如'id')
    private getKey(item: any): string | number {
        return item.id;
    }
}

代码说明

  • 缓存与请求去重:cachedItems存储已获取的额外对象,pendingRequests跟踪正在进行的请求,确保相同键值仅发起一次请求;
  • 数据合并逻辑:每次生成完整数据集时,自动合并初始options和缓存内容,初始数据更新后,缓存的额外对象也会被包含进去;
  • 管道链简化:不再通过追加missingValuesObservable的管道来累积状态,每次请求都基于当前的完整数据源处理,避免管道链无限变长。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 17:37:08