如何优化NodeJS中flatMap处理大数据时的内存占用
Node.js API 内存不足问题解决(分批处理优化)
问题场景
在Node.js API中处理约12000条数据时,原代码因**内存不足(out of memory error)**报错。原逻辑先一次性生成所有处理后的输入数据allInputs,再分批保存到数据库,导致内存峰值过高。
原代码核心问题
原代码通过flatMap+map一次性将所有days和sites的组合处理成完整的allInputs数组,所有处理后的对象都驻留在内存中,数据量较大时直接触发内存溢出。即使后续分批保存,生成阶段的内存占用已经超出限制。
原核心代码片段:
const allInputs = days .flatMap(day => { const isoDay = day.toISOString(); return sites.map(site => { const { studyId, siteId } = site; const data = { staleCandidates: getStaleCandidatesFromData({ data: analyticsData, site, day, noContact: studyGoals?.inactive?.rules?.noContact, }), ...replayData.find(x => x.day === isoDay && x.studyId === studyId && x.siteId === siteId), }; const calculator = new SiteScoreCalculator({ site, studyGoals, data, day }); return siteToInput({ days: [day], site, calculator }); }); }) .filter(removeNullScores);
优化方案:边生成边分批处理
核心思路是不一次性缓存所有处理结果,而是逐个生成输入项,积累到指定批次大小后立即保存并释放内存,同时优化数据查找效率:
- 预构建replayData查找Map:将
replayData转换成以${day}-${studyId}-${siteId}为key的Map,把每次O(n)的find操作优化为O(1),提升性能并减少内存开销。 - 逐一生成+分批积累:遍历days和sites的组合,生成单个输入项后加入当前批次,达到批次大小就执行保存,清空批次。
- 收尾剩余数据:遍历结束后,检查是否有未保存的剩余数据,执行最后一次保存。
修改后的完整代码
import logger from '@trialbee/logger'; import fetchStudyGoals from '../../projectGoals'; import { getAnalyticsSiteScoreReplayData, saveSiteScore } from '../../requests/analytics'; import { getStaleCandidatesFromData, removeNullScores, siteToInput } from '../../utils/site'; import { daysFromInterval } from '../../utils/times'; import SiteScoreCalculator from '../site/SiteScoreCalculator'; const log = logger(); const BATCH_SIZE = 100; // 根据服务器内存情况调整批次大小 const studiesToSites = studies => studies.flatMap(({ studyId, siteIds }) => siteIds.map(siteId => ({ studyId, siteId }))); async function replaySiteScores({ interval, studies, prune, data: analyticsData }) { const start = Date.now(); log.debug('Replay site scores start'); const days = daysFromInterval(interval); const sites = studiesToSites(studies); const [studyGoals, replayData] = await Promise.all([ fetchStudyGoals(), getAnalyticsSiteScoreReplayData({ interval, studies }), ]); // 预构建replayData的查找Map,优化查询效率 const replayDataMap = new Map(); replayData.forEach(item => { const key = `${item.day}-${item.studyId}-${item.siteId}`; replayDataMap.set(key, item); }); let currentBatch = []; let batchIndex = 0; let isFirstBatch = true; // 遍历days和sites组合,边生成边积累批次 for (const day of days) { const isoDay = day.toISOString(); for (const site of sites) { const { studyId, siteId } = site; const replayItem = replayDataMap.get(`${isoDay}-${studyId}-${siteId}`); const data = { staleCandidates: getStaleCandidatesFromData({ data: analyticsData, site, day, noContact: studyGoals?.inactive?.rules?.noContact, }), ...replayItem, }; const calculator = new SiteScoreCalculator({ site, studyGoals, data, day }); const inputItem = siteToInput({ days: [day], site, calculator }); // 过滤空分数项 if (removeNullScores(inputItem)) { currentBatch.push(inputItem); // 达到批次大小,执行保存 if (currentBatch.length >= BATCH_SIZE) { const cnt = await saveSiteScore({ input: currentBatch, prune: isFirstBatch ? prune : null }); batchIndex++; log.debug('Replay site score batch %d (%d data points) saved', batchIndex, cnt); currentBatch = []; // 清空批次释放内存 isFirstBatch = false; } } } } // 处理剩余未达批次大小的数据 if (currentBatch.length > 0) { const cnt = await saveSiteScore({ input: currentBatch, prune: isFirstBatch ? prune : null }); batchIndex++; log.debug('Replay site score batch %d (%d data points) saved', batchIndex, cnt); } log.info('Replay site scores completed in %ds', (Date.now() - start) / 1e3); } export default replaySiteScores;
额外优化建议
- 调整
BATCH_SIZE:根据服务器内存情况测试最优值(如50、150),平衡内存占用与请求频率。 - 内存监控:加入内存使用日志验证优化效果:
log.debug('Current memory usage: %o', process.memoryUsage());
内容的提问来源于stack exchange,提问作者Jakub
相关产品推荐
相关产品推荐

