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')我的所有数据完成处理之前就被命中了。有谁知道我可能做错了什么?
提前致谢!
.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)
我最近遇到了同样的问题。我通过使用一系列承诺来修复它,并等待所有承诺在.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)
| 归档时间: |
|
| 查看次数: |
4563 次 |
| 最近记录: |