如何用RxJS创建请求队列处理外部模块的回调数据加载?
使用RxJS实现请求队列避免并发HTTP请求
我有一个外部模块,它会在运行过程中通过回调函数加载数据:
const module = useModule(callback); function callback(params:{id:number, takeLast:boolean}):Promise<any>;
该回调函数接收params参数并返回一个Promise,从缓存或通过HTTP请求加载数据。
我希望创建一个队列来逐个处理请求并返回数据,避免生成大量并发HTTP请求。如何使用RxJS实现这一需求?
我尝试了以下代码:
const queue = new Subject<{id:number, takeLast:boolean}>(); const moduleHandler = useModule((params) => { return new Promise<any>((resolve, reject) => { queue.next(params); console.log('request for data'); queue // .asObservable() .pipe( concatMap((xx) => of( requestToHttpAndCache(xx) ).pipe(delay(1000)) ) ) .subscribe( { next: (result: any) => { debugger; resolve(result); }, error: (err: Error) => { debugger; }, } ); }); });
但这段代码会针对同一数据发起多个HTTP请求,我无法找到正确的RxJS队列实现方式。
问题分析
你之前的代码核心问题有两个:
- 每次回调触发时都会创建新的
subscribe订阅,同一个请求参数会被所有订阅接收并处理,导致重复请求; of(requestToHttpAndCache(xx))只是把Promise包装成Observable,没有等待Promise完成,delay也只是延迟传递Promise对象,无法实现真正的串行执行。
正确实现方案
我们需要构建一个单例的串行处理流,将所有请求统一送入队列,确保逐个执行,同时为每个请求绑定对应的Promise回调,保证结果能正确返回给发起请求的回调函数。
完整代码
import { Subject, concatMap, tap } from 'rxjs'; // 定义队列项类型,绑定请求参数和Promise的resolve/reject type QueueItem = { params: { id: number; takeLast: boolean }; resolve: (value: any) => void; reject: (error: any) => void; }; // 创建请求队列Subject const requestQueue = new Subject<QueueItem>(); // 初始化单例处理流,仅订阅一次 requestQueue.pipe( // concatMap保证上一个请求完成后再处理下一个,实现串行队列 concatMap((item) => // 执行实际的请求逻辑,requestToHttpAndCache返回Promise会被RxJS自动处理 requestToHttpAndCache(item.params).pipe( // 请求完成后调用对应的resolve/reject tap({ next: (result) => item.resolve(result), error: (err) => item.reject(err) }) ) ) ).subscribe(); // 传给外部模块的回调函数 const moduleHandler = useModule((params) => { return new Promise<any>((resolve, reject) => { // 将参数和Promise回调送入队列 requestQueue.next({ params, resolve, reject }); }); });
关键优化点
- 单例订阅:整个队列只订阅一次,避免重复处理请求;
- concatMap串行执行:确保请求按顺序逐个完成,不会产生并发HTTP请求;
- 请求与回调绑定:每个队列项携带自己的
resolve/reject,请求结果直接对应到发起请求的Promise; - 可选缓存增强:如果
requestToHttpAndCache未实现缓存,可以在流中加入去重逻辑,避免重复请求相同参数:
import { distinctUntilChanged } from 'rxjs'; requestQueue.pipe( // 根据id和takeLast判断是否为重复请求 distinctUntilChanged((prev, curr) => prev.params.id === curr.params.id && prev.params.takeLast === curr.params.takeLast ), concatMap(/* ... */) ).subscribe();
内容的提问来源于stack exchange,提问作者lkurylo
相关产品推荐
相关产品推荐

