浏览器被 400 个请求拖垮后,我用原生 JS 写了可取消、可重试的并发队列

事情起于一个做商品图片迁移的小工具。运营把 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 }));

这里用了 Errorcause 选项,把 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 个自动重试成功。运营同事点了两次「取消」试了一下,两个请求瞬间停掉,然后她说了句「哦这个还挺跟手」。

浏览器被 400 个请求拖垮后,我用原生 JS 写了可取消、可重试的并发队列
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请发送邮件1506151422@qq.com联系,将会及时下架删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 javascript 浏览器被 400 个请求拖垮后,我用原生 JS 写了可取消、可重试的并发队列 https://www.taomawang.com/web/javascript/2771.html

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务