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

编写JS/TS函数:基于线程数与任务ID限制并发执行AsyncIterator任务

带ID限制的异步并发执行函数实现

实现思路

要满足三个执行规则,需要同时处理全局并发数限制和同ID任务串行执行两个核心约束:

  • 用Map维护每个ID的执行链:同一ID的任务必须按迭代器顺序依次执行,新任务会追加到对应ID的Promise链末尾
  • 用计数器跟踪全局活跃任务数,确保同时运行的任务不超过设定线程限制
  • 针对无限长的AsyncIterator,采用动态取任务+调度的方式,避免一次性加载所有任务

最终实现代码

async function limitedConcurrentHandler<T>(
  handler: (task: T) => Promise<void>,
  taskIterator: AsyncIterator<T>,
  concurrencyLimit: number = Infinity
): Promise<void> {
  // 存储每个ID的执行链,保证同ID任务串行执行
  const idExecutionMap = new Map<string | number, Promise<void>>();
  // 当前全局活跃任务数
  let activeTasks = 0;
  // 标记迭代器是否已取完所有任务
  let iteratorDone = false;

  // 调度单个任务执行的逻辑
  async function scheduleTask(task: T) {
    const taskId = (task as any).id;
    // 获取当前ID的执行链,初始为已完成的Promise
    const currentChain = idExecutionMap.get(taskId) || Promise.resolve();

    // 构建新的执行链:等当前链完成后,再执行当前任务,同时遵守全局并发限制
    const newChain = currentChain
      .then(async () => {
        // 等待全局并发数有空位
        while (activeTasks >= concurrencyLimit) {
          await new Promise(resolve => setTimeout(resolve, 10));
        }
        activeTasks++;
        try {
          await handler(task);
        } finally {
          activeTasks--;
        }
      })
      .catch(err => {
        // 单个任务失败不中断同ID后续任务,可根据需求调整错误处理逻辑
        console.error(`Task ${taskId} failed:`, err);
      });

    // 更新当前ID的执行链
    idExecutionMap.set(taskId, newChain);
    // 执行完成后清理Map项(可选优化,避免内存占用)
    newChain.finally(() => {
      if (idExecutionMap.get(taskId) === newChain) {
        idExecutionMap.delete(taskId);
      }
    });
  }

  // 循环从迭代器取任务并调度
  while (!iteratorDone) {
    try {
      const { value, done } = await taskIterator.next();
      if (done) {
        iteratorDone = true;
        break;
      }
      await scheduleTask(value);
    } catch (err) {
      console.error('Failed to fetch task:', err);
      iteratorDone = true;
    }
  }

  // 等待所有剩余任务执行完成
  await Promise.all(idExecutionMap.values());
}

代码说明

  1. ID串行控制:每个ID对应一条Promise链,新任务会在链的末尾等待执行,确保同一ID的任务严格按迭代器顺序执行,不会出现并发
  2. 全局并发限制:通过activeTasks计数器和等待循环,确保同时运行的任务数始终不超过设定的线程限制
  3. 无限迭代器兼容:动态从AsyncIterator获取任务,直到迭代器标记done,完美适配无限长任务序列场景
  4. 错误隔离:单个任务执行失败不会中断同ID的后续任务,也不会影响其他ID的任务执行

使用示例

// 模拟异步任务处理函数
async function mockHandler(task) {
  console.log(`Start task ${task.id} (${task.type})`);
  await new Promise(resolve => setTimeout(resolve, 100));
  console.log(`Finish task ${task.id} (${task.type})`);
}

// 模拟异步任务迭代器(可无限生成任务)
async function* mockTaskIterator() {
  let id = 1;
  while (true) {
    yield { id: id % 3, type: id % 2 === 0 ? 'init' : 'cleanup' };
    id++;
    await new Promise(resolve => setTimeout(resolve, 20));
    // 执行10个任务后停止,可移除该行实现无限任务
    if (id > 10) break;
  }
}

// 调用函数:限制全局并发数为2
limitedConcurrentHandler(mockHandler, mockTaskIterator(), 2);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 12:55:20