使用Agenda时多个定时任务丢失问题排查求助
你的代码存在以下几个关键疏漏,直接导致了任务随机丢失的问题:
1. 异步任务并发执行未做等待
triggerJobs里用forEach遍历任务并调用异步的scheduleJob,但forEach不会等待每个异步任务完成,100个任务的scheduleJob会同时并发执行。Agenda实例在处理大量并发的任务定义、定时设置操作时,会出现数据库竞态,部分任务的持久化操作被覆盖或失败。
修复:将forEach改为for...of循环,确保每个任务的异步操作完成后再处理下一个:
public async triggerJobs(): Promise<void> { try { const jobs: Job[] = await this.repostory.findJobs(); for (const job of jobs) { await this.scheduleJob(job); } } catch (err: any) { this.logger.error(err); } }
2. 重复调用agenda.start
每个任务都调用await this.agenda.start({ name: job.id }),但agenda.start是用来启动整个Agenda调度器的方法,只需要调用一次即可。多次并发调用会打乱Agenda的内部状态,导致任务注册异常。
修复:在服务初始化时调用一次agenda.start,移除scheduleJob里的重复调用:
public async triggerJobs(): Promise<void> { try { // 仅启动一次调度器 await this.agenda.start(); const jobs: Job[] = await this.repostory.findJobs(); for (const job of jobs) { await this.scheduleJob(job); } } catch (err: any) { this.logger.error(err); } } public async scheduleJob(job: Job) { await this.agenda.cancel({ name: job.id }); if (job.status === JobStatus.ACTIVE) { this.agenda.define(job.id, (agendaJob, done) => { this.logger.debug(`--- Job Defined: Job ${job.name}-${job.schedule}-${job.id}`); this.executeJob(job); done(); }); await this.agenda.every(job.schedule, job.id); } }
3. agenda.cancel未等待异步完成
原代码中this.agenda.cancel({ name: job.id })是异步操作,但未加await,导致取消旧任务的操作还没完成,就开始定义新任务和设置定时,可能出现旧任务残留或者新任务被旧任务覆盖的情况。
修复:给agenda.cancel加上await,确保旧任务取消完成后再执行后续操作。
4. 任务定义与定时设置的逻辑冗余
每次处理任务都重新define,如果任务已经定义过,重复定义可能引发问题。可以先批量定义所有激活的任务,再批量设置定时,减少并发操作:
public async triggerJobs(): Promise<void> { try { const jobs: Job[] = await this.repostory.findJobs(); // 先批量定义所有激活任务 const activeJobs = jobs.filter(job => job.status === JobStatus.ACTIVE); activeJobs.forEach(job => { this.agenda.define(job.id, (agendaJob, done) => { this.logger.debug(`--- Job Defined: Job ${job.name}-${job.schedule}-${job.id}`); this.executeJob(job); done(); }); }); await this.agenda.start(); // 批量处理任务的取消与定时设置 for (const job of jobs) { await this.agenda.cancel({ name: job.id }); if (job.status === JobStatus.ACTIVE) { await this.agenda.every(job.schedule, job.id); } } } catch (err: any) { this.logger.error(err); } }
5. 异步错误未被捕获
原triggerJobs中,forEach里的scheduleJob如果出错,错误不会被外层catch捕获,导致部分任务失败但无日志可查。改成for...of配合try/catch后,错误可以被正确捕获和记录。
内容的提问来源于stack exchange,提问作者Developer

