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

如何正确实现NestJS调度器/Cron?解决批量记录处理中途停止问题

NestJS调度器处理大量记录时中途停止的问题

我在NestJS框架中实现了一个调度器,它会获取所有未处理的记录,执行业务逻辑后将结果保存至另一张表。目前调度器运行基本正常,但在处理十万余条未处理记录时,仅完成约5500条后便停止。

this.processLogic()包含所有业务逻辑,同时会将unprocessedTable中的isProcessed字段更新为true。

我会跳过数据摄入阶段不完整的记录(例如用户尚未生成登出记录或仍在工作的情况),以此避免生成错误输出。

预计每日会有数千条数据接入,我不希望该流程在生产环境中崩溃。

目前我怀疑代码中可能存在意外创建的无限循环,但尚未完全确认。

@Timeout(0)
async processRecords() {
  const projectCodes = await this.codes.findMany();
  const now = new Date();
  const lilo =
    await this.unprocessedTable.findMany({
      where: {
        isProcessed: false,
      },
      include: {
        a: true,
        b: true,
        c: true,
      },
      orderBy: {
        d: 'desc',
      },
    });

  for (const record of lilo) {
    const GRACE_PERIOD = 5 * 60 * 60 * 1000; // +5 hours

    if (record.shiftDate && record.shiftFrom && record.shiftTo) {
      const shiftEnd = new Date(record.shiftDate);

      const hours = record.shiftTo.getUTCHours();
      const minutes = record.shiftTo.getUTCMinutes();
      const seconds = record.shiftTo.getUTCSeconds();

      shiftEnd.setUTCHours(hours, minutes, seconds, 0);

      if (record.shiftFrom > record.shiftTo) {
        shiftEnd.setDate(shiftEnd.getDate() + 1);
      }

      const processAfter = new Date(shiftEnd.getTime() + GRACE_PERIOD);

      if (now < processAfter) {
        console.log('Date to SKIP:', processAfter);
        continue;
      }

      if (!record.vcLastLogout || !record.proLastLogout) {
        continue;
      }
      console.log(`${record.recordId} timelog success!`);


      await this.processLogic(record, projectCodes);
    }
  }


  console.log('Cron executed');
}

问题排查与解决方案

1. 一次性加载大量数据引发内存溢出

当前代码通过findMany()一次性拉取所有未处理记录,十万条数据会占用大量内存,直接导致Node.js进程因内存不足被系统终止,这是最可能的原因。

修复方案:分页分批处理
修改查询逻辑,每次只拉取固定数量的记录,处理完一批再取下一批:

@Timeout(0)
async processRecords() {
  const projectCodes = await this.codes.findMany();
  const now = new Date();
  const batchSize = 100; // 可根据服务器配置调整
  let offset = 0;
  let hasMoreRecords = true;

  while (hasMoreRecords) {
    // 分页获取未处理记录
    const lilo = await this.unprocessedTable.findMany({
      where: { isProcessed: false },
      include: { a: true, b: true, c: true },
      orderBy: { d: 'desc' },
      take: batchSize,
      skip: offset,
    });

    if (lilo.length === 0) {
      hasMoreRecords = false;
      break;
    }

    // 处理当前批次的记录
    for (const record of lilo) {
      const GRACE_PERIOD = 5 * 60 * 60 * 1000; // +5小时

      if (record.shiftDate && record.shiftFrom && record.shiftTo) {
        const shiftEnd = new Date(record.shiftDate);

        const hours = record.shiftTo.getUTCHours();
        const minutes = record.shiftTo.getUTCMinutes();
        const seconds = record.shiftTo.getUTCSeconds();

        shiftEnd.setUTCHours(hours, minutes, seconds, 0);

        if (record.shiftFrom > record.shiftTo) {
          shiftEnd.setDate(shiftEnd.getDate() + 1);
        }

        const processAfter = new Date(shiftEnd.getTime() + GRACE_PERIOD);

        if (now < processAfter) {
          console.log('跳过日期:', processAfter);
          continue;
        }

        if (!record.vcLastLogout || !record.proLastLogout) {
          continue;
        }
        console.log(`${record.recordId} 工时记录处理成功!`);

        // 包裹异常捕获,避免单条记录失败终止整个流程
        try {
          await this.processLogic(record, projectCodes);
        } catch (error) {
          console.error(`处理记录${record.recordId}失败:`, error);
          // 可选:标记记录为处理失败,便于后续重试
          // await this.unprocessedTable.update({
          //   where: { recordId: record.recordId },
          //   data: { isProcessed: true, processFailed: true }
          // });
        }
      }
    }

    offset += batchSize;
  }

  console.log('调度任务执行完成');
}

2. 未捕获异常终止进程

如果processLogic()内部抛出未捕获的异常,会直接终止整个循环甚至进程,导致任务中途停止。

修复方案:添加异常捕获
如上面的代码所示,在调用processLogic时用try/catch包裹,确保单个记录处理失败不会影响整体流程。

3. 无限循环排查

当前代码本身没有明显的无限循环,但需要检查processLogic()中的更新逻辑:

  • 确认处理完成后,unprocessedTable中的isProcessed字段被正确设为true
  • 避免逻辑错误导致已处理记录被重新标记为未处理,引发重复处理

4. 超时限制问题

@Timeout(0)是项目启动后立即执行一次,如果处理时间过长,可能被系统或框架的超时机制中断。建议改用@Cron定时任务(比如每小时执行一次),并确保任务有足够的执行时间,或者拆分任务为更小的批次。


内容的提问来源于stack exchange,提问作者Gad Ashell

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.02 07:23:10