Node.js 服务中的异步任务并发控制:从限流到失败重试的实现方法

AI智能摘要
Node.js异步批量任务中,直接用Promise.all会因同时启动全部调用而冲击下游并放大失败。可靠的并发控制应基于信号量模型,对正在执行的任务设硬上限并让其余排队;并发值应按下游承受能力设定,队列还需有长度边界和排队超时,并将排队超时与执行超时区分处理,重试也必须与限流统一设计。
— 此摘要由AI分析文章内容生成,仅供参考。

批量调用外部接口,或把一批本地任务丢进异步流程时,最省事的写法是先把 Promise 全部建出来,再用 Promise.all 等结果。任务少时看不出问题;量一上来,下游开始拒请求,本进程里同时挂起的调用把连接和内存占满,超时后又在同一时刻重试,失败被放大成第二波冲击。真正要设计的不是“怎么把任务并发出去”,而是给正在执行的数量设硬上限,让其余任务排队,再决定失败后如何重试、如何避免同一条业务被跑两遍。

把并发上限做成硬约束

并发控制的对象是正在执行的任务,不是已经创建的 Promise。Promise.all 在调用那一刻就把全部任务启动了,上限从一开始就不存在。更稳妥的模型是信号量:最多允许 N 个任务占用槽位,有空槽才启动下一个,跑完再释放。

下面这个限流器把等待队列和并发计数放在一起,接口保持成“包一层异步函数”,方便嵌进现有代码:

function createLimiter(concurrency) {
  if (concurrency < 1) throw new Error("concurrency must be >= 1");

  let active = 0;
  const waiting = [];

  const pump = () => {
    while (active < concurrency && waiting.length > 0) {
      const { task, resolve, reject } = waiting.shift();
      active += 1;
      Promise.resolve()
        .then(task)
        .then(resolve, reject)
        .finally(() => {
          active -= 1;
          pump();
        });
    }
  };

  const limit = (task) =>
    new Promise((resolve, reject) => {
      waiting.push({ task, resolve, reject });
      pump();
    });

  limit.stats = () => ({ active, waiting: waiting.length });
  return limit;
}

async function fetchUser(id) {
  // 这里换成真实的外部调用
  return { id };
}

async function runBatch(ids, concurrency) {
  const limit = createLimiter(concurrency);
  const jobs = ids.map((id) => limit(() => fetchUser(id)));
  return Promise.all(jobs);
}

concurrency 先按下游能承受的程度来定,而不是按本机还能再开多少 Promise。对外部 HTTP 来说,上限通常应对齐对方的配额、连接数和你这边客户端连接池;对 CPU 型任务,上限更接近可用核心数。一开始宁可比“感觉能跑满”更保守,因为后面的重试会把瞬时流量再抬一层。

limit.stats 不是装饰。队列长度持续上涨,说明生产速度超过消费速度;active 长期顶满而队列很短,说明瓶颈在下游耗时,继续加并发往往只是把超时和失败率一起抬上去。

排队要有边界,不能只是无限等待

限流器内部的 waiting 已经是队列。缺的是边界:队列要不要封顶、任务最多等多久、调用方要不要感知背压。没有边界的队列只是把打满下游,换成打满内存。

可以在入队处加上长度限制和等待超时:

