并发控制是异步工程的"压强管理":当一次性发起 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 管兜底。错误传播模式可以总结为三句话:
- 让错误带着上下文向上走(TaskError);
- 在恰当的一层统一收集(批处理层),不要在叶子层吞掉;
- 重试只做幂等操作,带退避与最大次数。
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 一梭子",而是能精确回答:同一时刻多少个、坏了怎么办、满了怎么办、用户走了怎么办。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。