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

AWS Lambda处理S3大文件同桶跨文件夹复制失败求助

AWS Lambda处理S3大文件复制问题排查与修复

核心问题分析

  • Async Handler与回调不兼容:你的handler声明为async,但内部使用了回调式的s3.createMultipartUpload。Lambda会在async函数执行到末尾时立即终止,不会等待回调中的异步操作完成,导致后续分片上传、日志输出都无法执行。
  • 错误读取S3对象:用fs.createReadStream(finalDestinationPath)尝试读取S3中的文件,但该路径是S3对象的键,并非Lambda本地文件系统路径,fs无法读取,流无数据触发后续逻辑。
  • 分片大小配置错误:maxchunksize设为文件总大小,导致curchunk.length > maxchunksize永远不成立,仅最后一次上传整个文件,大文件场景下易触发Lambda临时存储溢出。
  • 未定义prepareData函数:代码调用了prepareData(main_folder)但未实现,会直接抛出错误,但因async与回调的冲突,错误未被正确捕获输出。
  • CopySource格式错误:S3复制操作的CopySource需格式为bucket/key,你直接使用srcKey会导致复制失败。

修复后的代码实现

以下是修正后的完整代码,改用Promise化AWS SDK方法配合async/await,直接从S3读取流并分片上传,解决所有上述问题:

// dependencies
const AWS = require('aws-sdk');
const util = require('util');

// 初始化S3客户端,启用Promise支持
const s3 = new AWS.S3();
// 将回调式方法转为Promise,适配async/await
const createMultipartUpload = util.promisify(s3.createMultipartUpload.bind(s3));
const uploadPart = util.promisify(s3.uploadPart.bind(s3));
const completeMultipartUpload = util.promisify(s3.completeMultipartUpload.bind(s3));
const abortMultipartUpload = util.promisify(s3.abortMultipartUpload.bind(s3));
const getObject = util.promisify(s3.getObject.bind(s3));

// 示例prepareData函数(请替换为你的业务逻辑)
const prepareData = (mainFolder) => {
  return {
    folder: 'target-folder',
    subfolder: 'sub-target',
    renamedFile: 'new-filename.csv'
  };
};

// S3分片最小5MB,固定分片大小为5MB
const CHUNK_SIZE = 5 * 1024 * 1024;

exports.handler = async (event) => {
  try {
    console.log("Reading options from event:\n", util.inspect(event, {depth: 5}));
    
    const srcBucket = event.Records[0].s3.bucket.name;
    const srcKey = decodeURIComponent(event.Records[0].s3.object.key.replace(/\+/g, " "));
    const fileSize = event.Records[0].s3.object.size;
    const dstBucket = srcBucket; // 同一存储桶,可按需修改
    
    console.log("SRC KEY: ", srcKey, ", File Size: ", ((fileSize / 1024) / 1024), " MB");

    // 修正文件类型匹配正则,原正则无法正确识别后缀
    const typeMatch = srcKey.match(/\.([^.]+)$/);
    if (!typeMatch) {
      console.log("Could not determine the file type.");
      return { statusCode: 400, body: "Could not determine the file type." };
    }

    const fileType = typeMatch[1].toLowerCase();
    if (fileType !== "csv") {
      console.log(`Unsupported file type: ${fileType}`);
      return { statusCode: 400, body: `Unsupported file type: ${fileType}` };
    }
    
    // 解析源路径,添加边界检查避免数组越界
    const URI_PARTS = srcKey.split('/');
    const TOTAL_PARTS = URI_PARTS.length;
    if (TOTAL_PARTS < 8) {
      console.log("Source key path format is invalid.");
      return { statusCode: 400, body: "Source key path format is invalid." };
    }
    
    const pre_file_folder = URI_PARTS[TOTAL_PARTS - 2];
    const hour = URI_PARTS[TOTAL_PARTS - 3];
    const day = URI_PARTS[TOTAL_PARTS - 4];
    const month = URI_PARTS[TOTAL_PARTS - 5];
    const year = URI_PARTS[TOTAL_PARTS - 6];
    const sub_folder = URI_PARTS[TOTAL_PARTS - 7];
    const main_folder = URI_PARTS[TOTAL_PARTS - 8];
    
    console.log("PATHS: ", URI_PARTS);
    
    const dst = prepareData(main_folder);
    const finalDestinationPath = `${dst.folder}/${dst.subfolder ? `${dst.subfolder}/` : ''}${dst.renamedFile}`;
    
    console.log("####1. INITIALIZE MULTIPART UPLOAD: ", finalDestinationPath);
    
    // 1. 创建分片上传任务
    const multipartUpload = await createMultipartUpload({
      Bucket: dstBucket,
      Key: finalDestinationPath,
      ContentType: 'text/csv'
    });
    console.log("##### 2. MULTIPART UPLOAD INITIALIZED: ", multipartUpload.UploadId);
    
    const uploadId = multipartUpload.UploadId;
    const parts = [];
    let partNumber = 1;
    let startByte = 0;

    try {
      // 2. 分块读取源文件并上传
      while (startByte < fileSize) {
        const endByte = Math.min(startByte + CHUNK_SIZE - 1, fileSize - 1);
        console.log(`##### Upload part ${partNumber}: bytes ${startByte}-${endByte}`);
        
        // 读取S3对象的指定分片
        const { Body } = await getObject({
          Bucket: srcBucket,
          Key: srcKey,
          Range: `bytes=${startByte}-${endByte}`
        });
        
        // 上传分片
        const uploadResult = await uploadPart({
          Body,
          Bucket: dstBucket,
          Key: finalDestinationPath,
          UploadId: uploadId,
          PartNumber: partNumber
        });
        
        parts.push({
          PartNumber: partNumber,
          ETag: uploadResult.ETag
        });
        
        partNumber++;
        startByte = endByte + 1;
      }

      // 3. 完成分片上传
      const completeResult = await completeMultipartUpload({
        Bucket: dstBucket,
        Key: finalDestinationPath,
        UploadId: uploadId,
        MultipartUpload: { Parts: parts }
      });
      
      console.log('###### 5. UPLOAD COMPLETED: ', JSON.stringify(completeResult));
      return { statusCode: 200, body: "File copied successfully." };
    } catch (uploadErr) {
      // 上传失败时终止分片任务,避免残留资源产生费用
      console.log('## ERROR DURING UPLOAD: ', uploadErr);
      await abortMultipartUpload({
        Bucket: dstBucket,
        Key: finalDestinationPath,
        UploadId: uploadId
      });
      throw uploadErr;
    }
  } catch (err) {
    console.log("RESULT ERROR: ", err);
    return { statusCode: 500, body: err.message };
  }
};

关键修复点说明

  1. Promise化AWS SDK方法:用util.promisify将回调式方法转为Promise,配合async/await确保Lambda等待所有异步操作完成后再终止,解决日志无输出问题。
  2. 正确读取S3对象:改用s3.getObject并指定Range头部分片读取源文件,避免使用fs读取本地不存在的文件。
  3. 合理设置分片大小:固定分片为5MB(符合S3最小分片要求),确保大文件正确分片上传,同时避免Lambda临时存储溢出。
  4. 添加错误处理与资源清理:上传失败时调用abortMultipartUpload清理未完成的分片任务,避免产生额外费用。
  5. 修正文件类型匹配正则:原正则无法正确匹配文件名后缀,改为/\.([^.]+)$/。
  6. 路径边界检查:添加对URI_PARTS长度的检查,避免数组越界错误。

内容的提问来源于stack exchange,提问作者Antonio Longo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 17:45:40