【问题标题】:Node.js close the writestream at the end of the readstreamNode.js 在读流结束时关闭写流
【发布时间】:2020-06-28 14:40:06
【问题描述】:

我正在阅读存储在 AWS S3 上的 csv 文件。这些 csv 记录由一个名为 filterLogic() 的函数评估。那些未通过测试的记录也需要写入 AWS S3 上的错误报告 csv 文件。对于 csv 解析,我使用的是 fast-csv。

但是,对于最后一条记录,我得到了 "Error [ERR_STREAM_WRITE_AFTER_END]: write after end"

我需要在哪里以及如何正确调用csvwritestream.end()

const AWS = require('aws-sdk');
const utils = require('./utils');
const csv = require('fast-csv');
const s3 = new AWS.S3();

exports.handler = async (event) => {
    console.log("Incoming Event: ", JSON.stringify(event));
    const bucket = event.Records[0].s3.bucket.name;
    const filename = decodeURIComponent(event.Records[0].s3.object.key.replace(/\+/g, ' '));
    const message = `File is uploaded in - ${bucket} -> ${filename}`;
    console.log(message);

    const splittedFilename = filename.split('.');
    const reportFilename = splittedFilename[0] + "Report." + splittedFilename[1];
    const reportBucket = 'external.transactions.reports';

    const csvwritestream = csv.format({ headers: true });
    csvwritestream
        .pipe(utils.uploadFromStream(s3, reportBucket, reportFilename))
        .on('end', function () {
            console.log("Report written to S3 " + reportFilename);
        });

    var request = s3.getObject({ Bucket: bucket, Key: filename });
    var stream = request.createReadStream({ objectMode: true })
        .pipe(csv.parse({ headers: true }))
        .on('data', async function (data) {
            stream.pause();
            console.log("JSON: " + JSON.stringify(data));
            var response = await utils.filterLogic(data);
            if (response.statusCode !== 200) {
                await csvwritestream.write(data);
                console.log("Data: " + JSON.stringify(data) + " written."); 
            }
            stream.resume();
        })
        .on('end', function(){
            csvwritestream.end();   
        });
    return new Promise(resolve => {
        stream.on('close', async function () {
            csvwritestream.end();
            resolve();
        });
    });
};

【问题讨论】:

    标签: node.js amazon-s3 async-await stream es6-promise


    【解决方案1】:

    您的错误原因并不简单。主要问题是node.js event emitter 中没有等待侦听器data - 因此您可能期望数据处理按顺序开始,但完成的时间不是。现在考虑到这一点,看到end 事件在最后一个data 事件之后立即触发,因此在写入流关闭后完成一些最后处理的项目的写入......是的,我知道你确实暂停了,但是在async function 中,您在节点处理了同步内容并且end 是其中之一之后这样做了。

    您可以做的是实现一个 Transform 流,但是,鉴于您的用例可能比仅使用另一个模块复杂得多 - 就像我的 scramjet 将允许您在数据过滤器中运行异步代码。

    const AWS = require('aws-sdk');
    const utils = require('./utils');
    const {StringStream} = require('scramjet');
    const s3 = new AWS.S3();
    
    exports.handler = async (event) => {
        console.log("Incoming Event: ", JSON.stringify(event));
        const bucket = event.Records[0].s3.bucket.name;
        const filename = decodeURIComponent(event.Records[0].s3.object.key.replace(/\+/g, ' '));
        const message = `File is uploaded in - ${bucket} -> ${filename}`;
        console.log(message);
    
        const splittedFilename = filename.split('.');
        const reportFilename = splittedFilename[0] + "Report." + splittedFilename[1];
        const reportBucket = 'external.transactions.reports';
    
        var request = s3.getObject({ Bucket: bucket, Key: filename });
        var stream = StringStream
            // create a StringStream from a scramjet stream
            .from(
                request.createReadStream()
            )
            // then just parse the data
            .CSVParse({headers: true})
            // then filter using asynchronous function as if it was an array
            .filter(async data => {
                var response = await utils.filterLogic(data);
                return response.statusCode === 200;
            })
            // then stringify
            .CSVStringify({headers: true})
            // then upload
            .pipe(
                utils.uploadFromStream(s3, reportBucket, reportFilename)
            );
    
        return new Promise(
            (res, rej) => stream
                .on("finish", res)
                .on("error", rej)
        );
    };
    

    Scramjet 将负责从上述方法创建流管道,因此您无需暂停/恢复 - 这一切都已完成。

    您可能想阅读:

    哦,超燃冲压发动机只增加了 3 个部门,所以你的节点模块不会成为那个相对论笑话的例子。 ;)

    【讨论】:

      猜你喜欢
      • 2013-10-17
      • 2016-11-24
      • 1970-01-01
      • 2019-11-26
      • 2018-05-16
      • 1970-01-01
      • 1970-01-01
      • 2016-01-29
      • 1970-01-01
      相关资源
      最近更新 更多