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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 13:18:31