app.py 2.9 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697
  1. from __future__ import annotations
  2. import logging
  3. from typing import TYPE_CHECKING
  4. from apscheduler.schedulers.background import BackgroundScheduler
  5. from apscheduler.triggers.cron import CronTrigger
  6. from supply_infra.config import get_infra_settings
  7. from supply_infra.pipeline.run_service import submit_pipeline_run
  8. from supply_infra.scheduler.constants import SUPPLY_PIPELINE_JOB_ID, SUPPLY_PIPELINE_JOB_NAME
  9. if TYPE_CHECKING:
  10. from apscheduler.schedulers.base import BaseScheduler
  11. logger = logging.getLogger(__name__)
  12. _scheduler: BackgroundScheduler | None = None
  13. def submit_scheduled_pipeline() -> dict:
  14. """Cron callback: create a durable batch and return immediately."""
  15. return submit_pipeline_run(
  16. trigger_type="cron",
  17. trigger_source="scheduler",
  18. ).to_dict()
  19. def create_scheduler() -> BackgroundScheduler:
  20. """Create a Scheduler that only submits durable pipeline batches."""
  21. settings = get_infra_settings()
  22. tz = settings.scheduler_timezone
  23. scheduler = BackgroundScheduler(timezone=tz)
  24. scheduler.add_job(
  25. submit_scheduled_pipeline,
  26. trigger=CronTrigger(
  27. hour=settings.scheduler_cron_hour,
  28. minute=settings.scheduler_cron_minute,
  29. timezone=tz,
  30. ),
  31. id=SUPPLY_PIPELINE_JOB_ID,
  32. name=SUPPLY_PIPELINE_JOB_NAME,
  33. replace_existing=True,
  34. max_instances=1,
  35. coalesce=True,
  36. misfire_grace_time=3600,
  37. )
  38. logger.info(
  39. "Scheduler configured with %d job(s) | timezone=%s | pipeline_cron=%02d:%02d",
  40. len(scheduler.get_jobs()),
  41. tz,
  42. settings.scheduler_cron_hour,
  43. settings.scheduler_cron_minute,
  44. )
  45. return scheduler
  46. def _log_next_runs(scheduler: BaseScheduler) -> None:
  47. for job in scheduler.get_jobs():
  48. # apscheduler 3.x: 在 scheduler.start() 之前,Job.next_run_time 访问会抛 AttributeError
  49. # 这里用 getattr 兜底,保证启动阶段不因日志而中断。
  50. next_run_time = getattr(job, "next_run_time", None)
  51. logger.info(" - %s | next run: %s", job.name, next_run_time)
  52. def start_scheduler() -> BackgroundScheduler | None:
  53. """Start the independent scheduler process (idempotent)."""
  54. global _scheduler
  55. settings = get_infra_settings()
  56. if not settings.scheduler_enabled:
  57. logger.warning("Scheduler is disabled (SCHEDULER_ENABLED=false)")
  58. return None
  59. if _scheduler is not None and _scheduler.running:
  60. return _scheduler
  61. _scheduler = create_scheduler()
  62. logger.info("Starting scheduler...")
  63. _scheduler.start()
  64. _log_next_runs(_scheduler)
  65. return _scheduler
  66. def stop_scheduler() -> None:
  67. """Shut down the background scheduler if running."""
  68. global _scheduler
  69. if _scheduler is None or not _scheduler.running:
  70. return
  71. logger.info("Stopping scheduler...")
  72. _scheduler.shutdown(wait=False)
  73. _scheduler = None
  74. logger.info("Scheduler stopped.")