NestJS结合fast-csv处理大CSV时的内存溢出问题排查
大CSV文件处理内存溢出排查与解决方案(NestJS v9 + fast-csv v4 + BigQuery)
核心排查方向
1. FileInterceptor的文件存储方式
默认FileInterceptor会将上传文件完全加载到内存中,25万行的CSV文件体积远超512MB服务器内存上限,必然导致溢出。需确认是否未配置磁盘存储,仍用内存存储接收文件。
2. fast-csv流式处理的背压控制缺失
若未在批量插入时暂停CSV流,BigQuery插入速度慢于CSV解析速度时,会导致大量未处理的行数据堆积在内存中,最终溢出。
3. 未释放的内存引用
检查是否存在全局数组、闭包持有已处理完的行数据,或事件监听器未及时移除,导致GC无法回收无用内存。
4. BigQuery插入的阻塞逻辑
若批量插入采用同步等待模式,且插入耗时较长,会导致CSV流持续推送数据至内存,形成堆积。
针对性解决方案
1. 改用磁盘存储接收上传文件
修改FileInterceptor配置,将文件暂存到磁盘而非内存,避免大文件直接占用内存:
import { diskStorage } from 'multer'; import { FileInterceptor } from '@nestjs/platform-express'; import * as fs from 'fs'; @Post('upload') @UseInterceptors(FileInterceptor('file', { storage: diskStorage({ destination: './tmp', // 提前创建临时目录并配置权限 filename: (req, file, cb) => { const uniqueName = `${Date.now()}-${Math.random().toString(36).slice(2)}`; cb(null, `${uniqueName}.csv`); } }) })) async uploadFile(@UploadedFile() file: Express.Multer.File) { try { const fileStream = fs.createReadStream(file.path); await this.csvProcessingService.handleStream(fileStream); } finally { // 处理完成后删除临时文件 fs.unlinkSync(file.path); } }
2. 为fast-csv添加背压控制
在解析CSV时,每处理一批数据就暂停流,等待BigQuery插入完成后再恢复,避免内存堆积:
import * as csv from 'fast-csv'; import * as fs from 'fs'; async handleStream(stream: fs.ReadStream) { let batch: Record<string, any>[] = []; const parser = csv.parse({ headers: true }) .on('data', async (row) => { parser.pause(); // 暂停流,阻止继续推送数据 const transformedRow = this.transformRow(row); // 自定义行数据转换逻辑 batch.push(transformedRow); if (batch.length >= 200) { await this.bigQueryService.insertBatch(batch); batch = []; // 清空批次,释放内存 } parser.resume(); // 恢复流,继续解析下一批 }) .on('end', async () => { // 处理剩余未批量的行 if (batch.length > 0) { await this.bigQueryService.insertBatch(batch); } }); stream.pipe(parser); // 等待解析与插入全部完成 await new Promise((resolve, reject) => { parser.on('finish', resolve); parser.on('error', reject); }); }
3. 优化BigQuery批量插入逻辑
确保BigQuery插入采用异步非阻塞模式,使用官方SDK的批量插入方法,避免单条插入或同步等待:
import { BigQuery } from '@google-cloud/bigquery'; async insertBatch(rows: Record<string, any>[]) { const bigquery = new BigQuery(); const dataset = bigquery.dataset('your-target-dataset'); const table = dataset.table('your-target-table'); // 使用流式批量插入,降低阻塞风险 await table.insert(rows, { ignoreUnknownValues: true, raw: false }); }
4. 内存监控与定位
启动应用时开启Node.js调试模式,通过Chrome DevTools分析内存快照:
node --inspect dist/main.js
打开Chrome浏览器输入chrome://inspect,连接到调试进程,录制内存快照,对比处理不同规模文件时的内存对象分布,定位未释放的内存来源。
内容的提问来源于stack exchange,提问作者Amos Isaila Lucian Onofrei
相关产品推荐
相关产品推荐