function createLimiter(concurrency, { maxQueue = Infinity, enqueueTimeout = 0 } = {}) {
  let active = 0;
  const waiting = [];

  const pump = () => {
    while (active < concurrency && waiting.length > 0) {
      const item = waiting.shift();
      if (item.settled) continue;
      active += 1;
      Promise.resolve()
        .then(item.task)
        .then(item.resolve, item.reject)
        .finally(() => {
          active -= 1;
          pump();
        });
    }
  };

  const limit = (task) =>
    new Promise((resolve, reject) => {
      if (waiting.length >= maxQueue) {
        reject(Object.assign(new Error("queue is full"), { code: "QUEUE_OVERFLOW" }));
        return;
      }

      const item = { task, resolve, reject, settled: false };
      waiting.push(item);

      if (enqueueTimeout > 0) {
        setTimeout(() => {
          if (item.settled) return;
          item.settled = true;
          reject(Object.assign(new Error("enqueue timeout"), { code: "ENQUEUE_TIMEOUT" }));
        }, enqueueTimeout);
      }

      const wrap = (handler) => (value) => {
        if (item.settled) return;
        item.settled = true;
        handler(value);
      };
      item.resolve = wrap(resolve);
      item.reject = wrap(reject);
      pump();
    });

  limit.stats = () => ({ active, waiting: waiting.length });
  return limit;
}

队列满了是该让调用方立刻失败,还是堵住生产端,取决于任务从哪来。HTTP 请求进来再扇出一批外部调用,队列溢出更适合快速失败并返回可重试的错误,避免请求线程被无限挂起。内部批处理任务则可以让生产者在队列高位时暂停拉取,这才是背压。两种场景不要混用同一套“无限排队”默认值。

排队超时和执行超时也要分开。前者表示还没轮到它;后者表示已经占用了并发槽却迟迟不结束。把两种超时混成一个 timeout,后面排障时很难判断是容量不够,还是下游本身就慢。

重试会改变真实并发,必须和限流绑在一起设计

失败重试看起来是任务自己的事,实际上它会改写限流器的负载。一次调用失败后立刻再打,等于在原计划之外又插入新流量;如果很多任务在同一时刻失败,退避时间又相同,第二波会整齐地打出去,形成重试风暴。

先把可重试判断、指数退避和抖动写清楚:

function sleep(ms) {
  return new Promise((resolve) => setTimeout(resolve, ms));
}

function isRetryable(err) {
  if (!err) return false;
  if (err.name === "AbortError") return false;
  if (err.code === "QUEUE_OVERFLOW") return false;
  // 按你的调用约定扩展:限流、暂时不可用、网络抖动通常可重试
  // 参数错误、鉴权失败、业务拒绝通常不应重试
  return err.retryable === true;
}

async function withRetry(task, { retries = 3, baseDelay = 200, maxDelay = 5000 } = {}) {
  let attempt = 0;
  while (true) {
    try {
      return await task(attempt);
    } catch (err) {
      if (attempt >= retries || !isRetryable(err)) throw err;
      const exp = Math.min(maxDelay, baseDelay * 2 ** attempt);
      const jitter = exp * (0.5 + Math.random() / 2); // 50%~100% 抖动
      await sleep(jitter);
      attempt += 1;
    }
  }
}

抖动比“精确的指数曲线”更重要。一批评任务若共享同一套 baseDelay * 2 ** attempt,它们会在几乎相同的时间点再次出发;给延迟加随机窗口,才能把重试摊开。

更关键的是重试发生在限流器的哪一侧。

重试放在限流器里面,退避期间仍占着并发槽。好处是实现简单,坏处是一次缓慢失败就能把槽位冻住,队列里健康的任务跟着饿死,有效吞吐比你设置的 concurrency 低得多。

重试放在限流器外面,每次尝试都重新排队。槽位只在真正发请求时占用,退避期间把位置让给别人。代价是任务的总等待时间变长,而且重试会再次经过队列边界——这通常是你想要的,因为重试不该拥有比首次执行更高的特权。

推荐默认采用“每次尝试单独占槽”:

async function runWithLimitAndRetry(limit, task, retryOptions) {
  return withRetry(() => limit(task), retryOptions);
}

如果某类任务失败后必须立刻在同一条连接上续跑,再考虑槽内重试,并单独限制这类任务的数量。不要把“偶发网络抖动”和“必须连续会话”用同一套重试包起来。

超时还要能取消。只 Promise.race 一个计时器、不中止底层调用,原请求仍可能在超时后成功,重试就会变成两次有效执行:

