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

如何优化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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 00:42:41