Skip to content

Latest commit

 

History

History
106 lines (80 loc) · 5.78 KB

File metadata and controls

106 lines (80 loc) · 5.78 KB

Core Scheduler

Progress

  • Status: connected
  • Done: schedule registry、provider、trigger log repository、持久 ScheduleState、cron misfire policy、锁保护的触发路径、scheduler 背景上下文和 trace_id handoff、core scheduler --run-oncecore scheduler --run cron due 本地运行入口、按 profile 选择 Tasks provider 的提交链路、scheduler profile 运行参数、APScheduler/Celery Beat provider adapter、scheduler audit gate 和分布式 lock provider 串通已落地。
  • Next: none

职责

Scheduler 模块负责“何时触发任务”。它不负责具体任务怎么执行,具体执行交给 Tasks 模块。

与 Celery 的关系

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_idplanned_attriggered_attask_idstatus
  • 错过触发策略必须显式声明:skiprun_oncecatch_up_limited
  • 周期任务提交必须绑定幂等键,例如 schedule_id + planned_at
  • scheduler provider 可注入 ScheduleTriggerRepositoryScheduleTriggerLog;同一 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_idtenant_idplanned_attriggered_attask_idtask_typestatusrequest_id 和错误信息。
  • ScheduleTriggerRepository.record_result() 使用 insert-first + 唯一约束记录触发历史;重复 trigger key 返回 replayed,不创建第二条历史。
  • ScheduleState 持久保存 schedule_id + tenant_idlast_planned_atlast_triggered_atlast_checked_at,避免 scheduler 重启后重复处理同一 cron slot。
  • ScheduleStateRepository.plan_cron_due_slots() 支持 skiprun_oncecatch_up_limited misfire policy;catch_up_limited 可通过 trigger_config.misfire_limit 限制单轮补偿数量。
  • scheduler provider 不直接调用业务函数,tenant lifecycle gate 仍由 task provider 执行。
  • scheduler trigger 执行期间会从 ScheduleTriggerRequest 注入冻结背景上下文,透传 request_idtrace_idtenant_id,避免继承外层 HTTP/CLI ContextVar。
  • core scheduler --run-once 可按 --installed-app 加载 ScheduleRegistry/TaskRegistry,通过当前 profile 选中的 Tasks provider 触发指定 schedule,写入 TaskRunScheduleTriggerLog,并输出稳定 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_at slot 跨重启只尝试一次,传 --instance-id 时写入 scheduler heartbeat。
  • profile 模板通过 SCHEDULER__PROVIDERSCHEDULER__IDLE_SLEEP_SECONDSSCHEDULER__LOCK_TTL_SECONDS 参数化 scheduler loop。

APScheduler 或 Celery Beat provider 必须读取同一份 ScheduleRegistry,复用 ScheduleTriggerRequest/ScheduleTriggerResult 语义,并把触发结果提交到 Tasks provider,而不是直接调用业务函数。private/cloud 多实例部署时,LockedScheduleProvider 必须注入 Redis、数据库 advisory lock 或等价的分布式 LockProvider;内存 lock 只适合 local/profile 和测试。