news 2026/9/24 17:05:51

Agenda 的 `stop()` 保留运行中任务锁:滚动重启下如何防止同一任务被重复并发执行

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Agenda 的 `stop()` 保留运行中任务锁:滚动重启下如何防止同一任务被重复并发执行

【免费下载链接】agenda

Lightweight job scheduling for Node.js

项目地址:https://gitcode.com/gh_mirrors/ag/agenda
点击查看免费下载

本篇技术指南围绕 Agenda 仓库中的变更记录.changeset/stop-preserves-running-locks.md展开,解析stop()语义变更的来龙去脉:为什么停止调度器时不再解锁正在运行的任务、锁的生命周期如何工作、源码是如何实现的,以及生产环境滚动重启时应当如何正确关闭实例。

变更背景:stop()的旧语义带来重复执行风险

agenda是一个面向 Node.js 的轻量级任务调度库(见 README.md),多个进程/实例可以共享同一个数据库后端(MongoDB、PostgreSQL 或 Redis),通过「数据库锁」来协调任务的互斥执行。

在变更之前,agenda.stop()的行为是:清除处理任务的定时器,并解锁所有当前被本实例锁定的任务(包括正在运行的任务)。这在单实例场景下没有问题,但在**滚动重启(rolling restart)**场景下会引发一个严重缺陷:

  1. 实例 A 锁定了任务 J 并开始执行;
  2. 部署过程中实例 A 被SIGTERM通知,随后调用stop()
  3. 旧版stop()会把任务 J 的数据库锁一并释放;
  4. 实例 B 扫描数据库时发现任务 J 的锁已被释放,立即重新锁定并执行;
  5. 结果:同一个任务出现(job occurrence)被两个实例并发执行,对于不具备幂等性的任务会造成数据重复处理、副作用叠加。

本次变更的核心正是修复这一点:stop()不再解锁仍在运行的任务。变更记录原文如下(见 .changeset/stop-preserves-running-locks.md):

stop()no longer unlocks jobs that are still running, preventing duplicate concurrent execution of the same job occurrence on another instance during rolling restarts. Only jobs that are locked locally but have not started running are unlocked; running jobs keep their database lock until they complete (orlockLifetimeexpires if the process dies).

翻译过来即:stop()现在只释放「本地已锁定但尚未开始运行」的任务;正在运行的任务保留其数据库锁,直到任务完成,或者进程死亡后lockLifetime到期(由其他实例接管)。

锁的生命周期:从lockedAtlockLifetime

要理解这次变更,需要先弄清 Agenda 的锁模型。

JobProcessor从数据库取出一个可执行任务时(findAndLockNextJob),它会通过后端仓库的getNextJobToRun以原子方式锁定任务并写入lockedAt时间戳:

const lockDeadline = new Date(Date.now().valueOf() - definition.lockLifetime); const result = await this.agenda.db.getNextJobToRun( jobName, this.nextScanAt, lockDeadline, undefined, { lastModifiedBy: this.agenda.attrs.name || undefined } );

(见 JobProcessor.ts)

这里的关键是lockLifetime:它定义了锁的有效期。只有lockedAt早于lockDeadline(即锁已过期)的任务才会被其他实例重新拾取,这是进程崩溃后任务能够恢复执行(stale-lock recovery)的基础。lockLifetime是任务定义(define)的一个配置项,单位毫秒。

而解锁的语义由JobRepository接口定义(见 JobRepository.ts):

/** * Attempt to lock a job for processing */ lockJob(job, options): Promise<JobParameters | undefined>; /** * Unlock a single job */ unlockJob(job: JobParameters): Promise<void>; /** * Unlock multiple jobs by ID */ unlockJobs(jobIds: (JobId | string)[]): Promise<void>;

unlockJobs正是Agenda.stop()释放锁时调用的底层方法。锁本质上就是lockedAt字段:置空即解锁,写入时间戳即锁定。

源码剖析:新stop()的精确实现

JobProcessor 层:只返回「已锁定但未运行」的任务

JobProcessor内部维护了两个关键数组:

  • runningJobs:当前真正在执行处理器回调的任务;
  • lockedJobs:已经从数据库锁定、进入本地队列但未必已开始运行的任务。

新的stop()实现(见 JobProcessor.ts)先构建运行中任务的 ID 集合,再从lockedJobs中过滤出不在运行集合里的任务:

