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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 23:27:17