事情起于一个做商品图片迁移的小工具。运营把 Excel 导进来,400 个 SKU,程序要给每个 SKU 拉一次远端图片、再传到自家对象存储。
我第一版写得很顺手:
const results = await Promise.all(
tasks.map(t => doUpload(t))
);
本地拿 10 条数据测,一切正常。到线上跑了 400 条,浏览器标签页直接卡死,控制台刷出一屏 net::ERR_INSUFFICIENT_RESOURCES,对象存储那边也顺手给我限了流——回头看日志,400 个请求在 800 毫秒内全发出去了。
问题不在 fetch,在这个写法本身没有节流概念。我需要四个东西:并发上限、失败隔离、能暂停、能取消。这四件事加起来大概一百来行原生代码,不需要 lodash,也不需要 p-limit。
先想清楚要约束什么
Promise.all 有三个毛病,平时写小工具不太显,一到批量就露出来:
- 没有并发上限。 有多少 task 就同时发多少个,浏览器对同域连接数有上限,超出的会被排队或者直接拒掉。
- 失败不隔离。 400 个里有一个 reject,整个
Promise.all立刻 reject,剩下的 399 个还在跑,但你已经拿不到它们的结果了。虽然可以用allSettled兜一部分,但它同样不限并发。 - 没法中途叫停。 用户点了取消,请求还在飞,网络面板里一排 pending 要挂好几分钟才结束。
下面这个队列就是对着这三条来的。我把完整实现拆成几段讲,最后拼起来。
用 Promise.withResolvers 把任务的「出口」提前拿出来
队列的本质是:调用方调 add(),队列把它塞进池子,什么时候执行由队列决定,但调用方要立刻拿到一个 Promise 去 await。
过去写这段要么用 new Promise(executor) 把 resolve 存到外部变量,要么用一个 events 库,都挺别扭。ES2024 加了 Promise.withResolvers(),专门解决这个:
const { promise, resolve, reject } = Promise.withResolvers();
// 这时候 promise 已经存在,resolver 也拿到手了
setTimeout(() => resolve('done'), 1000);
await promise;
放在队列里就是这样:
add(task, { name = '', signal } = {}) {
const { promise, resolve, reject } = Promise.withResolvers();
this.#queue.push({ task, name, taskSignal: signal, resolve, reject });
this.#pump();
return promise;
}
名字里的 Resolvers 是复数,因为它一次给你 resolve 和 reject 两个函数。我第一次看到还以为它返回数组,查了下文档才发现是对象。
兼容性上,Node 22、Chrome 119、Safari 17.4、Firefox 121 开始都有。要兼容更老的浏览器,polyfill 就三行:
if (!Promise.withResolvers) {
Promise.withResolvers = function () {
let resolve, reject;
const promise = new Promise((res, rej) => {
resolve = res;
reject = rej;
});
return { promise, resolve, reject };
};
}
并发槽位与调度循环
队列内部维护三个数字:并发上限、当前在跑的数量、等待队列。任务完成时减一,然后立刻去捞下一个。
#pump() {
if (this.#paused) return;
while (this.#running 0) {
const entry = this.#queue.shift();
this.#running++;
this.#run(entry);
}
}
注意 this.#run(entry) 前面没有 await。这里是故意的,调度循环要一次性把四个槽位填满然后立刻返回,不能等第一个任务跑完。因为 #run 是 async 函数,它内部所有的异常都会变成 rejected promise,不会同步抛出来,所以这里不会有「未捕获异常中断调度」的风险。
AbortSignal.any:一个任务,两个取消来源
取消这件事有两个层次:用户点「取消整个任务」要停掉所有请求,同时单个任务可能还有自己的超时信号。
以前要把两个信号合成一个,得手写 addEventListener('abort') 转发。现在一行:
const signals = [this.#controller.signal];
if (taskSignal) signals.push(taskSignal);
const signal = signals.length === 1 ? signals[0] : AbortSignal.any(signals);
AbortSignal.any 返回一个新信号,任意一个源信号 abort,它跟着 abort,reason 取最先触发的那个。Node 20.3、Chrome 116、Firefox 124、Safari 17.4 起可用。
有个细节得提:AbortSignal 是一次性的。一个 controller 调过一次 abort,它的信号就永久是 aborted 状态,后面的新任务如果还挂在这个信号上,会在启动的第一时间就被判定为已取消。所以队列的 abort 是「终止整个队列」,不是「跳过当前这批再接着用」。要重新开工,得 new 一个队列。这个后面还会说。
重试:指数退避加一点随机抖动
function backoffDelay(attempt, base = 300, cap = 5000) {
const exp = Math.min(cap, base * 2 ** (attempt - 1));
return exp + Math.random() * 0.3 * exp;
}
第一次重试等 300ms 上下,第二次 600ms,第三次 1200ms,最高封在 5 秒。后面那个 Math.random() * 0.3 是关键,假如 20 个任务同时因为对象存储 503 失败,没有抖动的话它们会在同一毫秒一起重试,把对方再打挂一次。
为什么不做固定间隔?因为我们面临的失败大多是限流和瞬时故障,固定间隔会让重试批次节奏完全对齐,越退避越撞车。
sleep 也必须能被中断
这个坑我在另一台机器上排查了半小时。现象是:用户点了取消,请求确实停了,但队列要等好几秒才真正安静下来。原因是重试逻辑里那个 setTimeout 没管取消,它在那儿老老实实睡满了退避时长。
function sleep(ms, signal) {
return new Promise((resolve, reject) => {
if (signal?.aborted) {
reject(signal.reason);
return;
}
const timer = setTimeout(() => {
signal?.removeEventListener('abort', onAbort);
resolve();
}, ms);
function onAbort() {
clearTimeout(timer);
reject(signal.reason);
}
signal?.addEventListener('abort', onAbort, { once: true });
});
}
定时器回调里那句 removeEventListener 别忘了。一个队列跑 400 个任务,每个任务平均重试一次,就是 400 个 abort 监听器挂在长生命周期的信号上,不摘掉会慢慢涨。
完整实现
把上面的东西拼起来,一百来行:
class TaskQueue {
#concurrency;
#maxRetries;
#queue = [];
#running = 0;
#paused = false;
#controller = new AbortController();
#onProgress;
constructor({ concurrency = 4, retries = 2, onProgress } = {}) {
this.#concurrency = Math.max(1, concurrency);
this.#maxRetries = Math.max(0, retries);
this.#onProgress = typeof onProgress === 'function' ? onProgress : () => {};
}
get signal() { return this.#controller.signal; }
get waiting() { return this.#queue.length; }
get active() { return this.#running; }
add(task, { name = '', signal } = {}) {
const { promise, resolve, reject } = Promise.withResolvers();
this.#queue.push({ task, name, taskSignal: signal, resolve, reject });
this.#pump();
return promise;
}
pause() {
this.#paused = true;
}
resume() {
if (!this.#paused) return;
this.#paused = false;
this.#pump();
}
clear(reason) {
const err = toAbortError(reason, '队列已清空');
for (const entry of this.#queue.splice(0)) {
entry.reject(err);
}
}
abort(reason) {
this.#paused = true;
const err = toAbortError(reason, '队列已终止');
this.#controller.abort(err);
this.clear(err);
}
#pump() {
if (this.#paused) return;
while (this.#running 0) {
const entry = this.#queue.shift();
this.#running++;
this.#run(entry);
}
}
async #run(entry) {
const { task, name, taskSignal } = entry;
const signals = [this.#controller.signal];
if (taskSignal) signals.push(taskSignal);
const signal = signals.length === 1 ? signals[0] : AbortSignal.any(signals);
let settled = false;
const finish = (ok, value) => {
if (settled) return;
settled = true;
this.#running--;
ok ? entry.resolve(value) : entry.reject(value);
this.#pump();
};
for (let attempt = 1; attempt this.#maxRetries;
const aborted = signal.aborted || err?.name === 'AbortError';
if (aborted || isLast) {
this.#onProgress({ type: 'fail', name, attempt, error: err });
finish(false, err);
return;
}
const delay = backoffDelay(attempt);
this.#onProgress({ type: 'retry', name, attempt, delay });
try {
await sleep(delay, signal);
} catch (abortErr) {
finish(false, abortErr);
return;
}
}
}
}
}
function toAbortError(reason, fallback) {
if (reason instanceof Error || reason instanceof DOMException) return reason;
return new DOMException(reason || fallback, 'AbortError');
}
function backoffDelay(attempt, base = 300, cap = 5000) {
const exp = Math.min(cap, base * 2 ** (attempt - 1));
return exp + Math.random() * 0.3 * exp;
}
function sleep(ms, signal) {
return new Promise((resolve, reject) => {
if (signal?.aborted) {
reject(signal.reason);
return;
}
const timer = setTimeout(() => {
signal?.removeEventListener('abort', onAbort);
resolve();
}, ms);
function onAbort() {
clearTimeout(timer);
reject(signal.reason);
}
signal?.addEventListener('abort', onAbort, { once: true });
});
}
里面那个 settled 标志位看着有点多余,但它防的是一种真实情况:任务在执行过程中被 abort 了,catch 里处理一次,如果外层再有兜底逻辑,就会重复扣减 #running,导致槽位被永久占用,队列跑着跑着就卡住了。加上这一句,代价为零。
用起来:批量拉图片
const queue = new TaskQueue({
concurrency: 3,
retries: 2,
onProgress: evt => {
if (evt.type === 'retry') {
console.log(`${evt.name} 第 ${evt.attempt} 次失败,${Math.round(evt.delay)}ms 后重试`);
}
},
});
const urls = [/* 400 条 */];
const jobs = urls.map((url, i) =>
queue.add(
async ({ signal, attempt }) => {
const res = await fetch(url, { signal });
if (!res.ok) {
throw new Error(`HTTP ${res.status}`, {
cause: { url, status: res.status, attempt },
});
}
return res.blob();
},
{ name: `img-${i}` }
)
);
const settled = await Promise.allSettled(jobs);
const ok = settled.filter(r => r.status === 'fulfilled').length;
const failed = settled.filter(r => r.status === 'rejected');
console.log(`成功 ${ok},失败 ${failed.length}`);
failed.slice(0, 5).forEach(r => console.dir(r.reason, { depth: 3 }));
这里用了 Error 的 cause 选项,把 url、状态码、第几次尝试挂在错误对象上。好处是排查的时候能直接 console.dir(err, { depth: 3 }) 看到完整上下文,不用另外维护一个日志数组。
取消操作也很直白:
cancelBtn.onclick = () => {
queue.abort('用户主动取消');
};
// 或者临时挂起,等一下再恢复
pauseBtn.onclick = () => queue.pause();
resumeBtn.onclick = () => queue.resume();
调完 abort() 之后,正在飞的那几个 fetch 会立刻 reject AbortError,还没出队的任务直接以同一个错误 reject,不会再有新请求发出去。整个过程在一帧之内结束。
七个真踩过的坑
1. abort 的 reason 别传字符串。
如果写 controller.abort('用户取消'),那 signal.reason 就是那个字符串,fetch 被中断时抛出来的也是这个字符串。err.name 是 undefined,任何基于 AbortError 名字的判断都会失效。要么不传参数,要么传 new DOMException('用户取消', 'AbortError')。上面代码里的 toAbortError 就是干这个的。
2. AbortSignal 不能复用。
这是最容易写出隐蔽 bug 的地方。第一版我在队列构造函数里创建 controller,想着「队列这一个信号管到底」。结果用户取消一次之后再往队列里 add 任务,新任务一启动就立刻被 abort,因为信号早就废了。
正确的做法是给「队列实例」划出生命周期边界:一个队列,一次任务批次。要重来就 new 一个。
3. 任务的 promise 没接住会报 unhandledrejection。
Promise.withResolvers() 出来的 promise 如果一直没人 await,reject 的时候浏览器会打印未处理的拒绝警告。这个不一定代表代码有 bug,但会让控制台很脏。稳妥起见,外面统一用 Promise.allSettled 兜住。
4. add 是同步启动的。
add() 里直接调了 #pump(),所以在并发没满的情况下,任务函数会同步进入执行路径。如果你期望「先收集完所有任务再统一开始」,得把 #pump 包一层 queueMicrotask。
5. 进度回调不要直接绑到 UI 更新上。
我一开始把 onProgress 直接写成了「每来一个事件就 setState 刷新进度条」。400 个任务,每个任务最少触发一次回调,多的三次,页面在跑任务期间几乎一直在重渲染。
后来改成两件事:成功和失败事件不做 UI 更新,只累加计数器;重试事件才输出日志。进度条用 setInterval 每 200ms 读一次计数。CPU 占用从 30% 掉到了个位数。
6. 并发数不是越大越好。
我一开始设了 8。测下来 400 个请求确实快了,但对象存储那边开始返 429。改成 3 之后总耗时反而缩短了,因为重试次数少了很多。并发数的合理值取决于下游能吃多少,不看本地机器性能。
7. 别把 CPU 密集型任务扔进这个队列。
这套东西是为 I/O 并发设计的。如果 task 里是纯计算(比如图片压缩、大 JSON 解析),并发数设多少都没意义,反而可能因为微任务堆积把主线程堵得更久。那类任务该丢给 Worker。
怎么验证它真的对
写完之后我用一个假任务测了两分钟,比人肉点页面靠谱:
const queue = new TaskQueue({ concurrency: 3, retries: 2, onProgress: e => log.push(e) });
const log = [];
const makeTask = (id, failTimes) => ({ signal, attempt }) =>
new Promise((resolve, reject) => {
if (signal.aborted) return reject(signal.reason);
const t = setTimeout(() => {
attempt {
clearTimeout(t);
reject(signal.reason);
}, { once: true });
});
// 20 个任务,前 5 个注定失败 2 次
const jobs = Array.from({ length: 20 }, (_, i) =>
queue.add(makeTask(i, i x.status === 'fulfilled').length);
console.log('失败', r.filter(x => x.status === 'rejected').length);
console.log('事件', log.map(e => e.type).join(','));
验证三个点:成功数是不是 20(因为重试次数够),失败任务的重试次数对不对,以及取消之后 #running 能不能回到 0。第三点特别重要,回不到 0 说明有槽位泄漏,队列跑久了必然卡死。
// 泄漏检测
const q2 = new TaskQueue({ concurrency: 2 });
q2.add(makeTask(1, 5), { name: 'never-succeed' });
setTimeout(() => q2.abort('test'), 100);
setTimeout(() => console.log('active =', q2.active), 5000); // 应该是 0
这段我跑了三个晚上,中间抓到过一次 #running 停在 1 不动的 bug,原因就是前面说的那个 settled 标志位还没加。
写在最后
这套代码不适合当通用库用,它没有做优先级队列,也没有任务依赖(A 完成才能跑 B)。真需要这些的话,现成库更合适。
它的价值在另一头:当你要做的只是一个后台小工具,或者一个功能页里的批量操作,一百行能搞定的事,没必要为了一个并发池去引入一个依赖。而且自己写的版本可控——重试策略想调就调,日志想打多细就打多细,出问题的时候你知道每一行在干什么。
回到开头,那 400 个 SKU 最后跑了 47 秒,中途失败 3 个自动重试成功。运营同事点了两次「取消」试了一下,两个请求瞬间停掉,然后她说了句「哦这个还挺跟手」。

