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

