如何在Lambda函数中按顺序使用async/await执行异步操作
问题根因
执行顺序随机错乱是Promise写法的典型错误:
- Promise构造函数内的逻辑是声明时立即同步执行的,和你什么时候写await没有关系。你在代码顶部同时定义
parserFcn和connectToDb两个Promise的时候,CSV解析和Mongo连接逻辑就已经同时启动了,根本不存在先后顺序 - 两个Promise都没写完整的状态流转:数据库连接的Promise里从来没调用过
resolve/reject,插入逻辑直接读闭包里的parsedData,很容易拿到解析未完成的空数组 - 最后那段
Promise.all完全是无效代码:你前面已经分别对两个Promise做了await,传给Promise.all的根本不是Promise实例,是两个await的返回值,起不到任何流程控制作用 - 漏了S3读取流的错误监听,流异常的时候会直接导致Promise挂起,触发Lambda超时
修复逻辑
严格按照依赖顺序执行:必须等S3流读取、CSV解析全流程完成,拿到完整的结构化数据之后,再启动数据库连接和写入逻辑,绝对不要提前初始化数据库连接的Promise。同时补全所有分支的resolve/reject,补全错误捕获,修正隐式全局变量问题。
修复后可直接运行的代码
const AWS = require("aws-sdk"); const csv = require("@fast-csv/parse"); const MongoClient = require("mongodb").MongoClient; const s3 = new AWS.S3(); exports.handler = async (event) => { const bucketName = event.Records[0].s3.bucket.name; const keyName = event.Records[0].s3.object.key; console.log("Bucket Name->", JSON.stringify(bucketName)); console.log("Bucket key->", JSON.stringify(keyName)); const params = { Bucket: bucketName, Key: keyName, }; const parsedData = []; const s3Contents = s3.getObject(params).createReadStream(); // 第一阶段:完成CSV解析,拿到全量数据 const parseResult = await new Promise((resolve, reject) => { // 补全S3流错误捕获 s3Contents.on('error', err => reject(`S3读取失败: ${err.message}`)); csv .parseStream(s3Contents, { headers: true }) .on("data", data => parsedData.push(data)) .on("end", rowCount => { console.log(`CSV解析完成,共${rowCount}行数据`); resolve(parsedData); }) .on("error", err => reject(`CSV解析失败: ${err.message}`)); }); // 第二阶段:解析完成后再连接数据库写入,此时parseResult一定是完整数据 const dbResult = await new Promise((resolve, reject) => { MongoClient.connect( "mongodb://user:pass@host/?ssl=true&retryWrites=false", { tlsCAFile: `/opt/rds-combined-ca-bundle.pem`, }, async (err, client) => { if (err) { return reject(`数据库连接失败: ${err.message}`); } console.log("DocumentDB连接成功,待写入数据长度:", parseResult.length); try { const db = client.db("database-name"); const col = db.collection("collection-name"); const insertRes = await col.insertMany(parseResult); console.log(`数据写入成功,共插入${insertRes.insertedCount}条`); resolve(insertRes); } catch (insertErr) { reject(`数据写入失败: ${insertErr.message}`); } finally { // 无论写入成功失败都关闭连接,避免连接泄漏 await client.close(); } } ); }); return { status: "success", parsedRowCount: parseResult.length, insertedCount: dbResult.insertedCount }; };
关键注意点
- 只要存在
new Promise()声明,构造函数内的逻辑就会立刻执行,不会等await触发。如果两个步骤有明确先后依赖,绝对不要提前声明后一个步骤的Promise - 所有Promise必须覆盖所有成功、失败分支的
resolve/reject调用,否则会出现await永久挂起、Lambda超时的问题 - async形式的Lambda handler不要混用callback参数,直接return返回结果、throw抛出错误即可,避免执行流混乱
- 流场景必须加error事件监听,否则流异常会直接抛出未捕获错误,导致函数异常退出
内容的提问来源于stack exchange,提问作者Fury
相关产品推荐
相关产品推荐

