| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859 |
- from __future__ import annotations
- import logging
- from apscheduler.schedulers.blocking import BlockingScheduler
- from apscheduler.triggers.cron import CronTrigger
- from supply_infra.config import get_infra_settings
- from supply_infra.scheduler.jobs.sync_global_tree_odps_to_mysql import sync_global_tree_odps_to_mysql
- from supply_infra.scheduler.jobs.sync_multi_demand_pool_odps_to_mysql import (
- sync_multi_demand_pool_odps_to_mysql,
- )
- logger = logging.getLogger(__name__)
- def create_scheduler() -> BlockingScheduler:
- """Create and configure the scheduler with all registered jobs."""
- settings = get_infra_settings()
- scheduler = BlockingScheduler(timezone=settings.scheduler_timezone)
- # 每天凌晨 2:30 从 ODPS 同步全局树元素与分类到 MySQL
- scheduler.add_job(
- sync_global_tree_odps_to_mysql,
- trigger=CronTrigger(hour=2, minute=30),
- id="sync_global_tree_odps_to_mysql",
- name="ODPS → MySQL 全局树同步",
- replace_existing=True,
- )
- # 每天 12:00 从 ODPS 同步策略需求天级表到 MySQL(当天 dt)
- scheduler.add_job(
- sync_multi_demand_pool_odps_to_mysql,
- trigger=CronTrigger(hour=12, minute=0),
- id="sync_multi_demand_pool_odps_to_mysql",
- name="ODPS → MySQL 策略需求池同步",
- replace_existing=True,
- )
- logger.info("Scheduler configured with %d job(s)", len(scheduler.get_jobs()))
- return scheduler
- def run_scheduler() -> None:
- """Start the blocking scheduler (CLI entry point)."""
- settings = get_infra_settings()
- if not settings.scheduler_enabled:
- logger.warning("Scheduler is disabled (SCHEDULER_ENABLED=false)")
- return
- scheduler = create_scheduler()
- logger.info("Starting scheduler...")
- for job in scheduler.get_jobs():
- logger.info(" - %s | next run: %s", job.name, job.next_run_time)
- try:
- scheduler.start()
- except (KeyboardInterrupt, SystemExit):
- logger.info("Scheduler stopped.")
|