嵌套Promise.all代码未按预期异步执行问题排查
Node.js嵌套Promise.all同步执行问题排查与解决
问题现象
处理目录下多个CSV文件时,文件读取阶段是异步的(文件乱序完成读取),但进入行写入数据库环节后,脚本变为同步执行——同一文件的行被连续处理,日志示例:
...
Done filtering row 3327 in file 1
Done filtering row 3328 in file 1
Done filtering row 3329 in file 1
...
预期的理想状态是不同文件的行交替处理:
...
Done filtering row 120 in file 1
Done filtering row 2 in file 3
Done filtering row 121 in file 1
...
原业务代码(简化版)
const start = async () => { try { // 连接数据库 await connectDB(process.env.DATABASE_URL); // 获取目录下的CSV文件列表 const file_list = await getCSVFilesFromFolder('data'); await Promise.all(file_list.map(async (file_path, i) => { // 读取文件 const file = await readFileFromPath(file_path); console.log(`Done reading file ${i}`); // 解析CSV内容 const parsed_rows = await parseCSV(file); console.log(`Done parsing file ${i}`); // 循环处理每行数据,写入数据库 return Promise.all(parsed_rows.map((row, j) => { new Promise(async (resolve, reject) => { try { // 过滤数据 const filtered_data = filterData(row); const vat_number = `BE${row.BTWNUMMER}`; console.log(`Done filtering row ${j} in file ${i}`); // 查询公司是否存在 const company = await Company.findOne({ vat_number }); // 不存在则创建公司 if(!company) { await Company.create({ ...filtered_data }).catch(err => { failed_company_creations.push({ ...filtered_data, error: err?.message }) }) console.log(`Created company for row ${j} in file ${i}`); } else { console.log(`Company found for row ${j} in file ${i}`); } resolve(); } catch(err) { reject(err); } }); })); })); console.log('Done writing to database'); // 将失败记录写入CSV const failed_company_creations_csv = new ObjectsToCsv(failed_company_creations); await failed_company_creations_csv.toDisk('./failed_company_creations.csv'); console.log('Failed company creations successfully saved to disk'); } catch(error) { console.log(`Crashed`); console.log(error); } }; start();
测试验证结果
将文件操作和数据库逻辑替换为带随机延迟的模拟异步函数后,脚本恢复预期的异步交替处理,日志符合要求。测试代码如下:
const stall = async (stallTime = 3000) => { await new Promise(resolve => setTimeout(resolve, stallTime)); } const getRandomInt = (min, max) => { return Math.floor(Math.random() * (max - min + 1)) + min; } const start = async () => { try { // 模拟文件列表 const arrays = Array(20).fill().map((v,i)=>i); await Promise.all(arrays.map(async (file_path, i) => { // 模拟读取文件 const file = await stall(getRandomInt(100, 3000)); console.log(`Done reading file ${i}`); // 解析数据 const parsed_rows = await stall(getRandomInt(100, 3000)); console.log(`Done parsing file ${i}`); // 模拟行数据列表 const sub_arrays = Array(20).fill().map((v,i)=>i); // 循环处理每行 return Promise.all(sub_arrays.map((row, j) => { new Promise(async (resolve, reject) => { try { // 模拟过滤数据 const filtered_data = await stall(getRandomInt(100, 1000)); console.log(`Done filtering row ${j} in file ${i}`); // 模拟查询公司 const company = await stall(getRandomInt(50, 4000)); // 模拟创建公司 await stall(getRandomInt(50, 2000)); console.log(`Created company for row ${j} in file ${i}`); resolve(); } catch(err) { reject(err); } }); })); })); console.log('Done writing to database'); // 模拟写入失败记录 console.log('Failed company creations successfully saved to disk'); } catch(error) { console.log(`Crashed`); console.log(error); } }; start();
问题根因
- 数据库连接池限制:Mongoose(或同类ODM)默认连接池大小有限(通常为5),同一文件的行批量发起数据库请求时,会占满连接池,其他文件的请求只能排队等待,表现为同步执行。
- 同步前置操作干扰:
filterData是同步函数,会快速完成同一文件所有行的过滤,紧接着批量发起数据库请求,进一步挤占连接池资源。 - 冗余Promise包装:异步函数本身会返回Promise,手动创建
new Promise属于冗余操作,且可能导致并发逻辑异常。
解决方案
方案1:限制行处理并发数
使用限流工具(如p-limit)控制每行处理的并发数,让不同文件的请求能穿插执行:
const pLimit = require('p-limit'); const limit = pLimit(5); // 并发数建议与数据库连接池大小一致 // ... await Promise.all(file_list.map(async (file_path, i) => { // 读取、解析文件逻辑保持不变 // 用p-limit限制行处理并发 await Promise.all(parsed_rows.map((row, j) => limit(async () => { try { const filtered_data = filterData(row); const vat_number = `BE${row.BTWNUMMER}`; console.log(`Done filtering row ${j} in file ${i}`); const company = await Company.findOne({ vat_number }); if(!company) { try { await Company.create({ ...filtered_data }); console.log(`Created company for row ${j} in file ${i}`); } catch(err) { failed_company_creations.push({ ...filtered_data, error: err?.message }); } } else { console.log(`Company found for row ${j} in file ${i}`); } } catch(err) { console.error(`Error processing row ${j} in file ${i}:`, err); } }))); }));
方案2:调整数据库连接池大小
增大Mongoose连接池的poolSize配置,允许更多并发请求:
await connectDB(process.env.DATABASE_URL, { poolSize: 20 // 根据服务器和数据库性能调整 });
注意:此方法会增加数据库负载,需根据实际资源情况调整,避免压垮数据库。
方案3:移除冗余Promise包装
异步函数本身返回Promise,无需手动创建new Promise,简化代码同时避免潜在逻辑问题:
// 原代码中row处理部分改为: return Promise.all(parsed_rows.map(async (row, j) => { try { const filtered_data = filterData(row); const vat_number = `BE${row.BTWNUMMER}`; console.log(`Done filtering row ${j} in file ${i}`); const company = await Company.findOne({ vat_number }); if(!company) { try { await Company.create({ ...filtered_data }); console.log(`Created company for row ${j} in file ${i}`); } catch(err) { failed_company_creations.push({ ...filtered_data, error: err?.message }); } } else { console.log(`Company found for row ${j} in file ${i}`); } } catch(err) { console.error(`Error processing row ${j} in file ${i}:`, err); } }));
内容的提问来源于stack exchange,提问作者Thore
相关产品推荐
相关产品推荐

