TypeScript 异步并发控制:Promise 语义、AbortController 取消、并发池与任务队列

深入 TypeScript 异步并发控制:Promise 状态机与微任务语义、AbortController 取消传播、手写 p-limit 并发池、任务队列与 async/await 错误传播模式。

并发控制是异步工程的"压强管理":当一次性发起 1000 个网络请求、批量处理 10 万条记录时,无约束的并发会把 CPU、文件描述符、数据库连接池全部压垮。真正的问题不是"能不能并发",而是"同一时刻并发多少、超时的怎么办、出错的如何收尾"。

本文以 TypeScript 为语言载体,从 Promise 的底层语义出发,覆盖取消(AbortController)、并发池(p-limit 的原理与实现)、任务队列与错误传播模式——全部配真实可运行的代码。


1. Promise 底层语义

1.1 状态机与微任务

Promise 是三个状态的有限状态机:pending → fulfilled | rejected,状态一旦迁移就不可逆。resolve/reject 后注册的 then 回调不会同步执行,而是进入微任务队列(microtask queue),在下一个宏任务之前按 FIFO 依次执行:

console.log('sync: 1');

const p = Promise.resolve('resolved');
p.then((v) => console.log('microtask:', v));

console.log('sync: 2');

// 输出顺序:sync:1 → sync:2 → microtask:resolved

理解微任务是理解"并发调度"的前提:Promise.all 的多个异步任务本质上是同步发起了多个 I/O,结果按各自的完成顺序通过微任务回归。

1.2 then 链与错误传播

then 返回的是新的 Promise,回调的返回值会被展平(flattening):返回 Promise 则等待其落地,返回普通值则直接包装。错误在链上一直向后传播,直到遇到 catch:

async function fetchUser(id: number) {
  if (id < 0) throw new Error('非法 id');
  return { id };
}

fetchUser(-1)
  .then((u) => {
    console.log('永远不会到这里');
    return u.id;
  })
  .catch((err) => {
    console.log('错误在这里被捕获:', err.message);
    return { fallback: true }; // catch 之后链继续,值是 fallback
  })
  .then((result) => console.log('链继续:', result));

catch 之后链可以继续是设计上最重要的特性:错误处理是链的一部分,不是终点。

1.3 并发组合原语

Promise.all / allSettled / race / any 是并发的四个基础"阀门":

// Promise.all —— 全部成功才算成功,一个失败立即整体失败(fail-fast)
const results = await Promise.all([
  fetch('/a').then((r) => r.json()),
  fetch('/b').then((r) => r.json()),
]);

// Promise.allSettled —— 每个结果都保留,适合"部分失败不影响整体"的批处理
const settled = await Promise.allSettled([
  Promise.resolve(1),
  Promise.reject(new Error('boom')),
]);
// => [{ status: 'fulfilled', value: 1 }, { status: 'rejected', reason: Error }]

四个原语的取舍表:

原语语义适用场景
all全成或全败组合请求,缺一不可
allSettled各自成败批量任务,逐条记录失败
race首个落定超时、竞态
any首个成功多个候选源取其一

2. AbortController 与取消

2.1 AbortSignal 基础

传统上 Promise 一旦发起就"无法取消",AbortController 给出了标准化的取消机制:通过 signal 把取消事件传播给异步操作,操作方监听 abort 事件并主动 reject:

const controller = new AbortController();
const { signal } = controller;

// 监听方:一旦 abort,立即 reject
const waitable = new Promise<void>((resolve, reject) => {
  const onAbort = () => {
    signal.removeEventListener('abort', onAbort);
    reject(new DOMException('操作被取消', 'AbortError'));
  };
  signal.addEventListener('abort', onAbort, { once: true });
  // ... 正常完成时调用 resolve()
});

// 触发取消
controller.abort();

signal.aborted 可以同步判断是否已取消,用于进入操作前就提前返回。

2.2 与 fetch 集成

fetch 原生支持 signal,取消会以 AbortError 的形式 reject:

async function fetchWithTimeout(url: string, ms = 8000) {
  const controller = new AbortController();
  const timer = setTimeout(() => controller.abort(), ms);
  try {
    const res = await fetch(url, { signal: controller.signal });
    if (!res.ok) throw new Error(`HTTP ${res.status}`);
    return await res.json();
  } catch (err) {
    if (err instanceof DOMException && err.name === 'AbortError') {
      throw new Error(`请求超时(${ms}ms): ${url}`);
    }
    throw err;
  } finally {
    clearTimeout(timer);
  }
}

