演示

定时任务

Quark ERP 支持两种 cron:进程内 @CronJob(基于 cron 包的本地定时器),以及分布式 CronJobService(BullMQ + Redis + Postgres)。共享/多实例调度用 BullMQ;单进程轻量 tick 用 @CronJob。

两种类型

按工作需要选择合适的方式。

  • 进程 cron(@CronJob)— @Service 被解析时启动的进程内调度;使用 Node cron 包(仅本进程)
  • BullMQ cron(CronJobService)— Redis 队列 + Postgres(cron_jobs)元数据;支持重试、重复上限、多实例安全

何时用哪种

  • BullMQ CronJobService — 领域清扫、提醒、对账等需跨重启且多实例正确的任务
  • 进程 @CronJob — 单进程上的轻量清理,或明确只要无 Redis 的内存 tick
  • 新工作不要用 @quan-erp-plugins/cron-schedular-backend——相同 API 已在核心以 CronJobService 提供

1. 进程 cron(@CronJob)

在 DI 解析的 @Service 方法上标记。类被解析时,核心按 expression 启动进程本地 CronJob。无数据库行、无 Redis 队列、无重试——每个加载该服务的进程各自 tick。

  • expression(必填)— cron 表达式,例如 0 * * * *
  • name(可选)— 日志标签
  • 类必须是 @Service()(或其它 DI 解析)并注册到模块 providers
