如何用Rx.js创建支持异步并发限制(最多3个任务)的可观察队列?
实现动态可观察队列并控制并发数
要实现动态可观察队列(能随时推送新任务并订阅结果),同时保持最多3个并发任务执行,你可以用RxJS的Subject作为任务入口——它既是Observable(可订阅输出)也是Observer(可推送新任务),再结合mergeAll(3)控制并发量。
完整实现代码
import { defer, Subject, Observable } from 'rxjs'; import { mergeAll } from 'rxjs/operators'; async function getData(x: string) { return new Promise((resolve) => { setTimeout(() => resolve(x), 1000); }); } // 1. 创建Subject作为动态任务队列 const taskQueue$ = new Subject<Observable<string>>(); // 2. 订阅队列,用mergeAll控制并发数为3 let data = [ { name: 'Jon start' } ]; taskQueue$ .pipe(mergeAll(3)) .subscribe((d) => { data.push({ name: d }); console.log('任务完成:', d); console.log('当前数据:', data); }); // 3. 定义推送任务的方法 function pushTask(id: number) { // 用defer包裹,确保每次推送才触发getData执行 const task$ = defer(() => getData('John ' + id)); taskQueue$.next(task$); } // 示例:初始推送20个任务 const ids = [...Array(20).keys()]; ids.forEach(id => pushTask(id)); // 示例:动态添加新任务(比如2秒后再加5个) setTimeout(() => { console.log('开始添加新任务'); for (let i = 20; i < 25; i++) { pushTask(i); } }, 2000);
关键部分说明
Subject作为队列:taskQueue$用来接收所有新的任务Observable,调用taskQueue$.next(task$)就能把任务推送到队列中。mergeAll(3):确保同时最多有3个任务在执行,当某个任务完成后,自动从队列中取出下一个任务执行。defer的作用:包裹getData调用,保证每次推送任务时才会触发异步请求,避免提前创建Promise导致失去队列控制的意义。- 动态添加任务:你可以在任何时机调用
pushTask添加新任务,队列会自动按照并发规则执行。
这样实现后,队列既支持初始批量任务,也支持后续动态推送新任务,同时保持3个并发的限制,所有任务的执行结果都会通过订阅统一处理。
内容的提问来源于stack exchange,提问作者Harsh
相关产品推荐
相关产品推荐

