Skip to content
Promise 对象池是一种控制异步任务并发数量的模式——给定一组返回 Promise 的工厂函数和一个最大并发数 n,对象池保证同一时刻最多有 n 个任务处于执行状态,并在所有任务结束后收拢结果。这一模式常用于批量请求节制、后端流量整形等场景。
基本概念
- 任务(Task):一个返回 Promise 的工厂函数(
() => Promise<any>)。任务在调度之前不会执行,由调度器决定何时调用。 - 并发数(Concurrency):允许同时执行的任务数量上限。超过上限的任务必须等待已有任务完成才能启动。
- 工作进程(Worker):对象池内部的活动调度单元。每个 worker 负责从任务队列中取任务并执行,任务完成后循环获取下一个,直至队列为空。
- 迭代器(Iterator):用于在多个 worker 之间无锁地分发任务。共享的迭代器保证每个任务恰好被分配一次。
工作原理
对象池不预先创建任务队列的副本,而是让所有 worker 直接竞争同一迭代器。functions[Symbol.iterator]() 返回一个具有状态的迭代器对象,每次调用 next() 都会推进内部索引,不同 worker 天然无法拿到同一个任务。
每个 worker 是一个 async function,内部以无限循环的形式不断从迭代器中取出任务:
while (迭代器未耗尽) {
取出下一个任务
等待任务完成
}当迭代器指示 done: true 时跳出循环,worker 自行结束。由于 await 会释放事件循环,当前 worker 在执行任务时,其他 worker 可以继续从迭代器索取新任务。最终,n 个 worker 并发运行,确保同时执行的任务数始终不超过 n。
完成调度后,使用 Promise.allSettled(workers) 等待所有 worker 退出。这里选择 allSettled 而非 all,是因为单个任务的失败不应中断调度流程——若采用 all,一旦某个 worker 异常退出,其余仍在执行的任务结果将无法被收集。
基本用法
基础实现只管理并发调度,不负责收集每个任务的业务结果。如果需要收集结果,应在 worker 内部将结果推入外部数组,并且必须用 try...catch 包裹 await fn(),否则一个失败的任务会直接终止当前 worker,导致该 worker 后续再也无法领取新任务——这是调度正确性的关键边界。
ts
type F = () => Promise<any>
async function promisePool(functions: F[], n: number): Promise<PromiseSettledResult<any>[]> {
const results: PromiseSettledResult<any>[] = []
const iterator = functions[Symbol.iterator]()
async function worker() {
while (true) {
const { value: fn, done } = iterator.next()
if (done) break
try {
const value = await fn()
results.push({ status: 'fulfilled', value })
} catch (reason) {
results.push({ status: 'rejected', reason })
}
}
}
const workers = Array.from({ length: n }, () => worker())
await Promise.allSettled(workers)
return results
}调用方:
ts
const sleep = (t: number) => new Promise(res => setTimeout(res, t))
const tasks = [
() => sleep(500),
() => sleep(400),
() => sleep(600)
]
promisePool(tasks, 2)如果函数需要传参,可以在工厂函数中完成绑定:
ts
const fetchUser = (id: number) => () =>
fetch(`/api/user/${id}`).then(r => r.json())示例
示例 1:时间线验证
ts
const start = Date.now()
const sleep = (t: number) => new Promise(res => setTimeout(res, t))
const log = (msg: string) =>
() => new Promise(res => {
const elapsed = Date.now() - start
console.log(`${elapsed}ms: ${msg}`)
res(undefined)
})
const tasks = [
log('task1 start'), // 耗时 500ms
() => sleep(500).then(() => log('task1 end')()),
log('task2 start'), // 耗时 400ms
() => sleep(400).then(() => log('task2 end')()),
log('task3 start'), // 耗时 600ms
() => sleep(600).then(() => log('task3 end')())
]
promisePool(tasks, 2)并发数为 2 时的典型输出(实际时间可能受事件循环调度影响,略有偏差):
约0ms: task1 start
约0ms: task2 start
约500ms: task1 end
约500ms: task3 start // task1 结束后立即启动 task3
约400ms: task2 end // task2 本身只需 400ms,但可能稍晚于 task1 end 出现
约1100ms: task3 end示例 2:错误处理
ts
const failTask = () => Promise.reject(new Error('fail'))
const successTask = () => Promise.resolve('ok')
const tasks = [failTask, successTask]
promisePool(tasks, 2).then(console.log)
// [
// { status: 'rejected', reason: Error('fail') },
// { status: 'fulfilled', value: 'ok' }
// ]worker 内部通过 try...catch 捕获任务异常并推入结果数组,单个任务失败不会中断其他任务的调度。执行完毕后,results 包含所有任务的履行状态,无论成功还是失败。若将 try...catch 移除而改用 Promise.all 等待 worker,则任意一个任务失败都会导致整个调度短路,successTask 的结果将丢失。
注意点
- 任务结果并不保证按输入顺序排列。worker 是并发执行的,结果数组中的顺序取决于各自完成时刻。
- 当
n >= functions.length时,所有任务几乎同时启动,行为近似Promise.all(functions.map(f => f())),但调度细节略有不同。 - 迭代器在 worker 间共享,因此
functions数组本身不应在调度期间被修改。 - 需要中途取消时,基本实现无法胜任。
fn()一旦被调用,即便外部调整并发需求,该任务也无法被取消。如需取消能力,应结合AbortController或类似机制实现协作式取消。
限制
- 调度自身几乎没有性能开销,但所有 worker 都在同一个线程中运行,因此仅适用于异步 I/O 密集型任务。对于 CPU 密集型任务,需要借助 Worker Threads(Node.js)或类似方案。
- 上述实现中结果收集依赖于外部数组
results,多个 worker 同时推入结果,实际上存在并发写入风险。由于 JavaScript 单线程模型中push操作是同步的,且await之间不会交错执行,所以此处不存在真正竞态,但在心智模型上应意识到这是一个共享状态。
其他语言的类似机制
Java
Java 中对应的模型是 ThreadPoolExecutor,通过有界队列和工作线程控制并发:
java
ExecutorService executor = Executors.newFixedThreadPool(n);
List<Future<Object>> futures = tasks.stream()
.map(task -> executor.submit(() -> task.run()))
.collect(Collectors.toList());
// 收集结果
executor.shutdown();Java 线程池基于预创建的线程,可以同时执行 CPU 密集型和 I/O 操作,但调度粒度更重。
Python
Python 中可使用 asyncio.Semaphore 控制协程并发数:
python
import asyncio
async def pool(tasks, n):
sem = asyncio.Semaphore(n)
async def worker(task):
async with sem:
await task
await asyncio.gather(*(worker(t) for t in tasks))Semaphore 在作用上与 Promise 对象池的 worker 模式等价,都是通过信号量限制并发协程数量。
应用
- 批量网络请求节制:对速率有限的 API,限制同时发起的请求数,避免触发 429 错误。
- 文件批量处理:在 Node.js 中限制同时打开的文件句柄数量,防止“Too many open files”错误。
- 浏览器端资源加载:限制同时进行的
<script>动态加载或图片预加载数量,避免占用过多带宽。