关键细节:超时用 AbortController 而不是 Promise.race。race 只是"忘记"慢请求,无法真正取消它背后的网络连接与资源占用;abort 才是真正的取消。

2.3 自定义异步操作的可取消封装

为自己的异步函数接入取消,核心模式是"监听 signal + 清理监听器":

interface CancellableOptions {
  signal?: AbortSignal;
}

function makeCancellable<T>(
  executor: (options: CancellableOptions) => Promise<T>,
): { promise: Promise<T>; cancel: () => void } {
  const controller = new AbortController();

  const promise = new Promise<T>((resolve, reject) => {
    const { signal } = controller;

    if (signal.aborted) {
      reject(new DOMException('已取消', 'AbortError'));
      return;
    }

    signal.addEventListener('abort', () => {
      reject(new DOMException('已取消', 'AbortError'));
    }, { once: true });

    executor({ signal }).then(resolve, reject);
  });

  return {
    promise,
    cancel: () => controller.abort(),
  };
}

// 用法:一个可能很慢的批量任务,可以随时取消
const { promise, cancel } = makeCancellable(async ({ signal }) => {
  const items = await fetchItems();
  return items.filter((it) => !signal.aborted); // 取消后跳过无用功
});

// 5 秒后用户离开页面 → 取消
setTimeout(cancel, 5000);

3. 并发池:p-limit 原理与手写实现

3.1 信号量思想

并发池的本质是计数器信号量(Semaphore):维护"当前在飞"的任务数,达到上限时后续任务排队,任务完成一个就放行一个。p-limit 的核心就是这个循环:

// 伪代码:信号量核心
let activeCount = 0;
const waitingQueue: Array<() => void> = [];

async function enqueue(task: () => Promise<void>) {
  if (activeCount >= limit) {
    // 排队:等待一个"空位"信号
    await new Promise<void>((resolve) => waitingQueue.push(resolve));
  }
  activeCount++;
  try {
    await task();
  } finally {
    activeCount--;
    // 释放一个空位
    waitingQueue.shift()?.();
  }
}

3.2 完整实现:带泛型返回值的并发池

生产级的并发池需要保留每个任务的返回值、传递拒绝、正确处理 finally 释放:

type Task<T = unknown> = () => Promise<T>;

interface PoolStats {
  active: number;
  waiting: number;
}

class ConcurrencyPool {
  private activeCount = 0;
  private waitingQueue: Array<() => void> = [];
  private readonly limit: number;

  constructor(limit: number) {
    if (!Number.isInteger(limit) || limit <= 0) {
      throw new Error('limit 必须是正整数');
    }
    this.limit = limit;
  }

  async run<T>(task: Task<T>): Promise<T> {
    // 排队直到有空位
    if (this.activeCount >= this.limit) {
      await new Promise<void>((resolve) => this.waitingQueue.push(resolve));
    }
    this.activeCount++;

    try {
      return await task();
    } finally {
      this.activeCount--;
      // 释放一个排队者
      const next = this.waitingQueue.shift();
      next?.();
    }
  }

  get stats(): PoolStats {
    return { active: this.activeCount, waiting: this.waitingQueue.length };
  }
}

这个 ConcurrencyPool 就是 p-limit 的等价物,只是把 p-limit 的"函数式入口"封装成了类。它的关键性质:

  • 每个 run 返回任务自身的 Promise,失败会原样向上传播;
  • finally 保证无论成功失败都释放空位,不会死锁;
  • 排队采用 FIFO,任务按到达顺序获得执行权。

3.3 p-limit 风格:函数式入口

p-limit 的 API 更函数化——传入 concurrency 返回一个"限流后的执行函数":

import pLimit from 'p-limit';

// 生产上直接用 p-limit
const limit = pLimit(5);

const urls = Array.from({ length: 100 }, (_, i) => `/api/item/${i}`);
const results = await Promise.all(
  urls.map((url) => limit(() => fetchJson(url))),
);

