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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 14:47:10