两种类型
按工作需要选择合适的方式。
- 进程 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
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 注入内置核心服务。
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 调度器。
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)。
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 — 卸载时的批量辅助方法
查询辅助
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 模式
- 01
注入
用 ContainerRegistryManager.BUILTIN_PLUGIN 注入 CronJobService。
- 02
启动时注册
在 @OnAllModuleLoaded 调用 register({ pluginName: metadata.name, jobName, cronExpression, … })。
- 03
每次启动挂接监听器
在 register 上传 callback 和/或 addEventListener——调度可跨重启,处理函数不能。
- 04
清理
卸载/禁用时 stop 或 remove,避免遗留 Redis 调度器。
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}