Node.js+AWS Lambda导出MySQL海量数据至XLSX超时问题求助
MySQL大数据量导出XLSX到S3的Lambda超时问题解决
我需要完成从MySQL提取数据并生成XLSX文件存储至S3的任务,但目标MySQL表含140+列,单月平均数据量超15万条,使用Node.js结合AWS Lambda执行时,即便拉满Lambda超时上限仍触发超时异常。原SQL查询语句与Node.js代码如下:
原SQL查询语句
SELECT filterTable.jobSequenceNumber, sntl_invoice_details_import_Ref.invoiceSlNo AS invoiceSequenceNumber, sntl_lineItem_import_Ref.itemSlNo AS itemSequenceNumber, sntl_lineItem_import_Ref.itemSlNo AS "Item Sno", sntl_lineItem_import_Ref.partCode AS "Part Code", sntl_lineItem_import_Ref.description AS "Description", sntl_lineItem_import_Ref.concatDescription AS "Concat Description", sntl_lineItem_import_Ref.hsn AS "HSN", sntl_lineItem_import_Ref.qty AS "Qty", sntl_lineItem_import_Ref.uom AS "UOM", sntl_lineItem_import_Ref.cusQty AS "Cus Qty", sntl_lineItem_import_Ref.cusUom AS "Cus UOM", sntl_lineItem_import_Ref.unitPrice AS "Unit Price", sntl_lineItem_import_Ref.amount AS "Amount", sntl_lineItem_import_Ref.specificQty1 AS "Specific Qty 1", sntl_lineItem_import_Ref.uom1 AS "UOM1", sntl_lineItem_import_Ref.uom2 AS "UOM2", sntl_lineItem_import_Ref.generalDescription AS "Generic Description", sntl_lineItem_import_Ref.specificQty2 AS "Specific Qty 2", sntl_lineItem_import_Ref.manufacturerName AS "Manufacturer Name", sntl_lineItem_import_Ref.manufacturerType AS "Manufacturer Type", sntl_lineItem_import_Ref.manufacturerCode AS "Manufacturer Code", sntl_lineItem_import_Ref.manufacturerAddress1 AS "Manufacturer Address1", sntl_lineItem_import_Ref.manufacturerAddress2 AS "Manufacturer Address2", sntl_lineItem_import_Ref.manufacturerCity AS "Manufacturer City", sntl_lineItem_import_Ref.manufacturerSubDivision AS "Manufacturer Sub Division", sntl_lineItem_import_Ref.manufacturerPin AS "Manufacturer Pin", sntl_lineItem_import_Ref.manufacturerCountry AS "Manufacturer Country", sntl_lineItem_import_Ref.model AS "Model", sntl_lineItem_import_Ref.endUse AS "End Use", sntl_lineItem_import_Ref.brand AS "Brand", sntl_lineItem_import_Ref.countryOfOrigin AS "Country Of Origin", sntl_lineItem_import_Ref.sourceCountry AS "Source Country", sntl_lineItem_import_Ref.transitCountry AS "Transit Country", sntl_lineItem_import_Ref.bcdNtfnNo AS "BCD Notification No.", sntl_lineItem_import_Ref.bcdNtfnSlNo AS "BCD Notification Sr.No.", sntl_lineItem_import_Ref.customsAdditionalDutyNtfnNo AS "Customs Additional Duty Notification No.", sntl_lineItem_import_Ref.customsAdditionalDutyNtfnSlNo AS "Customs Additional Duty Notification Sr.No.", sntl_lineItem_import_Ref.chCessNtfnNo AS "Health Cess Notification No.", sntl_lineItem_import_Ref.chCessNtfnSlNo AS "Health Cess Notification Sr.No.", sntl_lineItem_import_Ref.caidcNtfnNo AS "CAIDC Notification No.", sntl_lineItem_import_Ref.caidcNtfnSlNo AS "CAIDC Notification Sr.No.", sntl_lineItem_import_Ref.swsNtfnNo AS "SWS Notification No.", sntl_lineItem_import_Ref.swsNtfnSlNo AS "SWS Notification Sr.No.", sntl_lineItem_import_Ref.cusEduCessNtfnNo AS "Cus.Edu. Cess Notification No.", sntl_lineItem_import_Ref.cusEduCessNtfnSlNo AS "Cus.Edu. Cess Notification Sr.No.", sntl_lineItem_import_Ref.nccdNtfnNo AS "NCCD Notification No.", sntl_lineItem_import_Ref.nccdNtfnSlNo AS "NCCD Notification Sr.No.", sntl_lineItem_import_Ref.saptaNtfnNo AS "SAPTA Notification No.", sntl_lineItem_import_Ref.saptaNtfnSlNo AS "SAPTA Notification Sr.No." FROM ((SELECT filterJobRefNo, row_number() over() as jobSequenceNumber FROM (SELECT DISTINCT sntl_job_details_import_filter.id as filterJobRefNo FROM (SELECT * FROM sntl_job_details_import WHERE isActiveJob="1" AND tenant_id="8e95996f-4de1-45b7-9e5d-d327483f239a" AND DATE(`jobCreationDate`) BETWEEN '2023-03-20' AND '2023-03-21' order by create_Date_Time desc) sntl_job_details_import_filter ) table1) filterTable INNER JOIN (SELECT * FROM sntl_invoice_details_import WHERE isActiveInvoice=1) sntl_invoice_details_import_Ref ON sntl_invoice_details_import_Ref.jobRefNo=filterTable.filterJobRefNo INNER JOIN (SELECT * FROM sntl_lineItem_import WHERE activeLineItem=1) sntl_lineItem_import_Ref ON sntl_lineItem_import_Ref.jobRefNo=filterTable.filterJobRefNo AND sntl_lineItem_import_Ref.invoiceRefNo=sntl_invoice_details_import_Ref.invoiceCreationDate) ORDER BY jobSequenceNumber, invoiceSequenceNumber, itemSequenceNumber;
原Node.js代码
const fs = require("fs"); const xlsx = require("xlsx"); module.exports.createReportFileS3 = async (report,sqlRes) => { try { let responseOutput = [] if(report.template?.sheetType==="multiple"){ sqlRes.forEach(item => { if( Array.isArray(item)){ responseOutput = [...responseOutput, item] } }) console.log(responseOutput) fs.writeFileSync("responseOutput.json",JSON.stringify(responseOutput)); const newWB = xlsx.utils.book_new(); for(const [index, column] of responseOutput.entries()){ const newWS = xlsx.utils.json_to_sheet(column); xlsx.utils.book_append_sheet(newWB, newWS, `${report.tableLabelName[index]}`); } xlsx.writeFile(newWB, `/tmp/${report.template.id}.xlsx`); //construct xlsx file with fetched data in local }else{ console.log('report..',report); let responseOutput = [] sqlRes.forEach(item => { if( Array.isArray(item)){ responseOutput = [...item] } }) const newWB = xlsx.utils.book_new(); const newWS = xlsx.utils.json_to_sheet(responseOutput); xlsx.utils.book_append_sheet(newWB, newWS, "USERS_LIST"); xlsx.writeFile(newWB, `/tmp/${report.template.id}.xlsx`); //construct xlsx file with fetched data in local } const reportCreationDate = report.reportCreationDate//Date.now() const readStream = fs.createReadStream(`/tmp/${report.template.id}.xlsx`) // a ReadStream let params = { Bucket: config.S3_BUCKET.REPORTS_BUCKET_NAME, //S3 bucket name Key: `sntl_reports/${report.template.id}/${report.template.templateName}${reportCreationDate}.xlsx`, //file storage path and name Body: readStream, ACL:"public-read", //acess control list ContentType:'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet' }; console.log(params) const command = new PutObjectCommand(params) const data = await client.send(command); console.log("data",data); fs.unlink(`/tmp/${report.template.id}.xlsx`, (err) => { //delete xlsx file in local if (err){ throw err; console.log("error delete file"); } }); return data } catch (error) { console.error("error in upload xlsx file in s3 bucket", error, report, sqlRes); } }
问题分析
超时核心原因是一次性加载全量数据到内存、同步IO阻塞事件循环、SQL查询效率低下,结合Lambda的资源限制,导致15万条+140列的数据处理无法在超时时间内完成。
优化方案
1. SQL查询优化
- 简化子查询结构:去除多层嵌套子查询,直接关联表并生成行号,减少查询执行时间:
SELECT ROW_NUMBER() OVER(ORDER BY job.create_Date_Time DESC) AS jobSequenceNumber, inv.invoiceSlNo AS invoiceSequenceNumber, line.itemSlNo AS itemSequenceNumber, line.itemSlNo AS "Item Sno", -- 其他字段保持原查询不变 line.saptaNtfnSlNo AS "SAPTA Notification Sr.No." FROM sntl_job_details_import job INNER JOIN sntl_invoice_details_import inv ON job.id = inv.jobRefNo AND inv.isActiveInvoice = 1 INNER JOIN sntl_lineItem_import line ON job.id = line.jobRefNo AND inv.invoiceCreationDate = line.invoiceRefNo AND line.activeLineItem = 1 WHERE job.isActiveJob = "1" AND job.tenant_id = "8e95996f-4de1-45b7-9e5d-d327483f239a" AND DATE(job.jobCreationDate) BETWEEN '2023-03-20' AND '2023-03-21' ORDER BY jobSequenceNumber, invoiceSequenceNumber, itemSequenceNumber; - 添加联合索引:针对过滤、关联字段创建索引,大幅提升查询速度:
sntl_job_details_import:(isActiveJob, tenant_id, jobCreationDate, create_Date_Time, id)sntl_invoice_details_import:(isActiveInvoice, jobRefNo, invoiceCreationDate, invoiceSlNo)sntl_lineItem_import:(activeLineItem, jobRefNo, invoiceRefNo, itemSlNo)
- 分批查询:用游标或
LIMIT+OFFSET分批拉取数据(例如每次1000条),避免一次性加载全量数据到内存。
2. Node.js代码优化
- 流式生成XLSX:替换
xlsx库为exceljs(支持流式写入),边读取数据库数据边写入XLSX,降低内存占用:const ExcelJS = require('exceljs'); const { PassThrough } = require('stream'); const { PutObjectCommand } = require('@aws-sdk/client-s3'); module.exports.createReportFileS3 = async (report, fetchDataBatch) => { try { const passThroughStream = new PassThrough(); const workbook = new ExcelJS.stream.xlsx.WorkbookWriter({ stream: passThroughStream, useStyles: false, useSharedStrings: false }); // 创建工作表 const worksheet = workbook.addWorksheet( report.template?.sheetType === "multiple" ? report.tableLabelName[0] : "USERS_LIST" ); // 配置表头(对应SQL返回的字段) worksheet.columns = [ { header: 'jobSequenceNumber', key: 'jobSequenceNumber' }, { header: 'invoiceSequenceNumber', key: 'invoiceSequenceNumber' }, // 其余字段依次添加 ]; // 分批写入数据 let batch; while ((batch = await fetchDataBatch()) !== null) { worksheet.addRows(batch); await worksheet.commit(); } await workbook.commit(); passThroughStream.end(); // 直接流式上传到S3,无需本地文件 const params = { Bucket: config.S3_BUCKET.REPORTS_BUCKET_NAME, Key: `sntl_reports/${report.template.id}/${report.template.templateName}${report.reportCreationDate}.xlsx`, Body: passThroughStream, ACL: "public-read", ContentType: 'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet' }; const command = new PutObjectCommand(params); return await client.send(command); } catch (error) { console.error("上传S3失败", error); } }; - 替换同步IO为异步操作:将
fs.writeFileSync、fs.unlink改为fs.promises方法,避免阻塞事件循环。 - 移除冗余操作:删除原代码中
responseOutput.json的写入步骤,直接处理数据即可。
3. Lambda配置优化
- 提升内存配置:Lambda内存与CPU、网络带宽正相关,建议设置为2048MB或更高,加快数据处理和传输速度。
- 申请延长超时:若分批处理后仍有压力,可联系AWS申请将Lambda超时延长至最高15分钟。
- 使用Lambda层:将
exceljs、@aws-sdk/client-s3等依赖打包为Lambda层,减少部署包大小,加快冷启动速度。
4. 架构层面优化
如果数据量持续增长,可考虑:
- 用AWS Glue替代Lambda:Glue专为ETL任务设计,支持更大数据量处理,无超时限制。
- 先导出CSV到S3再转换:将MySQL数据直接导出为CSV到S3,再用Athena或Glue异步转换为XLSX,效率更高。
内容的提问来源于stack exchange,提问作者Satyam Kumar
相关产品推荐
相关产品推荐

