如何使用 async/await 在 node.js 中异步 createReadStream

Joh*_*son 3 stream node.js async-await

我无法使用fs.creadReadStream异步处理我的 csv 文件:

async function processData(row) {
    // perform some asynchronous function
    await someAsynchronousFunction();
}

fs.createReadStream('/file')
    .pipe(parse({
        delimiter: ',',
        columns: true
    })).on('data', async (row) => {
        await processData(row);
    }).on('end', () => {
        console.log('done processing!')
    })
Run Code Online (Sandbox Code Playgroud)

我想在createReadStream到达之前逐条读取每条记录后执行一些异步功能on('end')

但是,在on('end')我的所有数据完成处理之前就被命中了。有谁知道我可能做错了什么?

提前致谢!

jfr*_*d00 6

.on('data, ...)不等你await。请记住,async函数会立即返回一个承诺,并且.on()不会关注该承诺,因此它只是继续愉快地进行。

await功能只里面等待,它不会立即返回停止你的功能,因此在流认为你处理数据,并保持发送更多的数据,并产生更多的data事件。

这里有几种可能的方法,但最简单的方法可能是暂停流直到processData()完成,然后重新启动流。

此外,是否processData()返回与异步操作完成相关的承诺?这也是await能够完成其工作所必需的。

可读流文档包含期间暂停流的示例data事件,然后一些异步操作结束后恢复它。这是他们的例子:

const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
  console.log(`Received ${chunk.length} bytes of data.`);
  readable.pause();
  console.log('There will be no additional data for 1 second.');
  setTimeout(() => {
    console.log('Now data will start flowing again.');
    readable.resume();
  }, 1000);
});
Run Code Online (Sandbox Code Playgroud)


TFi*_*her 6

我最近遇到了同样的问题。我通过使用一系列承诺来修复它,并等待所有承诺在.on("end")被触发时解决。

import parse from "csv-parse";

export const parseCsv = () =>
  new Promise((resolve, reject) => {
    const promises = [];
    fs.createReadStream('/file')
      .pipe(parse({ delimiter: ',', columns: true }))
      .on("data", row => promises.push(processData(row)))
      .on("error", reject)
      .on("end", async () => {
        await Promise.all(promises);
        resolve();
      });
  });
Run Code Online (Sandbox Code Playgroud)