You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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强制结束了。

具体问题点

  1. s3.getObject的回调版本是异步执行的,Lambda函数执行到return语句时,这个回调还没触发,所以不会输出内容。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 17:19:04