app.py 2.0 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859
  1. from __future__ import annotations
  2. import logging
  3. from apscheduler.schedulers.blocking import BlockingScheduler
  4. from apscheduler.triggers.cron import CronTrigger
  5. from supply_infra.config import get_infra_settings
  6. from supply_infra.scheduler.jobs.sync_global_tree_odps_to_mysql import sync_global_tree_odps_to_mysql
  7. from supply_infra.scheduler.jobs.sync_multi_demand_pool_odps_to_mysql import (
  8. sync_multi_demand_pool_odps_to_mysql,
  9. )
  10. logger = logging.getLogger(__name__)
  11. def create_scheduler() -> BlockingScheduler:
  12. """Create and configure the scheduler with all registered jobs."""
  13. settings = get_infra_settings()
  14. scheduler = BlockingScheduler(timezone=settings.scheduler_timezone)
  15. # 每天凌晨 2:30 从 ODPS 同步全局树元素与分类到 MySQL
  16. scheduler.add_job(
  17. sync_global_tree_odps_to_mysql,
  18. trigger=CronTrigger(hour=2, minute=30),
  19. id="sync_global_tree_odps_to_mysql",
  20. name="ODPS → MySQL 全局树同步",
  21. replace_existing=True,
  22. )
  23. # 每天 12:00 从 ODPS 同步策略需求天级表到 MySQL(当天 dt)
  24. scheduler.add_job(
  25. sync_multi_demand_pool_odps_to_mysql,
  26. trigger=CronTrigger(hour=12, minute=0),
  27. id="sync_multi_demand_pool_odps_to_mysql",
  28. name="ODPS → MySQL 策略需求池同步",
  29. replace_existing=True,
  30. )
  31. logger.info("Scheduler configured with %d job(s)", len(scheduler.get_jobs()))
  32. return scheduler
  33. def run_scheduler() -> None:
  34. """Start the blocking scheduler (CLI entry point)."""
  35. settings = get_infra_settings()
  36. if not settings.scheduler_enabled:
  37. logger.warning("Scheduler is disabled (SCHEDULER_ENABLED=false)")
  38. return
  39. scheduler = create_scheduler()
  40. logger.info("Starting scheduler...")
  41. for job in scheduler.get_jobs():
  42. logger.info(" - %s | next run: %s", job.name, job.next_run_time)
  43. try:
  44. scheduler.start()
  45. except (KeyboardInterrupt, SystemExit):
  46. logger.info("Scheduler stopped.")