TSX@CronJob
1import { CronJob, Module, Service } from "@quan-erp/shared-backend-core"; 2import metadata from "../../module.metadata.json" with { type: "json" }; 3 4@Service() 5export class LocalHousekeepingService { 6 @CronJob({ expression: "0 * * * *", name: "HourlyLocalCleanup" }) 7 async cleanup() { 8 // runs on this Node process only 9 } 10} 11 12@Module({ 13 name: metadata.name, 14 providers: [LocalHousekeepingService], 15 controllers: [], 16 entities: [], 17}) 18export class MyPluginModule {}

2. BullMQ cron(CronJobService)

任务经 Redis 上的 BullMQ 入队。调度元数据持久化在 Postgres(cron_jobs)。启动时会把活跃调度载入 Redis。回调仅在内存中——每次启动都要重新挂接监听器。

注入 CronJobService

用 ContainerRegistryManager.BUILTIN_PLUGIN 注入内置核心服务。

TSXinject
1import { 2 ContainerRegistryManager, 3 CronJobService, 4 Inject, 5 Service, 6} from "@quan-erp/shared-backend-core"; 7 8@Service() 9export class MaturedSweepService { 10 @Inject(CronJobService, ContainerRegistryManager.BUILTIN_PLUGIN) 11 private cronJobService: CronJobService; 12}

注册 BullMQ 任务

用唯一的 (pluginName, jobName) 调用 register。pluginName 使用 metadata.name。register 会 upsert 数据库行并(重新)创建 Redis 调度器。

TSXregister
1import { 2 ContainerRegistryManager, 3 CronJobService, 4 Inject, 5 Service, 6} from "@quan-erp/shared-backend-core"; 7import metadata from "../../module.metadata.json" with { type: "json" }; 8 9@Service() 10export class MaturedSweepService { 11 @Inject(CronJobService, ContainerRegistryManager.BUILTIN_PLUGIN) 12 private cronJobService: CronJobService; 13 14 async registerJob() { 15 await this.cronJobService.register({ 16 cronExpression: "0 2 * * *", 17 pluginName: metadata.name, 18 jobName: "matured-daily-sweep", 19 status: "active", 20 startDate: new Date(), 21 data: { reason: "daily-sweep" }, 22 retry: { attempt: 3, delay: 5_000, type: "exponential" }, 23 callback: async (job) => { 24 await this.runSweep(job.data); 25 }, 26 }); 27 } 28 29 private async runSweep(data: unknown) { /* … */ } 30}

重启后重新挂接监听器

启动时会把调度从 Postgres 载入 Redis。回调不会持久化——务必在 @OnAllModuleLoaded 重新注册(或通过 register 再次传入 callback)。

TSXlistener
1import { 2 ContainerRegistryManager, 3 CronJobService, 4 Inject, 5 OnAllModuleLoaded, 6 Service, 7} from "@quan-erp/shared-backend-core"; 8import metadata from "../../module.metadata.json" with { type: "json" }; 9 10@Service() 11export class MaturedSweepService { 12 @Inject(CronJobService, ContainerRegistryManager.BUILTIN_PLUGIN) 13 private cronJobService: CronJobService; 14 15 @OnAllModuleLoaded() 16 async onAllModuleLoaded() { 17 this.cronJobService.addEventListener( 18 metadata.name, 19 "matured-daily-sweep", 20 async (job) => { await this.runSweep(job.data); }, 21 ); 22 } 23 24 private async runSweep(data: unknown) { /* … */ } 25}

停止 / 移除(BullMQ)

  • stop(pluginName, jobName) — 暂停调度器并标为 inactive
  • remove(pluginName, jobName) — 删除 Redis 调度器、数据库行与监听器
  • stopAllByPluginName / removeAllByPluginName — 卸载时的批量辅助方法

查询辅助

TSXquery
1import { 2 ContainerRegistryManager, 3 CronJobService, 4 Inject, 5 Service, 6} from "@quan-erp/shared-backend-core"; 7import metadata from "../../module.metadata.json" with { type: "json" }; 8 9@Service() 10export class MaturedSweepService { 11 @Inject(CronJobService, ContainerRegistryManager.BUILTIN_PLUGIN) 12 private cronJobService: CronJobService; 13 14 async listJobs() { 15 const job = await this.cronJobService.getByPluginNameWithJobName( 16 metadata.name, 17 "matured-daily-sweep", 18 ); 19 20 const jobs = await this.cronJobService.getByPluginName({ 21 pluginName: metadata.name, 22 skip: 0, 23 limit: 50, 24 status: "active", 25 }); 26 27 return { job, jobs }; 28 } 29}

CreateCronJobDTO

  • cronExpression(必填)— 例如 0 2 * * *
  • pluginName / jobName(必填)— key 为 `${pluginName}/${jobName}`
  • status — active | inactive
  • data — BullMQ job.data 上的载荷
  • startDate — 最早开始时间(初始延迟)
  • repeat.limit — 最大执行次数
  • retry.attempt / delay / type(fixed | exponential)
  • callback — 注册时可选的内存监听器

推荐的 BullMQ 模式

  1. 01

    注入

    用 ContainerRegistryManager.BUILTIN_PLUGIN 注入 CronJobService。

  2. 02

    启动时注册

    在 @OnAllModuleLoaded 调用 register({ pluginName: metadata.name, jobName, cronExpression, … })。

  3. 03

    每次启动挂接监听器

    在 register 上传 callback 和/或 addEventListener——调度可跨重启,处理函数不能。

  4. 04

    清理

    卸载/禁用时 stop 或 remove,避免遗留 Redis 调度器。

TSXend-to-end
1import { 2 ContainerRegistryManager, 3 CronJobService, 4 Inject, 5 OnAllModuleLoaded, 6 Service, 7} from "@quan-erp/shared-backend-core"; 8import metadata from "../../module.metadata.json" with { type: "json" }; 9 10const JOB_NAME = "inventory-nightly-reconcile"; 11 12@Service() 13export class InventoryReconcileCronService { 14 @Inject(CronJobService, ContainerRegistryManager.BUILTIN_PLUGIN) 15 private cronJobService: CronJobService; 16 17 @OnAllModuleLoaded() 18 async onAllModuleLoaded() { 19 await this.cronJobService.register({ 20 cronExpression: "0 3 * * *", 21 pluginName: metadata.name, 22 jobName: JOB_NAME, 23 startDate: new Date(), 24 callback: async () => { await this.reconcile(); }, 25 }); 26 } 27 28 private async reconcile() { /* domain work */ } 29}