【问题标题】:How to stream read an S3 JSON file to postgreSQL using async/await in a NodeJS 12 Lambda function?如何在 NodeJS 12 Lambda 函数中使用 async/await 将 S3 JSON 文件流式读取到 postgreSQL?
【发布时间】:2020-11-28 05:08:39
【问题描述】:

我没有意识到这么简单的任务会有多么危险。 我们正在尝试流式读取存储在 S3 中的 JSON 文件——我认为我们已经完成了这部分工作。 我们的 .on('data') 回调被调用,但 Node 选择并选择它想要运行的位 - 似乎是随机的。

我们设置了一个流阅读器。

stream.on('data', async x => { 
  await saveToDb(x);  // This doesn't await.  It processes saveToDb up until it awaits.
});

有时 db 调用会到达 db ——但大多数时候它不会。 我得出的结论是 EventEmitter 在异步/等待事件处理程序方面存在问题。 只要您的代码是同步的,它似乎就会与您的异步方法一起播放。但是,在您等待时,它会随机决定是否实际执行此操作。

它流式传输各种块,我们可以console.log 将它们输出并查看数据。但是,一旦我们尝试触发 await/async 调用,我们就会停止看到可靠消息。

我在 AWS Lambda 中运行它,有人告诉我有一些特殊注意事项,因为它们显然在某些情况下会停止处理?

我尝试在 IFFY 中包围 await 调用,但这也不起作用。

我错过了什么? 有没有办法告诉 JavaScript——“好的,我需要你同步运行这个异步任务。我的意思是——也不要再触发任何事件通知。就坐在这里等着。”?

【问题讨论】:

    标签: node.js amazon-s3 aws-lambda async-await jsonstream


    【解决方案1】:

    TL;DR:

    • 使用异步迭代器从流管道的末尾拉取!
    • 不要在任何流代码中使用异步函数!

    详情:

    关于async/await 和流的生命之谜的秘密似乎包含在Async Iterators 中!

    简而言之,我将一些流通过管道连接在一起,最后,我创建了一个异步迭代器来将内容拉出,以便我可以异步调用数据库。 ChunkStream 为我做的唯一一件事就是排队多达 1,000 个来调用数据库,而不是每个项目。我是队列新手,所以可能已经有更好的方法了。

    // ...
    const AWS = require('aws-sdk');
    const s3 = new AWS.S3();
    const JSONbigint = require('json-bigint');
    JSON.parse = JSONbigint.parse; // Let there be proper bigint handling!
    JSON.stringify = JSONbigint.stringify;
    const stream = require('stream');
    const JSONStream = require('JSONStream');
    
    exports.handler = async (event, context) => {
        // ...
        let bucket, key;
        try {
            bucket = event.Records[0].s3.bucket.name;
            key = event.Records[0].s3.object.key;
            console.log(`Fetching S3 file: Bucket: ${bucket}, Key: ${key}`);
            const parser = JSONStream.parse('*'); // Converts file to JSON objects
            let chunkStream = new ChunkStream(1000); // Give the db a chunk of work instead of one item at a time
            let endStream = s3.getObject({ Bucket: bucket, Key: key }).createReadStream().pipe(parser).pipe(chunkStream);
            
            let totalProcessed = 0;
            async function processChunk(chunk) {
                let chunkString = JSON.stringify(chunk);
                console.log(`Upserting ${chunk.length} items (starting with index ${totalProcessed}) items to the db.`);
                await updateDb(chunkString, pool, 1000); // updateDb and pool are part of missing code
                totalProcessed += chunk.length;
            }
            
            // Async iterator
            for await (const batch of endStream) {
                // console.log(`Processing batch (${batch.length})`, batch);
                await processChunk(batch);
            }
        } catch (ex) {
            context.fail("stream S3 file failed");
            throw ex;
        }
    };
    
    class ChunkStream extends stream.Transform {
        constructor(maxItems, options = {}) {
            options.objectMode = true;
            super(options);
            this.maxItems = maxItems;
            this.batch = [];
        }
        _transform(item, enc, cb) {
            this.batch.push(item);
            if (this.batch.length >= this.maxItems) {
                // console.log(`ChunkStream: Chunk ready (${this.batch.length} items)`);
                this.push(this.batch);
                // console.log('_transform - Restarting the batch');
                this.batch = [];
            }
            cb();
        }
        _flush(cb) {
            // console.log(`ChunkStream: Flushing stream (${this.batch.length} items)`);
            if (this.batch.length > 0) {
                this.push(this.batch);
                this.batch = [];
            }
            cb();
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2022-01-27
      • 2019-05-30
      • 1970-01-01
      • 2020-04-10
      • 1970-01-01
      • 2019-03-30
      • 1970-01-01
      • 2018-06-13
      • 2019-02-17
      相关资源
      最近更新 更多