async function withTimeout(factory, ms) {
  const controller = new AbortController();
  const timer = setTimeout(() => controller.abort(), ms);
  try {
    return await factory(controller.signal);
  } finally {
    clearTimeout(timer);
  }
}

factory 必须把 signal 传到实际的 HTTP 客户端或子任务里,超时才不是摆设。取消后的错误应视为不可盲目重试,除非你能确认对方没有落地副作用。

重复执行往往比失败更难收场

并发、超时、重试叠在一起时,同一条业务被执行两次是常态,不是意外。常见路径包括:第一次调用其实已经成功,但客户端超时后重试;两个请求同时通过了“尚未处理”的检查;队列里积了同一 id 的两条任务。外部写操作、扣减、发送通知这类动作,失败可以再补,重复却可能制造错账或骚扰。

先挡住飞行中的重复。同一进程里,相同幂等键的任务只允许一份在跑,后来者共用同一个 Promise:

function createInflight() {
  const inflight = new Map();

  return function runOnce(key, task) {
    const existing = inflight.get(key);
    if (existing) return existing;

    const pending = Promise.resolve()
      .then(task)
      .finally(() => {
        inflight.delete(key);
      });

    inflight.set(key, pending);
    return pending;
  };
}

这只能合并“此刻正在跑”的重复,挡不住“上次已经成功,这次又被重新投递”。所以业务层还需要幂等:把 idempotencyKey 传到下游,或在本地用任务键记录终态。读取类请求相对好办;写入类请求应假设至少一次投递,把处理函数写成再执行一次结果不变。

超时重试前要先问:能不能确认上次的结果?能查到成功就不要再发;查不到再带着同一幂等键重试。把“重试”写成“再提交一条新任务”,重复执行几乎必然出现。

任务键的粒度也要小心。用用户 ID 当键,会把该用户互不相关的操作串行化;用“用户 ID + 操作类型 + 业务单据号”才是在去重,而不是在制造新的全局锁。

参数如何拉动稳定性和吞吐量

这几个参数不是越大越好,它们在抢同一组资源:下游配额、本进程连接、任务可接受的等待时间。

参数调大时通常发生什么调小时通常发生什么
并发上限吞吐上升,下游压力、超时和连接占用一起上升更稳,队列变长,尾延迟变差
队列上限更能吸收突发,内存和等待时间变差更快背压,突发期失败率上升
最大重试次数瞬时抖动更容易扛过,重复执行和额外流量上升失败更快暴露,成功率可能下降
退避基数 / 上限重试更分散,单任务完成时间变长恢复快,但容易对齐成重试风暴
执行超时慢任务更少被误杀,槽位被慢调用占得更久槽位周转快,误杀后的重试变多

没有一组通用最佳值,但有比较稳的调节顺序。先把并发上限降到下游几乎不再返回限流或连接错误;再看队列长度,决定是给调用方快速失败,还是让批处理生产者减速;最后才加重试。反过来先把重试次数加到很高,只会用重复流量掩盖容量不足,排障时看到的全是“偶发失败”,根因其实是上限设错了。

吞吐不该只看每秒完成数。把重试计入后的有效成功数、以及 P95 完成时间,更能说明系统有没有在用等待和重复换成功率。并发从 5 调到 20,若成功数几乎不变、超时却明显增加,多出来的 15 个槽位只是在排队等待失败。

重试次数对稳定性的帮助是递减的。第一次重试常常能吃掉瞬时抖动;第三次、第四次多半已经撞上持续故障。持续 5xx 或持续限流时,继续重试是在延长故障面。这时应停止该批任务、把错误冒泡,而不是让限流器一直被重试占满。

可观察、可维护的任务流程

策略散落在调用点里,很快就会出现有的地方重试 3 次、有的地方不限流、有的地方超时后既不取消也不去重。把上限、排队、重试、去重收成一个 runner,调用方只提交“键 + 任务 + 超时”,日志和统计从同一处打出来。

