You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

嵌套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();

问题根因

  1. 数据库连接池限制:Mongoose(或同类ODM)默认连接池大小有限(通常为5),同一文件的行批量发起数据库请求时,会占满连接池,其他文件的请求只能排队等待,表现为同步执行。
  2. 同步前置操作干扰:filterData是同步函数,会快速完成同一文件所有行的过滤,紧接着批量发起数据库请求,进一步挤占连接池资源。
  3. 冗余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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 00:54:13