【问题标题】:nodejs async await inside createReadStreamnodejs async await 在 createReadStream 中
【发布时间】:2020-02-02 03:48:06
【问题描述】:

我正在逐行读取 CSV 文件并在 MongoDB 中插入/更新。预期的输出将是 1. 控制台日志(行); 2. 控制台日志(光标); 3.console.log("流");

但是得到像这样的输出 1. 控制台日志(行); 控制台.log(行);控制台.log(行);控制台.log(行);控制台.log(行); …………………… 2. 控制台日志(光标); 3.console.log("流"); 请让我知道我在这里缺少什么。

const csv = require('csv-parser');
const fs = require('fs');

var mongodb = require("mongodb");

var client = mongodb.MongoClient;
var url = "mongodb://localhost:27017/";
var collection;
client.connect(url,{ useUnifiedTopology: true }, function (err, client) {

  var db = client.db("UKCompanies");
  collection = db.collection("company");
  startRead();
});
var cursor={};

async function insertRec(row){
  console.log(row);
  cursor = await collection.update({CompanyNumber:23}, row, {upsert: true});
  if(cursor){
    console.log(cursor);
  }else{
    console.log('not exist')
  }
  console.log("stream");
}



async function startRead() {
  fs.createReadStream('./data/inside/6.csv')
    .pipe(csv())
    .on('data', async (row) => {
      await insertRec(row);
    })
    .on('end', () => {
      console.log('CSV file successfully processed');
    });
}

【问题讨论】:

    标签: node.js async-await nodejs-stream


    【解决方案1】:

    在您的 startRead() 函数中,await insertRec() 不会在 insertRec() 处理时阻止更多 data 事件的流动。因此,如果您不希望在 insertRec() 完成之前运行下一个 data 事件,则需要暂停,然后恢复流。

    async function startRead() {
      const stream = fs.createReadStream('./data/inside/6.csv')
        .pipe(csv())
        .on('data', async (row) => {
          try {
            stream.pause();
            await insertRec(row);
          } finally {
            stream.resume();
          }
        })
        .on('end', () => {
          console.log('CSV file successfully processed');
        });
    }
    

    仅供参考,如果insertRec() 失败,您还需要一些错误处理。

    【讨论】:

    • 谢谢@jfriend00
    【解决方案2】:

    在这种情况下这是预期的行为,因为您的 on 数据侦听器会在数据流中可用时异步触发 insertRec。这就是为什么您的第一行插入方法正在并行执行的原因。如果您想控制此行为,您可以在创建读取流时使用highWaterMark (https://nodejs.org/api/stream.html#stream_readable_readablehighwatermark) 属性。这样您一次将获得 1 条记录,但我不确定您的用例是什么。

    类似的东西

    fs.createReadStream(`somefile.csv`, {
      "highWaterMark": 1
    })
    

    您也没有等待您的startRead 方法。我会将它包装在 Promise 中并在 end 侦听器中解决它,否则您将不知道处理何时完成。类似的东西

    function startRead() {
      return new Promise((resolve, reject) => {
        fs.createReadStream(`somepath`)
          .pipe(csv())
          .on("data", async row => {
            await insertRec(row);
          })
          .on("error", err => {
            reject(err);
          })
          .on("end", () => {
            console.log("CSV file successfully processed");
            resolve();
          });
      });
    
    }
    

    【讨论】:

    • 设置 highWaterMark 不会让您限制 data 事件的速率。相反,OP 应该实现一个流 Writable,可以将其配置为 write 逐个文档或 writev 大量文档。 highWaterMark 可让您控制内存压力。
    • @jorgenkg 这是真的。感谢您的澄清。
    • @jorgenkg - “对于在对象模式下运行的流,highWaterMark 指定对象总数” - nodejs.org/api/stream.html#stream_buffering
    • 是 - 将在(读/写)流的内部缓冲区中缓冲的对象数。对象将始终使用write 一次处理一个。highWaterMark 指示可以为流实例缓冲多少对象。
    猜你喜欢
    • 2022-11-23
    • 2021-06-15
    • 1970-01-01
    • 2018-09-10
    • 2018-06-13
    • 1970-01-01
    • 2018-10-13
    • 2016-10-15
    相关资源
    最近更新 更多