| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697 |
- from __future__ import annotations
- import logging
- from typing import TYPE_CHECKING
- from apscheduler.schedulers.background import BackgroundScheduler
- from apscheduler.triggers.cron import CronTrigger
- from supply_infra.config import get_infra_settings
- from supply_infra.pipeline.run_service import submit_pipeline_run
- from supply_infra.scheduler.constants import SUPPLY_PIPELINE_JOB_ID, SUPPLY_PIPELINE_JOB_NAME
- if TYPE_CHECKING:
- from apscheduler.schedulers.base import BaseScheduler
- logger = logging.getLogger(__name__)
- _scheduler: BackgroundScheduler | None = None
- def submit_scheduled_pipeline() -> dict:
- """Cron callback: create a durable batch and return immediately."""
- return submit_pipeline_run(
- trigger_type="cron",
- trigger_source="scheduler",
- ).to_dict()
- def create_scheduler() -> BackgroundScheduler:
- """Create a Scheduler that only submits durable pipeline batches."""
- settings = get_infra_settings()
- tz = settings.scheduler_timezone
- scheduler = BackgroundScheduler(timezone=tz)
- scheduler.add_job(
- submit_scheduled_pipeline,
- trigger=CronTrigger(
- hour=settings.scheduler_cron_hour,
- minute=settings.scheduler_cron_minute,
- timezone=tz,
- ),
- id=SUPPLY_PIPELINE_JOB_ID,
- name=SUPPLY_PIPELINE_JOB_NAME,
- replace_existing=True,
- max_instances=1,
- coalesce=True,
- misfire_grace_time=3600,
- )
- logger.info(
- "Scheduler configured with %d job(s) | timezone=%s | pipeline_cron=%02d:%02d",
- len(scheduler.get_jobs()),
- tz,
- settings.scheduler_cron_hour,
- settings.scheduler_cron_minute,
- )
- return scheduler
- def _log_next_runs(scheduler: BaseScheduler) -> None:
- for job in scheduler.get_jobs():
- # apscheduler 3.x: 在 scheduler.start() 之前,Job.next_run_time 访问会抛 AttributeError
- # 这里用 getattr 兜底,保证启动阶段不因日志而中断。
- next_run_time = getattr(job, "next_run_time", None)
- logger.info(" - %s | next run: %s", job.name, next_run_time)
- def start_scheduler() -> BackgroundScheduler | None:
- """Start the independent scheduler process (idempotent)."""
- global _scheduler
- settings = get_infra_settings()
- if not settings.scheduler_enabled:
- logger.warning("Scheduler is disabled (SCHEDULER_ENABLED=false)")
- return None
- if _scheduler is not None and _scheduler.running:
- return _scheduler
- _scheduler = create_scheduler()
- logger.info("Starting scheduler...")
- _scheduler.start()
- _log_next_runs(_scheduler)
- return _scheduler
- def stop_scheduler() -> None:
- """Shut down the background scheduler if running."""
- global _scheduler
- if _scheduler is None or not _scheduler.running:
- return
- logger.info("Stopping scheduler...")
- _scheduler.shutdown(wait=False)
- _scheduler = None
- logger.info("Scheduler stopped.")
|