Lambda无法读取S3触发事件中的CSV数据问题排查
问题描述
我用Node.js编写Lambda从S3存储桶加载数据,架构为S3(创建事件)>SNS>SQS>Lambda,触发流程完全正常,CloudWatch日志显示每次S3 PUT事件都能成功调用Lambda,且已配置Lambda读取SQS和S3的IAM角色,未对S3存储桶或对象使用SSE。但代码中两种读取S3对象的方式均无法在控制台输出CSV文件内容,CloudWatch日志仅输出到参数部分便结束,请问是否S3 API有更新导致该问题?
代码
const readline = require('readline'); const AWS = require("aws-sdk"); const s3 = new AWS.S3({apiVersion: '2006-03-01'}); exports.handler = async (event) => { console.log('###### Received event:', JSON.stringify(event, null, 2)); for (const record of event.Records) { try { const body = JSON.parse((record.body)); // console.log(`###### Record BODY ${JSON.stringify(body, null, 2)}`); const message = JSON.parse(body.Message); //console.log(`##### Record BODY MESSAGE ${JSON.stringify(message, null, 2)}`); const bucket = message.Records[0].s3.bucket.name; const key = message.Records[0].s3.object.key; console.log(`###### Record BODY bucket ${bucket}`); console.log(`###### Record BODY key ${key}`); const params = { Bucket: bucket, Key: key }; console.log(`###### THE PARAMS ${JSON.stringify(params, null, 2)}`); s3.getObject(params , function (err, data) { if (err) { console.log(`###### ERR ${err}`); throw err; } else { console.log(`###### ${data.Body.toString()}`); } }) const rl = readline.createInterface({ input: s3.getObject(params).createReadStream() }); rl.on('line', function(line) { console.log(`####### RECORD DATA LINE ${line}`); }) .on('close', function() { console.log(`####### RECORD DATA CLOSED!!!!!`); }); } catch (e ) { console.log(`###### Error ${JSON.stringify(e, null, 2)}`); } } return `Successfully processed ${event.Records.length} messages.`; };
IAM角色
{ "Version": "2012-10-17", "Statement": [ { "Sid": "VisualEditor0", "Effect": "Allow", "Action": [ .. "s3:GetObject", .. ], "Resource": "arn:aws:s3:::my-REDACTED-bucket" } ] }
CloudWatch日志
2022-08-31T12:00:07.897Z 0191e48d-f3ac-5af8-9334-bf04bcd815d8 INFO ###### Record BODY bucket my-REDACTED-bucket 2022-08-31T12:00:07.897Z 0191e48d-f3ac-5af8-9334-bf04bcd815d8 INFO ###### Record BODY key test.csv 2022-08-31T12:00:07.897Z 0191e48d-f3ac-5af8-9334-bf04bcd815d8 INFO THE PARAMS { "Bucket": "my-REDACTED-bucket", "Key": "test.csv" } END RequestId: 0191e48d-f3ac-5af8-9334-bf04bcd815d8
问题原因及解决办法
这不是S3 API更新导致的问题,核心原因是Lambda的async函数在异步操作完成前就返回,进程提前终止,异步的S3调用还没执行完就被Lambda强制结束了。
具体问题点
s3.getObject的回调版本是异步执行的,Lambda函数执行到return语句时,这个回调还没触发,所以不会输出内容。readline的流操作也是异步的,同样在Lambda返回时还没完成读取,导致line和close事件都没机会触发。
修复方案
将异步操作包装成Promise,用await等待所有异步任务完成后再让Lambda返回。
修改后的代码
const readline = require('readline'); const AWS = require("aws-sdk"); const s3 = new AWS.S3({apiVersion: '2006-03-01'}); // 把getObject包装为Promise const getS3Object = async (params) => { return new Promise((resolve, reject) => { s3.getObject(params, (err, data) => { if (err) reject(err); else resolve(data); }); }); }; // 把readline流读取包装为Promise const readS3Stream = async (params) => { return new Promise((resolve) => { const rl = readline.createInterface({ input: s3.getObject(params).createReadStream() }); rl.on('line', (line) => { console.log(`####### RECORD DATA LINE ${line}`); }) .on('close', () => { console.log(`####### RECORD DATA CLOSED!!!!!`); resolve(); }); }); }; exports.handler = async (event) => { console.log('###### Received event:', JSON.stringify(event, null, 2)); for (const record of event.Records) { try { const body = JSON.parse((record.body)); const message = JSON.parse(body.Message); const bucket = message.Records[0].s3.bucket.name; const key = message.Records[0].s3.object.key; console.log(`###### Record BODY bucket ${bucket}`); console.log(`###### Record BODY key ${key}`); const params = { Bucket: bucket, Key: key }; console.log(`###### THE PARAMS ${JSON.stringify(params, null, 2)}`); // 等待getObject完成并输出内容 const data = await getS3Object(params); console.log(`###### ${data.Body.toString()}`); // 等待流读取完成 await readS3Stream(params); } catch (e ) { console.log(`###### Error ${JSON.stringify(e, null, 2)}`); } } return `Successfully processed ${event.Records.length} messages.`; };
额外权限修正
当前IAM角色的Resource配置有误,s3:GetObject需要授权到存储桶内的对象级别,而非仅存储桶本身。修改后的IAM角色:
{ "Version": "2012-10-17", "Statement": [ { "Sid": "VisualEditor0", "Effect": "Allow", "Action": [ "s3:GetObject" ], "Resource": "arn:aws:s3:::my-REDACTED-bucket/*" } ] }
内容的提问来源于stack exchange,提问作者Frederick Scott Smith
相关产品推荐
相关产品推荐

