app.py 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138
  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 apscheduler.triggers.interval import IntervalTrigger
  7. from supply_infra.config import get_infra_settings
  8. from supply_infra.pipeline.run_service import submit_pipeline_run
  9. from supply_infra.scheduler.constants import (
  10. AIGC_PUBLISH_JOB_ID,
  11. AIGC_PUBLISH_JOB_NAME,
  12. SUPPLY_PIPELINE_JOB_ID,
  13. SUPPLY_PIPELINE_JOB_NAME,
  14. )
  15. if TYPE_CHECKING:
  16. from apscheduler.schedulers.base import BaseScheduler
  17. logger = logging.getLogger(__name__)
  18. _scheduler: BackgroundScheduler | None = None
  19. def submit_scheduled_pipeline() -> dict:
  20. """Cron callback: create a durable batch and return immediately."""
  21. return submit_pipeline_run(
  22. trigger_type="cron",
  23. trigger_source="scheduler",
  24. ).to_dict()
  25. def submit_aigc_publish() -> dict:
  26. """Interval callback: publish any unpublished primary discovery candidates."""
  27. from supply_infra.scheduler.jobs.publish_videos_from_discovery import (
  28. publish_videos_from_discovery,
  29. )
  30. result = publish_videos_from_discovery(
  31. skip_published=True,
  32. any_biz_dt=True,
  33. )
  34. logger.info(
  35. "AIGC publish job finished: success=%s candidates=%s batches=%s failed=%s",
  36. result.get("success"),
  37. result.get("candidate_count"),
  38. result.get("batch_count"),
  39. result.get("failed_batch_count"),
  40. )
  41. return result
  42. def create_scheduler() -> BackgroundScheduler:
  43. """Create a Scheduler for the daily pipeline and AIGC publish polling."""
  44. settings = get_infra_settings()
  45. tz = settings.scheduler_timezone
  46. scheduler = BackgroundScheduler(timezone=tz)
  47. scheduler.add_job(
  48. submit_scheduled_pipeline,
  49. trigger=CronTrigger(
  50. hour=settings.scheduler_cron_hour,
  51. minute=settings.scheduler_cron_minute,
  52. timezone=tz,
  53. ),
  54. id=SUPPLY_PIPELINE_JOB_ID,
  55. name=SUPPLY_PIPELINE_JOB_NAME,
  56. replace_existing=True,
  57. max_instances=1,
  58. coalesce=True,
  59. misfire_grace_time=3600,
  60. )
  61. scheduler.add_job(
  62. submit_aigc_publish,
  63. trigger=IntervalTrigger(
  64. minutes=settings.scheduler_aigc_publish_interval_minutes,
  65. timezone=tz,
  66. ),
  67. id=AIGC_PUBLISH_JOB_ID,
  68. name=AIGC_PUBLISH_JOB_NAME,
  69. replace_existing=True,
  70. max_instances=1,
  71. coalesce=True,
  72. misfire_grace_time=300,
  73. )
  74. logger.info(
  75. "Scheduler configured with %d job(s) | timezone=%s | pipeline_cron=%02d:%02d "
  76. "| aigc_publish_interval=%dm",
  77. len(scheduler.get_jobs()),
  78. tz,
  79. settings.scheduler_cron_hour,
  80. settings.scheduler_cron_minute,
  81. settings.scheduler_aigc_publish_interval_minutes,
  82. )
  83. return scheduler
  84. def _log_next_runs(scheduler: BaseScheduler) -> None:
  85. for job in scheduler.get_jobs():
  86. # apscheduler 3.x: 在 scheduler.start() 之前,Job.next_run_time 访问会抛 AttributeError
  87. # 这里用 getattr 兜底,保证启动阶段不因日志而中断。
  88. next_run_time = getattr(job, "next_run_time", None)
  89. logger.info(" - %s | next run: %s", job.name, next_run_time)
  90. def start_scheduler() -> BackgroundScheduler | None:
  91. """Start the independent scheduler process (idempotent)."""
  92. global _scheduler
  93. settings = get_infra_settings()
  94. if not settings.scheduler_enabled:
  95. logger.warning("Scheduler is disabled (SCHEDULER_ENABLED=false)")
  96. return None
  97. if _scheduler is not None and _scheduler.running:
  98. return _scheduler
  99. _scheduler = create_scheduler()
  100. logger.info("Starting scheduler...")
  101. _scheduler.start()
  102. _log_next_runs(_scheduler)
  103. return _scheduler
  104. def stop_scheduler() -> None:
  105. """Shut down the background scheduler if running."""
  106. global _scheduler
  107. if _scheduler is None or not _scheduler.running:
  108. return
  109. logger.info("Stopping scheduler...")
  110. _scheduler.shutdown(wait=False)
  111. _scheduler = None
  112. logger.info("Scheduler stopped.")