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

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记录。


问题原因与解决方案

核心问题

  1. 回调模式覆盖流模式:创建csvParser时传入了回调函数,这会让csv-parse切换到回调模式——解析完成后所有记录会一次性传给回调,流本身不会再输出数据,后续for await...of自然拿不到内容。
  2. 遍历已结束的流:await pipeline执行完成时,csvParser流已经完全结束,此时再尝试遍历流,不会有任何数据输出。
  3. 返回流对象而非数据:修改后返回的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 13:48:25