es-toolkit forEachAsync 深入指南:异步遍历数组与并发控制完整实战
【免费下载链接】es-toolkitA modern JavaScript utility library that's 2-3 times faster and up to 97% smaller, a major upgrade to lodash.项目地址: https://gitcode.com/GitHub_Trending/es/es-toolkit
导读
forEachAsync是 es-toolkit 数组模块(es-toolkit/array)提供的异步遍历工具:它对数组的每个元素执行一个返回 Promise 的非同步回调函数,并在所有异步操作全部完成时返回一个 Promise。与原生forEach不同,它会真正"等待"每个异步任务结束;同时通过可选的concurrency选项,你可以在不引入额外依赖的情况下限制并发数,轻松控制服务器或数据库的负载。读完本文,你将掌握forEachAsync的完整 API、并发控制原理(信号量 Semaphore)、错误传播行为,以及它在日志记录、文件上传、数据库批量更新等场景中的最佳实践。
一、为什么需要 forEachAsync
JavaScript 原生数组的forEach是一个同步方法,它不会等待回调函数中返回的 Promise:
// 原生 forEach:不会等待异步操作完成 users.forEach(async user => { await updateUser(user.id); // 这里的 await 对外层毫无意义 }); console.log('可能在任何用户更新完成前就执行了');这会导致"回调仍在执行,代码却已继续往下走"的经典问题。forEachAsync解决了这一痛点——它返回一个 Promise,只有所有元素的异步操作都完成后才会 resolve,从而保证后续代码(如"更新完成"提示、数据刷新)在正确的时机执行。
二、基本用法
2.1 函数签名
await forEachAsync(array, callback);完整签名如下(来自 src/array/forEachAsync.ts):
export async function forEachAsync<T>( array: readonly T[], callback: (item: T, index: number, array: readonly T[]) => Promise<void>, options?: ForEachAsyncOptions ): Promise<void>2.2 入门示例:批量更新用户信息
import { forEachAsync } from 'es-toolkit/array'; // 更新所有用户信息 const users = [{ id: 1 }, { id: 2 }, { id: 3 }]; await forEachAsync(users, async user => { await updateUser(user.id); }); // 到这里时,所有用户更新操作均已全部完成注意回调函数会收到与原生forEach一致的三个参数:当前元素item、元素索引index、原始数组array。测试 forEachAsync.spec.ts 中验证了这一点:
const arr = [1, 2, 3]; await forEachAsync(arr, callback); // callback.mock.calls[0] === [1, 0, arr] // callback.mock.calls[1] === [2, 1, arr] // callback.mock.calls[2] === [3, 2, arr]2.3 限制并发数
const items = [1, 2, 3, 4, 5]; await forEachAsync(items, async item => await processItem(item), { concurrency: 2 }); // 同一时刻最多只有 2 个 item 被处理concurrency选项用于控制同时运行的最大操作数,防止瞬时大量请求压垮服务器或数据库。它非常适合日志记录、文件上传、数据库更新这类不需要返回值的副作用操作。
2.4 串行执行
将concurrency设为1即可实现严格的顺序执行——一次只处理一个元素,处理完上一个才开始下一个:
import { forEachAsync } from 'es-toolkit/array'; // 顺序上传文件 const files = ['file1.txt', 'file2.txt', 'file3.txt']; await forEachAsync(files, async file => await uploadFile(file), { concurrency: 1 }); // 每个时刻只有一个文件在上传这在处理存在顺序依赖(如依赖前一步结果的状态流转)或需要温和地控制速率(rate limiting)的场景下非常实用。
三、参数详解
| 参数 | 类型 | 是否必填 | 说明 |
|---|---|---|---|
array | readonly T[] | 必填 | 需要遍历的数组 |
callback | (item: T, index: number, array: readonly T[]) => Promise<void> | 必填 | 对每个元素执行的异步函数,接收当前元素、索引与原始数组 |
options | ForEachAsyncOptions | 可选 | 控制并发的配置对象 |
options.concurrency | number | 可选 | 同时运行的最大操作数;不指定时所有操作同时执行(即完全并发) |
返回值:Promise<void>——当所有操作都完成时 resolve 的 Promise。这意味着你可以直接await forEachAsync(...),也可以将它放入Promise.all与其他异步任务并行编排。
四、源码级原理:一行代码背后的并发控制
forEachAsync的实现非常精简(见 src/array/forEachAsync.ts):
export async function forEachAsync<T>( array: readonly T[], callback: (item: T, index: number, array: readonly T[]) => Promise<void>, options?: ForEachAsyncOptions ): Promise<void> { if (options?.concurrency != null) { callback = limitAsync(callback, options.concurrency); } await Promise.all(array.map(callback)); }核心逻辑只有两步:
- 若不指定
concurrency:直接Promise.all(array.map(callback)),所有回调同时启动,等价于完全并发; - 若指定
concurrency:先用limitAsync把回调包装成"限流版本",再交给Promise.all调度。
4.1 limitAsync:用信号量包装回调
limitAsync位于 src/promise/limitAsync.ts,它创建一个计数信号量(Semaphore)并返回包装后的函数:
export function limitAsync<F extends (...args: any[]) => Promise<any>>(callback: F, concurrency: number): F { const semaphore = new Semaphore(concurrency); return async function (this: ThisType<F>, ...args: Parameters<F>): Promise<ReturnType<F>> { try { await semaphore.acquire(); return await callback.apply(this, args); } finally { semaphore.release(); } } as F; }关键设计点:release()被放在finally中,即使回调抛出异常,信号量也一定会被归还,不会因单个任务失败而"卡死"整个并发池。注意limitAsync也被 src/array/index.ts 单独导出,你可以直接复用它来限制任意异步函数的并发数。
4.2 Semaphore:FIFO 公平调度
信号量实现位于 src/promise/semaphore.ts,维护两个核心状态:
available:当前可用许可数量;deferredTasks:等待队列,采用 FIFO(先入先出)顺序,保证公平性。
acquire()的逻辑:若有可用许可则直接扣减并立即返回;否则把当前任务的 resolve 推入等待队列挂起。release()的逻辑:若等待队列非空,则取出队首任务并唤醒它(把许可"接力"给下一位),否则在不超过capacity的前提下归还许可。
这套机制保证了:同一时刻最多只有concurrency个回调在真正执行,其余回调按调用顺序排队等待。
五、行为边界与测试验证
仓库中的测试(forEachAsync.spec.ts)覆盖了五个关键行为,可作为使用时的行为契约参考:
| 测试场景 | 验证结论 |
|---|---|
| 对每个元素异步执行回调 | 回调被调用arr.length次,参数顺序为(item, index, array) |
| 空数组 | 回调调用次数为 0,函数正常 resolve |
| 任一回调抛出异常 | forEachAsync会 reject(错误传播) |
指定concurrency: 2 | 实测最大并发数maxRunning <= 2,且所有元素都被处理 |
不指定concurrency | 10 个元素的实测最大并发数为 10,即完全并发 |
其中两个边界值得特别注意:
- 空数组安全:
forEachAsync([], callback)不会调用任何回调,直接 resolve,可放心对可能为空的数组调用; - 错误快速失败(fail-fast):由于底层使用
Promise.all,一旦某个回调 reject,forEachAsync立即以该错误 reject,符合"任一失败即整体失败"的语义。如果你的场景需要"单个失败不中断整体",应先在回调内部自行捕获异常。
六、典型实战场景
结合concurrency的取值,forEachAsync可以覆盖从串行到完全并发的完整谱系:
- 数据库批量更新:
concurrency: 5左右,避免同时建立大量数据库连接导致连接池耗尽; - 文件批量上传/下载:
concurrency: 1或较小值,保护带宽与目标服务器; - 外部 API 调用:设置合理的并发上限,遵守第三方服务的速率限制(rate limit);
- 日志批量写入:无需返回值,遍历日志数组逐条落盘,可用较高并发提升吞吐;
- 清理任务、通知推送:不需要收集结果、只需"全部做完再继续"的批处理任务。
七、相关资源
- 中文文档主体来源:docs/ja/reference/array/forEachAsync.md(英文版见 docs/reference/array/forEachAsync.md)
- 源码实现:src/array/forEachAsync.ts
- 单元测试:src/array/forEachAsync.spec.ts
- 并发控制基础:src/promise/limitAsync.ts、src/promise/semaphore.ts
- 模块导出入口:src/array/index.ts
同类异步工具中,filterAsync 支持用异步谓词过滤数组并同样支持concurrency限制;如果你需要逆序遍历数组,可参考同步版本的 forEachRight。需要说明的是,forEachAsync的设计目标是"执行副作用并等待完成",它不收集回调的返回值——若你需要对结果做变换或收集,应优先选择mapAsync、filterAsync等返回数据集的工具。
【免费下载链接】es-toolkitA modern JavaScript utility library that's 2-3 times faster and up to 97% smaller, a major upgrade to lodash.项目地址: https://gitcode.com/GitHub_Trending/es/es-toolkit
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考