JavaScript 有界并发任务池:把任务延后启动,才能真正限制同时运行的数量
批量生成缩略信息时,直接把一百次请求交给 Promise.all,往往会在很短时间内同时启动。聚合结果的方法本身不会替你控制连接数。真正需要限制的是“任务开始到结束”这段区间里,有多少项工作处于活跃状态。
解决办法是先保存尚未执行的任务函数,再让固定数量的工作循环领取任务。一次领取后等待它完成,无论成功还是失败,才继续领取下一项。这样控制点位于启动之前,也能把错误收集在各自的位置,方便最后逐条报告。
用可观察的计数器验证上限
下面保存为 demo.mjs,用 Node.js 运行。例子只使用本地计时器,不访问网络。四个任务的延迟不同,其中一个主动失败;active 记录当前活跃数量,peak 记录峰值。断言同时核对结果位置、失败状态和最终归零,避免只凭日志先后判断并发。
输入要求是稠密数组,每一项都必须是返回值或 Promise 的函数。上限必须为正整数。同步抛出的异常也会进入对应失败记录,因为函数调用放在 try 内。示例把 Error 取为 message,其他常见抛出值转成字符串,并额外检验 null 与同步抛出的字符串,防止读取不存在的 message 时再次失败。
AI概念示意图:以抽象物件说明本文主题,不代表真实界面或运行结果。
import assert from 'node:assert/strict';
const wait = ms => new Promise(resolve => setTimeout(resolve, ms));
async function runPool(jobs, limit) {
if (!Number.isInteger(limit) || limit < 1) {
throw new RangeError('limit must be a positive integer');
}
let next = 0;
const results = new Array(jobs.length);
async function worker() {
while (next < jobs.length) {
const index = next++;
try {
results[index] = {ok: true, value: await jobs[index]()};
} catch (error) {
results[index] = {ok: false, error: error instanceof Error ? error.message : String(error)};
}
}
}
await Promise.all(Array.from(
{length: Math.min(limit, jobs.length)}, () => worker()
));
return results;
}
let active = 0, peak = 0;
const jobs = [30, 5, 10, 1].map((ms, index) => async () => {
active++;
peak = Math.max(peak, active);
try {
await wait(ms);
if (index === 2) throw new Error('demo failure');
return `item-${index}`;
} finally {
active--;
}
});
const results = await runPool(jobs, 2);
assert.deepEqual(results.map(x => x.ok), [true, true, false, true]);
assert.equal(results[3].value, 'item-3');
assert.equal(peak, 2);
assert.equal(active, 0);
assert.deepEqual(await runPool([], 2), []);
const unusual = await runPool([
() => { throw null; },
() => { throw 'plain'; }
], 2);
assert.deepEqual(unusual, [
{ok: false, error: 'null'}, {ok: false, error: 'plain'}
]);
console.log(JSON.stringify(results));
console.log(`peak=${peak}, active=${active}`);为什么领取任务不需要额外加锁
最后一行应为 peak=2, active=0。上一行依次列出 item-0、item-1、demo failure 和 item-3。结果写入领取时的下标,所以完成顺序改变也不会错配。计时器只保证达到可执行条件后的调度,不承诺精确毫秒数;测试关注的是上限与归属。
这里的 next++ 和读取索引之间没有 await。在单个 JavaScript 执行代理内,这一小段同步代码会连续执行,另一个工作循环不会插进来领取同一项。若改成跨工作线程共享计数,或者把领取过程改成异步远程调用,就不能照搬这个解释。
把代码换成已经创建好的 Promise 数组会失去限流意义。调用请求函数时工作通常已经开始,后来再排队等待只是控制观察结果的时机。需要把发起动作包在函数里,只有工作循环实际调用这个函数,外部操作才开始。
并发数量与请求速率是两个约束
上限为二不等于每秒最多两次。任务如果很快完成,一秒仍可能启动大量请求。服务若限制单位时间内的调用量,还需要单独的速率控制,并结合服务返回的等待信息;不能把并发池当成通用的请求配额实现。
这个小实现会保留全部结果,也会继续执行剩余任务。它适合有限批次,不适合无限数据流。若任务量很大,应按批处理或及时消费结果;若业务要求首次失败后停止,需要明确哪些尚未开始、哪些正在运行,以及正在运行的动作能否取消。
接入真实业务前,再加入一个永不结束的任务进行设计检查:它会一直占用一个名额。超时应由任务自己明确实现,并在完成清理后才算释放名额。把“调用方不再等待”当成“底层工作已停止”,会让实际并发悄悄超过约定。