function createRunner({ concurrency, retry, timeoutMs }) {
  const limit = createLimiter(concurrency, { maxQueue: 1000, enqueueTimeout: 10_000 });
  const runOnce = createInflight();

  async function exec(key, task) {
    const started = Date.now();
    const result = await runOnce(key, () =>
      withRetry(
        (attempt) =>
          limit(() =>
            withTimeout((signal) => task({ key, attempt, signal }), timeoutMs)
          ),
        retry
      )
    );
    return { result, ms: Date.now() - started };
  }

  exec.stats = limit.stats;
  return exec;
}

const runner = createRunner({
  concurrency: 8,
  timeoutMs: 3000,
  retry: { retries: 2, baseDelay: 200, maxDelay: 2000 },
});

async function syncUsers(ids) {
  const jobs = ids.map((id) =>
    runner(`user:${id}`, async ({ attempt, signal }) => {
      // 把 signal、attempt 带进实际请求,便于取消和打日志
      return fetchUser(id, { signal, attempt });
    })
  );

  const settled = await Promise.allSettled(jobs);
  const rejected = settled.filter((item) => item.status === "rejected");
  if (rejected.length) {
    throw new Error(`${rejected.length} tasks failed`);
  }
}

日志至少要能把一次执行还原成一条时间线:任务键、第几次尝试、入队等待了多久、实际调用花了多久、以成功、失败、取消还是去重合并结束。没有这些字段,线上只能看到“有的请求慢”,无法判断是队列堵住、下游变慢,还是重试在空转。

Promise.allSettled 适合批处理收口:先让能完成的完成,再决定失败项是记录后跳过,还是整批失败。Promise.all 在第一处拒绝时就短路,已经占用的槽位和飞行中请求并不会因此消失,只是调用方更早失去对整批结果的掌控。

可维护性还体现在错误分类稳定。QUEUE_OVERFLOWENQUEUE_TIMEOUT、不可重试的业务错误、可重试的暂时失败,应在 runner 边界上变成明确的 code,而不是一串无结构的 Error: failed。后续改并发或重试参数时,你才能从失败码分布判断是该降并发、加退避,还是修业务校验。

设计这类流程时,先写清四个决定:同时允许多少个在跑、队列满了怎么办、哪些失败值得重试以及重试是否重新排队、同一业务键如何避免双写。代码可以很短,但这四个决定含糊时,限流和重试只会把偶发问题变成稳定的重复执行。把 runner 的统计和失败码接到现有日志里,改参数时看队列长度、重试次数和有效成功率,而不是凭感觉把并发再加一档。

发布者:jacky,转转请注明出处:https://kubiyun.com/archives/4495

(0)
jacky的头像jacky
上一篇 2026-08-24 10:41
下一篇 2026-09-01 09:10

相关推荐

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注

评论列表(9条)

  • 瑶光夫人的头像
    瑶光夫人 2026-09-05 11:33

    退避不加抖动这点太真实了,之前重试全挤在一秒里打过去

  • 银月之翼的头像
    银月之翼 2026-09-05 13:39

    那个 runner 封装思路太香了,抄作业

  • 琴师小周的头像
    琴师小周 2026-09-06 12:01

    之前一直用 Promise.all,量大时确实容易把下游打挂

    • 天坛回音的头像
      天坛回音 2026-09-06 12:08

      @琴师小周同款踩坑,后来给调用加了并发上限才稳下来,你是用哪种方式限流的?

  • ToastyMarshmallow的头像
    ToastyMarshmallow 2026-09-06 15:36

    写操作一定要把幂等键带上

  • 糖人董的头像
    糖人董 2026-09-07 11:09

    只做 Promise.race 超时,底层请求其实还没停

    • 白玫瑰的头像
      白玫瑰 2026-09-07 11:51

      @糖人董对,其实只是把拿结果的那一层截断了,请求本身还在跑,所以超时+重试如果不控并发,下游压力会更大。