Agenda 的 `stop()` 保留运行中任务锁:滚动重启下如何防止同一任务被重复并发执行
【免费下载链接】agendaLightweight 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**场景下会引发一个严重缺陷实例 A 锁定了任务 J 并开始执行部署过程中实例 A 被SIGTERM通知随后调用stop()旧版stop()会把任务 J 的数据库锁一并释放实例 B 扫描数据库时发现任务 J 的锁已被释放立即重新锁定并执行结果同一个任务出现job occurrence被两个实例并发执行对于不具备幂等性的任务会造成数据重复处理、副作用叠加。本次变更的核心正是修复这一点stop()不再解锁仍在运行的任务。变更记录原文如下见 .changeset/stop-preserves-running-locks.mdstop()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到期由其他实例接管。锁的生命周期从lockedAt到lockLifetime要理解这次变更需要先弄清 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): PromiseJobParameters | undefined; /** * Unlock a single job */ unlockJob(job: JobParameters): Promisevoid; /** * Unlock multiple jobs by ID */ unlockJobs(jobIds: (JobId | string)[]): Promisevoid;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 层只对返回的任务调用unlockJobsAgenda.stop()见 index.ts在调用jobProcessor.stop()拿到需要解锁的列表后仅对该列表执行数据库解锁async stop(closeConnection?: boolean): Promisevoid { 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得以保留在数据库中。运行中任务何时释放锁运行中的任务在三种情况下会释放数据库锁正常完成runOrRetry的finally块中任务完成/失败后会把该任务同时从runningJobs和lockedJobs移除并由Job.run()内部的完成逻辑处理后序状态见 JobProcessor.tslockLifetime到期进程仍存活runOrRetry中的checkIfJobIsStillAlive周期性检查间隔取processEvery / 2与lockLifetime / 2的较大值一旦检测到job.isExpired()执行时长超过lockLifetime会抛出异常终止执行此时锁因超时失效可被其他实例接管见 JobProcessor.ts。因此长任务必须调用job.touch()续期进程死亡进程退出后不再续期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为 0stop()后查询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后过期行为不变工程上的建议滚动重启优先使用drain() 超时兜底stop()兼顾「不丢任务」与「不重复执行」为每个任务定义显式配置lockLifetime并在长任务中周期性调用job.touch()续期这是锁语义正确工作的前提任务尽量设计为幂等锁机制降低的是并发概率而非绝对消除幂等仍是分布式任务处理的最后防线。对于想深入研究的读者建议继续阅读任务锁定的核心判定逻辑 JobProcessor.ts、stop()/drain()的完整实现 index.ts、锁相关接口定义 JobRepository.ts以及后端仓库对getNextJobToRun/unlockJobs的具体实现如 MongoJobRepository.ts、PostgresJobRepository.ts、RedisJobRepository.ts它们共同构成了完整的任务互斥与故障恢复机制。赞分享【免费下载链接】agendaLightweight job scheduling for Node.js项目地址https://gitcode.com/gh_mirrors/ag/agenda点击查看免费下载相关推荐Miniflux 2 任务调度分布式锁防止重复执行Miniflux 2 任务调度分布式锁防止重复执行 在多实例部署的 Miniflux 2 环境中任务调度重复执行会导致资源浪费和数据不一致。本文将解析 Mi后端Conductor任务执行重启重新开始任务的执行Conductor任务执行重启重新开始任务的执行 你是否遇到过任务执行失败后需要重新启动的情况在微服务编排过程中任务执行中断或失败是常见问题。本文将详细介后端流程编排工作流自动化微服务vibe-tools并发执行如何同时运行多个AI任务vibe tools并发执行如何同时运行多个AI任务 想要让AI团队工作效率翻倍vibe tools的并发执行功能正是你需要的解决方案 这款强大的工具AI Agent开发工具MCP 服务浏览器控制上一篇如何用 Lottie 在网页跑通 After Effects 动画5 步上手下一篇BT下载加速3分钟配好Tracker的完整方法创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考