- 同时运行的任务数永远不超过 2;
- 一个任务结束后,队列里下一个任务要立刻补位;
- 每个
add出去的任务,调用方还能拿到它自己的结果或错误。
先看最终效果。5 个耗时不同的任务,限制并发为 2:
start A
start B
end B
start C
end A
start D
end C
start E
end E
end D
results: [ 'A', 'B', 'C', 'D', 'E' ]
peak concurrent: 2
任意时刻最多只有 2 个 start 没有对应 end,并且结果按添加顺序返回。下面从零把它拆出来。
一、先明确需求
把模糊的口语题翻译成可验收的规格:
- 提供一个
add(task)方法,task是一个返回 Promise 的函数(不是 Promise 本身,这点很重要,后面会讲)。 add返回一个 Promise,task成功时它 resolve,失败时它 reject。- 全局任意时刻正在执行的任务数 ≤
limit。 - 前一个任务结束(无论成功还是失败),都要触发下一次调度。
这里最容易踩的坑:传函数而不是传 Promise。 如果直接传一个 Promise 进来,它在你 add 之前就已经开始执行了,调度器根本没机会控制它。
二、错误示范:Promise.all 不是并发控制
很多人第一反应是 Promise.all:
await Promise.all(tasks.map((task) => task()))
这段代码会把所有任务同时启动,Promise.all 只是等它们全部结束,它不做任何限流。任务一多,并发请求瞬间打满,接口限流或浏览器同域连接数限制就会教做人。
另一个常见的"伪并发控制"是分批:
for (let i = 0; i < tasks.length; i += 2) {
await Promise.all(tasks.slice(i, i + 2).map((t) => t()))
}
它确实做到了「每批最多 2 个」,但问题是慢任务会拖住整批。假设前一批里有一个 10 秒的任务,另一个任务 100ms 就结束了,那个空出来的槽位必须等满 10 秒才能接收下一批。这不是调度器,这是批处理。
真正的并发控制是:任何一个槽位空出来,就马上拉一个排队任务补进去,不用等其他人。
三、核心思路:队列 + 运行计数
整个调度器只需要两个状态:
queue:等待执行的任务队列(FIFO)。running:当前正在执行的任务数量。
调度逻辑就是一句话:只要 running < limit 且队列非空,就取出一个任务执行,并把 running 加一;任务结束时 running 减一,再递归地尝试调度。
把「任务结束」和「尝试调度」绑定在一起,就自然实现了"谁结束谁补位",不依赖批次。
四、实现:Scheduler 类
class Scheduler {
constructor(limit = 2) {
this.limit = limit;
this.running = 0;
this.queue = [];
}
add(task) {
return new Promise((resolve, reject) => {
this.queue.push({ task, resolve, reject });
this._next();
});
}
_next() {
while (this.running < this.limit && this.queue.length) {
const { task, resolve, reject } = this.queue.shift();
this.running++;
Promise.resolve()
.then(task)
.then(resolve, reject)
.finally(() => {
this.running--;
this._next();
});
}
}
}
几处关键设计:
add把task和它对应的resolve/reject一起存进队列,这样每个任务都能独立地把结果交还给调用方。add里先入队再调用_next(),保证无论当前有没有空闲槽位,逻辑都统一。_next()用while而不是if:一次调用尽量填满所有空槽位。调用时要么是刚add进来一个任务,要么是正好空出一个槽位,两种情况while都能正确处理。Promise.resolve().then(task)而不是直接task():一是能捕获task内部的同步抛错,二是保证任务在微任务里异步启动,add的同步代码不会被打断。
五、执行过程复盘
用 5 个任务验证,各自耗时 A=300ms B=200ms C=100ms D=400ms E=150ms:
const scheduler = new Scheduler(2);
const tasks = [
['A', 300],
['B', 200],
['C', 100],
['D', 400],
['E', 150],
];
const results = await Promise.all(
tasks.map(([name, ms]) =>
scheduler.add(async () => {
console.log(`start ${name}`);
await sleep(ms);
console.log(` end ${name}`);
return name;
})
)
);
执行时间线:
| 时刻 | 事件 | 正在运行 |
|---|---|---|
| 0 | A、B 启动,填满 2 个槽位 | A, B |
| 200 | B 结束,C 补位 | A, C |
| 300 | A 结束,D 补位 | C, D |
| 400 | C 结束,E 补位 | D, E |
| 550 | E 结束 | D |
| 700 | D 结束 | 空 |
任意时刻并发都不超过 2,且每个任务一结束,排队的下一个立刻顶上。这就是调度器和分批的本质区别。
六、错误处理:失败不能拖垮队列
面试的第二个考察点:如果某个任务 reject 了,队列会不会卡死?
答案是不会,因为回收计数和触发下一次调度放在 .finally() 里,而不是 .then() 的成功分支里。.finally 无论成功失败都会执行,所以 running-- 和 _next() 一定会跑到。
验证一下,第 2 个任务抛错:
const settled = await Promise.allSettled([
scheduler.add(async () => { await sleep(50); return 'ok-1' }),
scheduler.add(async () => { await sleep(10); throw new Error('boom-2') }),
scheduler.add(async () => { await sleep(10); return 'ok-3' }),
]);
console.log(
settled.map((x) => (x.status === 'fulfilled' ? x.value : x.reason.message))
);
输出:
[ 'ok-1', 'boom-2', 'ok-3' ]
boom-2 只让对应的那个 Promise reject,失败任务的槽位被正常回收,ok-3 照常执行。
一个重要细节:这里用的是
Promise.allSettled而不是Promise.all。如果用Promise.all,任意一个add的 Promise reject,整个await会立刻抛出,虽然队列实际还在跑,但你的等待逻辑已经提前退出了。要感知全部任务结果,用allSettled更合适。
七、面试进阶:泛化 limit + 保序返回
如果面试官接着问「那再实现一个 runPool(tasks, limit),一次性跑完并按原顺序返回结果」,可以用固定 worker 数的写法:
async function runPool(tasks, limit) {
const results = new Array(tasks.length);
let cursor = 0;
async function worker() {
while (cursor < tasks.length) {
const index = cursor++; // 先取号,保证不会重复消费
results[index] = await tasks[index]();
}
}
await Promise.all(
Array.from({ length: Math.min(limit, tasks.length) }, worker)
);
return results;
}
验证保序和并发:
const results = await runPool(tasks, 2);
console.log(results); // [ 'task-0', 'task-1', 'task-2', 'task-3', 'task-4' ]
console.log(peak); // 2
两种实现思路的差别:
- Scheduler 类:支持运行中动态
add,适合任务源源不断、边生产边消费的场景(比如上传队列)。 - runPool 固定 worker:任务集合已知,一次性并发执行并收集结果,代码更短,也能自然保序。
cursor++ 这一句是保序的关键:index 在自增前被取走,多个 worker 并发时每个序号只会被消费一次,results[index] 又保证结果落在正确位置。
八、常见追问
Q:为什么 add 一定要返回 Promise?
调度器本质是延迟执行,调用方需要知道「我这个任务什么时候好、结果是什么、有没有失败」。返回 Promise 把控制权交还给调用方,而不是让调度器默默吞掉结果。
Q:为什么传函数而不是传 Promise? Promise 在构造的那一刻就已经开始执行了,调度器拦不住。传函数(thunk)才能把启动时机交给调度器。
Q:_next() 里为什么用 while 而不是 if?
初始 add 一批任务时,第一次调用可能需要连续填满多个槽位;用 while 一次填满,逻辑更简单。单个补位时它也只循环一次,没有副作用。
Q:如果 task 内部是同步代码会怎样?
Promise.resolve().then(task) 会把它推入微任务执行,既能捕获同步抛错,也不会阻塞当前同步逻辑。
九、完整代码
class Scheduler {
constructor(limit = 2) {
this.limit = limit;
this.running = 0;
this.queue = [];
}
add(task) {
return new Promise((resolve, reject) => {
this.queue.push({ task, resolve, reject });
this._next();
});
}
_next() {
while (this.running < this.limit && this.queue.length) {
const { task, resolve, reject } = this.queue.shift();
this.running++;
Promise.resolve()
.then(task)
.then(resolve, reject)
.finally(() => {
this.running--;
this._next();
});
}
}
}
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
const scheduler = new Scheduler(2);
const results = await Promise.all(
['A', 'B', 'C', 'D', 'E'].map((name, i) =>
scheduler.add(async () => {
console.log(`start ${name}`);
await sleep((i + 1) * 100);
console.log(` end ${name}`);
return name;
})
)
);
console.log(results);
十、总结
Promise.all/ 分批Promise.all都不是并发控制;前者全量并发,后者会被慢任务拖住整批。- 正确的模型是队列 + 运行计数,谁结束谁补位,任意时刻运行数不超过
limit。 - 回收计数和触发下一次调度必须放在
.finally(),保证失败任务也能释放槽位。 add返回 Promise、任务传函数(thunk)而不是 Promise,是接口设计上的两个关键点。- 需要动态追加任务用
Scheduler类;任务集合已知、要保序返回用固定 worker 的runPool。 - 面试加分项:主动提错误隔离、保序、以及
limit泛化,说明你真的在生产里用过,而不只是背了段代码。