Node.js 18.7中async/await结合Stream解析CSV无输出问题排查
问题:Node.js结合Stream与async/await解析CSV无数据输出
我在Ubuntu环境使用Node.js 18.7,需将大量CSV文件解析为对象(借助csv-parse库)并最终存入数据库。因文件数量较多,选择使用Stream并希望结合async/await语法。
原代码
const { parse } = require('csv-parse'); const path = __dirname + '/file1.csv'; const opt = { columns: true, relax_column_count: true, skip_empty_lines: true, skip_records_with_error: true }; console.log(path); const { pipeline } = require('stream'); // const pipeline = stream.pipeline; async function readByLine(path, opt) { const readFileStream = fs.createReadStream(path); var csvParser = parse(opt, function (err, records) { if (err) throw err; }); await pipeline(readFileStream, csvParser, (err) => { if (err) { console.error('Pipeline failed.', err); } else { console.log('Pipeline succeeded.'); } }); for await (const record of csvParser) { console.log(record); } } readByLine(path, opt)
原代码运行结果
Pipeline succeeded.
解析后的对象并未输出,于是我修改了代码:
修改后的代码
async function readByLine(path, opt) { const readFileStream = fs.createReadStream(path); var csvParser = parse(opt, function (err, records) { if (err) throw err; }); await pipeline(readFileStream, csvParser, (err) => { if (err) { console.error('Pipeline failed.', err); } else { console.log('Pipeline succeeded.'); } }); // for await (const record of csvParser) { // console.log(record); // } return csvParser; } (async function () { const o = await readByLine(path, opt); console.log(o); })();
修改后的运行结果
输出的是一个包含大量内部属性的流对象,而非解析后的CSV记录。
问题原因与解决方案
核心问题
- 回调模式覆盖流模式:创建
csvParser时传入了回调函数,这会让csv-parse切换到回调模式——解析完成后所有记录会一次性传给回调,流本身不会再输出数据,后续for await...of自然拿不到内容。 - 遍历已结束的流:
await pipeline执行完成时,csvParser流已经完全结束,此时再尝试遍历流,不会有任何数据输出。 - 返回流对象而非数据:修改后返回的
csvParser是流实例本身,不是解析后的记录,所以打印出来的是流的内部属性。
修正后的代码
要让csv-parse工作在流模式,需去掉回调函数,直接通过for await...of遍历解析流获取数据,同时正确使用Promise化的pipeline:
const { parse } = require('csv-parse'); const fs = require('fs'); const path = __dirname + '/file1.csv'; const opt = { columns: true, relax_column_count: true, skip_empty_lines: true, skip_records_with_error: true }; // 引入Promise化的pipeline,更适配async/await const { pipeline } = require('stream/promises'); async function readByLine(path, opt) { const readFileStream = fs.createReadStream(path); // 不要传入回调,让parse返回可读流 const csvParser = parse(opt); try { // 启动pipeline连接读文件流与解析流 await pipeline(readFileStream, csvParser); // 遍历解析流获取每条记录 for await (const record of csvParser) { console.log(record); // 在此处添加存入数据库的逻辑 } console.log('解析完成'); } catch (err) { console.error('解析失败:', err); } } readByLine(path, opt);
更简洁的写法(无需单独pipeline)
如果不需要pipeline的自动销毁错误处理,也可以直接用pipe连接流后遍历:
const { parse } = require('csv-parse'); const fs = require('fs'); const path = __dirname + '/file1.csv'; const opt = { columns: true, relax_column_count: true, skip_empty_lines: true, skip_records_with_error: true }; async function readByLine(path, opt) { const readFileStream = fs.createReadStream(path); const csvParser = parse(opt); // 连接读文件流与解析流 readFileStream.pipe(csvParser); try { // 遍历解析后的每条记录 for await (const record of csvParser) { console.log(record); // 执行数据库写入操作 } console.log('解析完成'); } catch (err) { console.error('解析出错:', err); // 出错时销毁读文件流 readFileStream.destroy(err); } } readByLine(path, opt);
内容的提问来源于stack exchange,提问作者user1592380
相关产品推荐
相关产品推荐

