编写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()); }
代码说明
- ID串行控制:每个ID对应一条Promise链,新任务会在链的末尾等待执行,确保同一ID的任务严格按迭代器顺序执行,不会出现并发
- 全局并发限制:通过
activeTasks计数器和等待循环,确保同时运行的任务数始终不超过设定的线程限制 - 无限迭代器兼容:动态从AsyncIterator获取任务,直到迭代器标记
done,完美适配无限长任务序列场景 - 错误隔离:单个任务执行失败不会中断同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
相关产品推荐
相关产品推荐