// 也可以自己写函数式版本,与上面类实现等价
function createLimiter(limit: number) {
  const pool = new ConcurrencyPool(limit);
  return <T,>(task: Task<T>): Promise<T> => pool.run(task);
}

3.4 并发池 + 超时 + 取消的三合一

生产场景并发池往往要叠加超时与取消。把第 2 节的取消封装接入池:

async function runWithPool<T>(
  tasks: Array<() => Promise<T>>,
  { limit, timeoutMs }: { limit: number; timeoutMs?: number },
): Promise<Array<T>> {
  const pool = new ConcurrencyPool(limit);
  return Promise.all(
    tasks.map((task) =>
      pool.run(() => (timeoutMs ? withTimeout(task(), timeoutMs) : task())),
    ),
  );
}

注意:Promise 无法被"强制打断",AbortController 必须由每个异步操作自愿监听才能彻底取消。并发池层面的取消,本质是"停止派发新任务 + 通知在飞任务"。


4. 任务队列

4.1 串行 FIFO 队列

并发池限的是"并行度",任务队列限的是"顺序"——按入队顺序逐个执行。最简串行队列:

class TaskQueue {
  private chain: Promise<unknown> = Promise.resolve();

  enqueue<T>(task: Task<T>): Promise<T> {
    const result = this.chain.then(() => task());
    // 无论成败,链都继续推进(错误由调用方各自处理)
    this.chain = result.catch(() => {});
    return result;
  }
}

const queue = new TaskQueue();
queue.enqueue(() => writeLog('A'));  // 串行执行
queue.enqueue(() => writeLog('B'));

这个实现的精髓是 this.chain = result.catch(() => {}):即使任务 A 失败,队列也不中断,B 照常执行;每个调用方拿到的是自己任务的结果 Promise。

4.2 带并发度的"队列 + 池"混合

真正的生产队列(如 kafka consumer、批处理 worker)需要"排队 + 限并发":

class ParallelQueue {
  private waiting: Task[] = [];
  private active = 0;
  constructor(private readonly concurrency: number) {}

  push<T>(task: Task<T>): Promise<T> {
    return new Promise<T>((resolve, reject) => {
      this.waiting.push(() =>
        task().then(resolve, reject),
      );
      this.pump();
    });
  }

  private pump() {
    while (this.active < this.concurrency && this.waiting.length > 0) {
      const next = this.waiting.shift()!;
      this.active++;
      next().finally(() => {
        this.active--;
        this.pump();
      });
    }
  }

  get length() { return this.waiting.length; }
}

4.3 背压(Backpressure)与暂停/恢复

当"生产者速度 > 消费者速度"时,需要背压:队列长度超过阈值就暂停消费,落回阈值之下再恢复。同时支持显式暂停/恢复:

class BackpressuredQueue extends ParallelQueue {
  private paused = false;
  private readonly highWaterMark: number;

  constructor(concurrency: number, highWaterMark = 1000) {
    super(concurrency);
    this.highWaterMark = highWaterMark;
  }

  push<T>(task: Task<T>): Promise<T> {
    if (!this.paused && this.length >= this.highWaterMark) {
      this.pause();
    }
    const result = super.push(task);
    void result.finally(() => this.maybeResume());
    return result;
  }

  pause() { this.paused = true; }
  resume() { this.paused = false; this.maybeResume(); }

  private maybeResume() {
    if (this.paused && this.length <= this.highWaterMark / 2) {
      this.resume();
    }
  }
}

背压的窗口(watermark 与恢复阈值之间)避免了"反复暂停-恢复"的抖动。这是处理"无限数据源 + 有限消费能力"的核心机制。


5. async/await 错误传播模式

5.1 错误分类与结构化捕获

并发编程的错误传播最大的问题是"错误失去了上下文"——不知道是哪个任务、哪个批次出的错。解决方式是结构化错误对象:

class TaskError<T> extends Error {
  constructor(
    public readonly index: number,
    public readonly input: T,
    public readonly cause: unknown,
  ) {
    super(`任务 ${index} 失败: ${cause instanceof Error ? cause.message : String(cause)}`);
    this.name = 'TaskError';
  }
}