stop(): JobWithId[] { log.extend('stop')('stop job processor', this.isRunning); this.isRunning = false; if (this.processInterval) { clearInterval(this.processInterval); this.processInterval = undefined; } // Unsubscribe from notifications if (this.notificationUnsubscribe) { log.extend('stop')('unsubscribing from notification channel'); this.notificationUnsubscribe(); this.notificationUnsubscribe = undefined; } const runningJobIds = new Set(this.runningJobs.map(job => job.attrs._id.toString())); // Only unlock jobs that are held locally but have not started running. // Running jobs keep their database locks until they complete or the lock expires. return this.lockedJobs.filter(job => !runningJobIds.has(job.attrs._id.toString())); }

注意stop()还做了两件伴随工作:

  • isRunning置为false,这会让后续的process()runOrRetry()提前返回(见 JobProcessor.ts 与 JobProcessor.ts),防止停下来的处理器继续拉取新任务;
  • 取消通知通道订阅,停止接收新任务到达的实时通知。

Agenda 层:只对返回的任务调用unlockJobs

Agenda.stop()(见 index.ts)在调用jobProcessor.stop()拿到需要解锁的列表后,仅对该列表执行数据库解锁:

async stop(closeConnection?: boolean): Promise<void> { if (!this.jobProcessor) { log('Agenda.stop called, but agenda has never started!'); return; } const lockedJobs = this.jobProcessor.stop(); const jobIds = lockedJobs?.map(job => job.attrs._id) || []; if (jobIds.length > 0) { log('about to unlock jobs with ids: %O', jobIds); await this.db.unlockJobs(jobIds); } // Unsubscribe from state notifications // Disconnect notification channel if configured // Close backend connection (defaults to backend.ownsConnection) this.jobProcessor = undefined; }

于是「解锁哪些任务」的判定完全交由JobProcessor完成,Agenda层只负责将筛选出的任务 ID 批量解锁。运行中的任务不在返回值里,它们的lockedAt得以保留在数据库中。

运行中任务何时释放锁?

运行中的任务在三种情况下会释放数据库锁:

  1. 正常完成runOrRetryfinally块中,任务完成/失败后会把该任务同时从runningJobslockedJobs移除,并由Job.run()内部的完成逻辑处理后序状态(见 JobProcessor.ts);
  2. lockLifetime到期(进程仍存活)runOrRetry中的checkIfJobIsStillAlive周期性检查(间隔取processEvery / 2lockLifetime / 2的较大值),一旦检测到job.isExpired()(执行时长超过lockLifetime),会抛出异常终止执行,此时锁因超时失效,可被其他实例接管(见 JobProcessor.ts)。因此长任务必须调用job.touch()续期;
  3. 进程死亡:进程退出后不再续期,lockLifetime到期后锁自然过期,其他实例通过过期锁恢复路径重新接管。

这也正是变更记录中「running jobs keep their database lock until they complete (orlockLifetimeexpires if the process dies)」的含义:即使stop()被调用,运行中的任务锁也一直保留,直到任务自行结束或锁超时,从而堵住了滚动重启期间另一实例提前接管同一任务的口子。

测试佐证:行为边界的精确锁定

仓库测试套件 agenda-test-suite.ts 用一组用例精确刻画了新语义的边界,可以在本地运行pnpm --filter agenda test验证:

  • 运行中的任务,stop()后锁仍然保留:先触发clear-lock-test任务开始执行,再调用agenda.stop(),随后查询数据库,result.jobs[0].lockedAt仍为真值(见 agenda-test-suite.ts);
  • 已锁定但未运行的任务,stop()后被解锁queued-lock-test场景中,队列里存在多个已锁定任务但只有一个正在运行,stop()之后数据库中保留lockedAt的任务数恰好为 1(即正在运行的那个),其余均被释放(见 agenda-test-suite.ts);
  • 仅锁定、从未运行的任务,stop()后锁被清除scheduled-queued-lock-test中任务被锁定但runningJobs为 0,stop()后查询lockedAt为假值(见 agenda-test-suite.ts)。

这三个用例分别对应「运行中 → 保留锁」「混合状态 → 只解锁未运行的」「仅锁定 → 全部解锁」,完整覆盖了变更记录的语义声明。

生产实践:滚动重启与优雅关闭的正确姿势

场景一:滚动重启(多实例部署)

在 K8s、ECS 等多实例部署下进行滚动发布时,编排系统会给旧实例发送SIGTERM再等待一段时间后强制终止。得益于本次变更,旧实例的stop()不再释放运行中任务的锁,新实例不会立刻重复执行同一任务:

process.on('SIGTERM', async () => { // stop() 立即停止调度;运行中的任务锁被保留, // 不会在滚动重启时被其他实例重复接管 await agenda.stop(); process.exit(0); });

需要权衡的是:stop()是立即返回的(不会等待运行中任务完成),因此运行中的任务会在本进程内被"遗留",其锁只能靠lockLifetime自然过期后由其他实例接管。这意味着:

  • 若任务执行时长通常短于lockLifetime,遗留任务的恢复会有一定延迟(需等锁过期);
  • 建议给任务定义设置合理的lockLifetime(默认值可参考 JobDefinition.ts),避免过长导致故障任务长时间无法被接管、过短导致慢任务频繁被判定过期中断。

场景二:优雅关闭(希望等待任务跑完)

如果业务允许停机等待,推荐使用drain()而非stop()drain()会停止接收新任务,但等待所有运行中的任务完成后再关闭,与stop()形成互补(见 JobProcessor.ts 与 index.ts):

// 等待所有任务完成(可带超时或 AbortSignal) await agenda.drain(); // 带超时的关闭 const result = await agenda.drain(30_000); if (result.timedOut) { console.log(`${result.running} jobs still running, forcing stop`); await agenda.stop(); }

drain()DrainResult会返回{ completed, running, timedOut, aborted }四类统计,便于在超时后决定是否强制stop()

完整的可运行示例见 graceful-shutdown.ts,它演示了SIGTERM/SIGINT下的优雅关闭、drain()等待、超时兜底stop()的完整模式,可用npx tsx examples/graceful-shutdown.ts直接运行体验(需要本地 MongoDB)。

场景三:进程崩溃(无stop()调用)

如果进程是直接崩溃或被kill -9强制终止,stop()根本不会被调用,此时依赖的正是lockLifetime过期机制:其他实例扫描到过期锁(lockedAt早于lockDeadline)后接管任务。本次变更不改变这一路径,崩溃恢复语义保持不变。

结论与选型建议

.changeset/stop-preserves-running-locks.md记录了一次小而关键的语义修正,把stop()的职责从「清空一切锁」收敛为「仅释放本地未运行任务的锁」,使滚动重启下的任务去重得到保证:

任务状态stop()行为stop()行为
正在运行(在runningJobs中)解锁 → 可能被其他实例重复执行保留锁,直至完成或lockLifetime过期
已锁定未运行(在lockedJobs中)解锁解锁(行为不变)
进程崩溃(未调用stop()锁在lockLifetime后过期锁在lockLifetime后过期(行为不变)

工程上的建议:

  1. 滚动重启优先使用drain()+ 超时兜底stop(),兼顾「不丢任务」与「不重复执行」;
  2. 为每个任务定义显式配置lockLifetime,并在长任务中周期性调用job.touch()续期,这是锁语义正确工作的前提;
  3. 任务尽量设计为幂等:锁机制降低的是并发概率而非绝对消除,幂等仍是分布式任务处理的最后防线。

对于想深入研究的读者,建议继续阅读:任务锁定的核心判定逻辑 JobProcessor.ts、stop()/drain()的完整实现 index.ts、锁相关接口定义 JobRepository.ts,以及后端仓库对getNextJobToRun/unlockJobs的具体实现(如 MongoJobRepository.ts、PostgresJobRepository.ts、RedisJobRepository.ts),它们共同构成了完整的任务互斥与故障恢复机制。

【免费下载链接】agenda

Lightweight job scheduling for Node.js

项目地址:https://gitcode.com/gh_mirrors/ag/agenda
点击查看免费下载

相关推荐

上一篇:如何用 Lottie 在网页跑通 After Effects 动画:5 步上手
下一篇:BT下载加速:3分钟配好Tracker的完整方法

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/24 17:02:52

大麦自动抢票指南:Python Selenium + Appium 双端抢票脚本完整教程

大麦自动抢票指南&#xff1a;Python Selenium Appium 双端抢票脚本完整教程 【免费下载链接】ticket-purchase 大麦自动抢票&#xff0c;支持人员、城市、日期场次、价格选择 项目地址: https://gitcode.com/GitHub_Trending/ti/ticket-purchase ticket-purchase 是一…

作者头像 李华
网站建设 2026/9/24 17:01:00

高速传输灵活交付 金士顿移动存储赋能项目全周期数据流转迁移

乙方项目归档交付应包含完整项目的原始素材、源文件、多版迭代稿件、最终成片与交付文档等等。而实际上&#xff0c;很多行业往往需要混合办公、跨地协作&#xff0c;依托网盘存储看似便利实际暗藏隐患&#xff0c;不仅容易出现版本错乱、链接过期、文件压缩损坏、画质音质失真…

作者头像 李华