- Status:
connected - Done: schedule registry、provider、trigger log repository、持久
ScheduleState、cron misfire policy、锁保护的触发路径、scheduler 背景上下文和 trace_id handoff、core scheduler --run-once和core scheduler --runcron due 本地运行入口、按 profile 选择 Tasks provider 的提交链路、scheduler profile 运行参数、APScheduler/Celery Beat provider adapter、scheduler audit gate 和分布式 lock provider 串通已落地。 - Next: none
Scheduler 模块负责“何时触发任务”。它不负责具体任务怎么执行,具体执行交给 Tasks 模块。
Celery 可以实现 scheduler,但 core 不应直接绑定 Celery。
Scheduler
定义触发时间、周期、错过触发策略、启停和注册。
Tasks
定义任务提交、执行、状态、重试和结果。
Celery Beat
可以作为 Scheduler provider。
Celery Worker
可以作为 Tasks provider。
如果项目选择 Celery 技术栈,推荐组合是:
core.scheduler provider = celery_beat
core.tasks provider = celery
本地开发或单机版可以使用:
core.scheduler provider = apscheduler
core.tasks provider = sync
src/core/scheduler/
provider.py
registry.py
external.py
interval
cron
date
manual
- 清理临时文件。
- 归档审计日志。
- 刷新外部配置缓存。
- 生成周期报表。
- 扫描超时任务。
- 调度器只负责触发,不直接写业务逻辑。
- 分布式部署时必须保证同一调度不会多实例重复触发。
- 可以通过 Locks 模块保护单例调度。
- app 通过 module 注册 schedule definitions。
- 生产环境需要能查看调度状态和最近执行记录。
- schedule trigger 必须记录
schedule_id、planned_at、triggered_at、task_id、status。 - 错过触发策略必须显式声明:
skip、run_once或catch_up_limited。 - 周期任务提交必须绑定幂等键,例如
schedule_id + planned_at。 - scheduler provider 可注入
ScheduleTriggerRepository写ScheduleTriggerLog;同一schedule_id + tenant_id + planned_at只能保留一条触发历史,重复触发返回已有历史并依赖 task idempotency 避免重复执行。 - scheduler 只提交任务或写 outbox,不直接执行业务逻辑。
第一版先提供 ScheduleRegistry:
- 从
AppModule.schedules收集 schedule definition。 - schedule_id 全局唯一。
- 每个 schedule 的 task_type 必须能在
TaskRegistry中找到。 ManualScheduleProvider提供本地/运维触发入口,读取ScheduleRegistry,构造带schedule_id + tenant_id + planned_at幂等键的TaskEnvelope,再提交给 Tasks provider。LockedScheduleProvider可包装任意 scheduler provider,在触发前获取scheduler:trigger:{schedule_id}:{tenant_id}:{planned_at}锁,触发完成或失败后释放锁,避免同一实例集内重复触发同一 planned slot。AuditedScheduleProvider可包装任意 scheduler provider,将scheduler.triggered写入 audit gate;成功记录 task、planned slot、lock key/fencing token,失败记录异常 reason。APSchedulerScheduleProvider读取同一份ScheduleRegistry导出 APScheduler job specs,并在外部 job 触发时复用ScheduleTriggerRequest/ScheduleTriggerResult。CeleryBeatScheduleProvider读取同一份ScheduleRegistry导出 Celery Beat entry,entry 指向统一core.scheduler.trigger桥接任务,不直接调用业务 handler。ScheduleTriggerLog保存schedule_id、tenant_id、planned_at、triggered_at、task_id、task_type、status、request_id和错误信息。ScheduleTriggerRepository.record_result()使用 insert-first + 唯一约束记录触发历史;重复 trigger key 返回replayed,不创建第二条历史。ScheduleState持久保存schedule_id + tenant_id的last_planned_at、last_triggered_at和last_checked_at,避免 scheduler 重启后重复处理同一 cron slot。ScheduleStateRepository.plan_cron_due_slots()支持skip、run_once和catch_up_limitedmisfire policy;catch_up_limited可通过trigger_config.misfire_limit限制单轮补偿数量。- scheduler provider 不直接调用业务函数,tenant lifecycle gate 仍由 task provider 执行。
- scheduler trigger 执行期间会从
ScheduleTriggerRequest注入冻结背景上下文,透传request_id、trace_id和tenant_id,避免继承外层 HTTP/CLI ContextVar。 core scheduler --run-once可按--installed-app加载ScheduleRegistry/TaskRegistry,通过当前 profile 选中的 Tasks provider 触发指定 schedule,写入TaskRun和ScheduleTriggerLog,并输出稳定 JSON。core scheduler --run会扫描 app 注册的 cron schedule definition,按持久ScheduleState规划 due slot;--provider local|apscheduler|celery_beat控制 scheduler provider adapter;--max-iterations用于 local/CI 有限轮验证,不传则常驻轮询;同一schedule_id + tenant_id + planned_atslot 跨重启只尝试一次,传--instance-id时写入 scheduler heartbeat。- profile 模板通过
SCHEDULER__PROVIDER、SCHEDULER__IDLE_SLEEP_SECONDS和SCHEDULER__LOCK_TTL_SECONDS参数化 scheduler loop。
APScheduler 或 Celery Beat provider 必须读取同一份 ScheduleRegistry,复用 ScheduleTriggerRequest/ScheduleTriggerResult 语义,并把触发结果提交到 Tasks provider,而不是直接调用业务函数。private/cloud 多实例部署时,LockedScheduleProvider 必须注入 Redis、数据库 advisory lock 或等价的分布式 LockProvider;内存 lock 只适合 local/profile 和测试。