// 用 allSettled 逐条收集失败,保留上下文
async function runBatch<T, R>(
  items: T[],
  fn: (item: T) => Promise<R>,
): Promise<{ succeeded: Array<{ index: number; value: R }>; failed: Array<TaskError<T>> }> {
  const settled = await Promise.allSettled(items.map((item, index) => fn(item)));
  const succeeded: Array<{ index: number; value: R }> = [];
  const failed: Array<TaskError<T>> = [];

  settled.forEach((result, index) => {
    if (result.status === 'fulfilled') {
      succeeded.push({ index, value: result.value });
    } else {
      failed.push(new TaskError(index, items[index], result.reason));
    }
  });

  return { succeeded, failed };
}

5.2 错误边界:在哪里 catch

一个常见错误是"在任务函数内部 catch 又 throw",把原始错误吞掉。正确的分层是:叶子函数只抛领域错误,批处理器统一收集,调用方决定是否重试:

// 叶子:只抛明确的错误
async function apiCall(id: number): Promise<Payload> {
  const res = await fetch(`/api/${id}`);
  if (!res.ok) throw new Error(`API 返回 ${res.status}`);
  return res.json();
}

// 批处理器:收集 + 可选重试(指数退避)
async function runWithRetry<T>(
  fn: () => Promise<T>,
  { retries = 3, baseDelay = 200 }: { retries?: number; baseDelay?: number } = {},
): Promise<T> {
  let lastErr: unknown;
  for (let attempt = 0; attempt <= retries; attempt++) {
    try {
      return await fn();
    } catch (err) {
      lastErr = err;
      if (attempt < retries) {
        await delay(baseDelay * 2 ** attempt); // 指数退避
      }
    }
  }
  throw lastErr;
}

5.3 并发任务的整体容错

并发 + 重试的完整组装:池限流 → 批收集 → 单任务重试:

const pool = new ConcurrencyPool(10);

const results = await Promise.all(
  items.map((item) =>
    pool.run(() =>
      runWithRetry(() => apiCall(item.id), { retries: 2 })
        .catch((err) => ({ ok: false as const, id: item.id, err })),
    ),
  ),
);

const succeeded = results.filter((r) => r.ok);
const failed = results.filter((r) => !r.ok);

这里每一层职责清晰:pool 管并行度,runWithRetry 管韧性,外层 catch 管兜底。错误传播模式可以总结为三句话:

  1. 让错误带着上下文向上走(TaskError);
  2. 在恰当的一层统一收集(批处理层),不要在叶子层吞掉;
  3. 重试只做幂等操作,带退避与最大次数。

6. 最佳实践清单

场景推荐手段避免
并发请求组合Promise.all + 池无界 Promise.all 裸奔
批量任务部分失败可容忍Promise.allSettled + 结构化错误一错全挂
超时控制AbortController + 定时器Promise.race 假超时
严格限制并行度并发池(p-limit / 自实现)手写计数器漏释放
严格串行任务队列(chain)递归地狱
生产速度 > 消费速度背压队列 + watermark无限堆积内存

最重要的几条原则:finally 释放资源(空位、定时器、监听器)——并发池死锁十有八九是漏了 finally;取消是"协作式"的,必须让异步操作自愿监听 signal;大任务拆小批次 + 结构化错误,可观测性远胜于一个模糊的 catch。

想深入 Promise 底层与运行时行为,可参考 https://plumephp.com/typescript-nodejs-backend/ 中对 Node 事件循环与异步 I/O 的讨论;把并发控制与运行时校验结合,可参考 https://plumephp.com/typescript-runtime-validation-typesafe/ 中的 API 边界实践。


7. 总结

异步并发控制的本质是在"吞吐量、延迟、资源、韧性"之间做显式取舍。Promise 给了原语,AbortController 给了取消,并发池给了压强管理,任务队列给了顺序与背压,错误传播模式给了可观测性。

把这几件工具装进工具箱,你就不再是"遇到并发就 Promise.all 一梭子",而是能精确回答:同一时刻多少个、坏了怎么办、满了怎么办、用户走了怎么办。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「frontend」更多文章

  1. CSS 架构与样式方案:从方法论到现代 CSS 新特性
  2. 可访问性与国际化:WCAG 2.2、ARIA 与 i18n 工程实践
  3. SSR/SSG 渲染模式全景:Next.js App Router、流式渲染与岛屿架构