如何正确实现NestJS调度器/Cron?解决批量记录处理中途停止问题
我在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

