NestJS定时任务持久化咨询:避免服务器崩溃丢失未完成任务
解决NestJS延迟任务服务器崩溃丢失问题的方案(无需独立存储服务)
核心思路
既然不想依赖独立存储服务器,可将未完成任务持久化到本地文件或应用内嵌的轻量存储中,服务重启时自动恢复未执行的任务。关键是不能直接存储函数,需把回调逻辑抽象为可序列化的任务标识+参数形式。
具体可行方案
1. 本地JSON文件持久化
这种方案实现简单,适合小型应用场景:
- 定义可序列化的任务结构,每次创建延迟任务时写入JSON文件;
- 服务启动时读取文件,计算任务剩余延迟时间并恢复执行;
- 任务完成后从文件中移除对应记录。
修改后的示例代码:
import { writeFileSync, readFileSync, existsSync } from 'fs'; import { join } from 'path'; import { Injectable } from '@nestjs/common'; import { SchedulerRegistry } from '@nestjs/schedule'; import { nanoid } from 'nanoid'; // 可序列化的待执行任务结构 interface PendingTask { taskId: string; taskType: string; // 任务类型标识,比如"send_user_notice" params: Record<string, any>; // 任务执行所需参数 delayMinutes: number; createdAt: number; // 任务创建时间戳 } @Injectable() export class TasksService { private readonly taskStoragePath = join(process.cwd(), 'pending-tasks.json'); constructor(private schedulerRegistry: SchedulerRegistry) { // 服务启动时自动恢复未完成任务 this.recoverPendingTasks(); } execAfterGivenMinutes(taskType: string, params: Record<string, any>, minutes: number) { const taskId = 'after_minutes_' + nanoid(); const createdAt = Date.now(); // 写入任务到本地文件 this.savePendingTask({ taskId, taskType, params, delayMinutes: minutes, createdAt }); const timeout = setTimeout(() => { // 根据任务类型执行对应逻辑 this.executeTask(taskType, params); this.deleteTimeout(taskId); // 任务完成后从文件移除 this.removePendingTask(taskId); }, minutes * 60000); this.schedulerRegistry.addTimeout(taskId, timeout); } // 抽象任务执行逻辑,替代直接传递回调函数 private executeTask(taskType: string, params: Record<string, any>) { switch (taskType) { case 'send_user_notice': // 执行发送用户通知逻辑,使用params中的参数 console.log(`发送通知给用户:${params.userId}`); break; case 'clean_temp_files': // 执行临时文件清理逻辑 console.log('清理临时文件'); break; // 可扩展更多任务类型 } } // 将任务写入JSON文件 private savePendingTask(task: PendingTask) { let tasks: PendingTask[] = []; if (existsSync(this.taskStoragePath)) { tasks = JSON.parse(readFileSync(this.taskStoragePath, 'utf-8')); } tasks.push(task); writeFileSync(this.taskStoragePath, JSON.stringify(tasks)); } // 移除已完成的任务记录 private removePendingTask(taskId: string) { if (!existsSync(this.taskStoragePath)) return; let tasks = JSON.parse(readFileSync(this.taskStoragePath, 'utf-8')); tasks = tasks.filter(task => task.taskId !== taskId); writeFileSync(this.taskStoragePath, JSON.stringify(tasks)); } // 恢复未完成的任务 private recoverPendingTasks() { if (!existsSync(this.taskStoragePath)) return; const tasks = JSON.parse(readFileSync(this.taskStoragePath, 'utf-8')); const now = Date.now(); tasks.forEach(task => { // 计算剩余延迟时间(转毫秒) const elapsedMinutes = (now - task.createdAt) / (1000 * 60); const remainingMinutes = task.delayMinutes - elapsedMinutes; if (remainingMinutes > 0) { // 剩余时间大于0,重新创建延迟任务 this.execAfterGivenMinutes(task.taskType, task.params, remainingMinutes); } else { // 已过执行时间,立即执行任务 this.executeTask(task.taskType, task.params); } }); // 清空文件,避免重复恢复任务 writeFileSync(this.taskStoragePath, JSON.stringify([])); } private deleteTimeout(name: string) { this.schedulerRegistry.deleteTimeout(name); } }
注意点:需处理文件读写的并发问题(比如加简单的文件锁),定期清理无效任务记录。
2. 内嵌式SQLite数据库
比JSON文件更可靠,支持事务和查询,适合任务量稍大的场景:
- 安装依赖:
@nestjs/typeorm和sqlite3; - 创建
PendingTask实体类,映射SQLite表; - 在
TasksService中注入Repository<PendingTask>,实现任务的增删查操作; - 服务启动时查询所有未完成任务,计算剩余延迟时间并恢复执行。
优势:数据操作更稳定,支持事务避免读写异常导致的数据丢失。
3. 内存映射文件(可选)
使用mmap-io等库实现内存映射文件存储,兼顾内存读写速度和磁盘持久化能力,适合对性能有要求的场景。核心逻辑与JSON文件方案类似,但读写效率更高,崩溃后数据不会丢失。
关键注意事项
- 避免直接序列化函数:必须将回调逻辑抽象为任务类型+参数的形式,否则无法持久化;
- 幂等性设计:任务恢复执行时,要确保重复执行不会产生副作用(比如重复发送通知),可给任务添加执行状态标记;
- 异常处理:捕获文件/数据库读写异常,避免影响主流程;任务执行失败时可添加重试机制。
内容的提问来源于stack exchange,提问作者vsay
相关产品推荐
相关产品推荐

