""" 定时任务调度器 基于 APScheduler,支持在 FastAPI 生命周期内运行定时任务 """ import os import logging from datetime import datetime from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger from .generator import run_creator from .optimizer import run_optimizer from .sync import sync_all_topics logger = logging.getLogger(__name__) class TaskScheduler: def __init__(self): self.scheduler = BackgroundScheduler() self._started = False def start(self): if self._started: logger.warning("Scheduler already started") return # 使用 CronTrigger 设置每日固定时间点 self.scheduler.add_job( self._run_sync, CronTrigger(hour=2, minute=30), id='scheduled_sync', replace_existing=True, max_instances=1, coalesce=True ) self.scheduler.add_job( self._run_generate, CronTrigger(hour=3, minute=30), id='scheduled_generate', replace_existing=True, max_instances=1, coalesce=True ) self.scheduler.add_job( self._run_optimize, CronTrigger(hour=4, minute=30), id='scheduled_optimize', replace_existing=True, max_instances=1, coalesce=True ) self.scheduler.start() self._started = True logger.info("Scheduler started with daily cron triggers (02:30, 03:30, 04:30)") def shutdown(self): if self.scheduler.running: self.scheduler.shutdown() logger.info("Scheduler shut down") def _run_generate(self): try: logger.info("[Scheduled] Starting content generation...") result = run_creator() logger.info("[Scheduled] Generation completed: %s", result) created_id = result.get("topic_id") if isinstance(result, dict) else None if created_id: logger.info("[Scheduled] Running compliance review on %s...", created_id) review_result = run_optimizer([created_id]) if review_result.get("ok"): logger.info("[Scheduled] Review completed for %s", created_id) else: logger.warning("[Scheduled] Review failed: %s", review_result.get("error")) except Exception as e: logger.exception("[Scheduled] Generation pipeline failed: %s", e) def _run_optimize(self): try: logger.info("[Scheduled] Starting compliance optimization...") result = run_optimizer() logger.info("[Scheduled] Optimization completed: %s", result) except Exception as e: logger.exception("[Scheduled] Optimization failed: %s", e) def _run_sync(self): try: logger.info("[Scheduled] Starting data sync...") sync_all_topics() logger.info("[Scheduled] Sync completed") except Exception as e: logger.exception("[Scheduled] Sync failed: %s", e) def get_jobs(self): """返回当前所有定时任务的状态""" jobs = [] for job in self.scheduler.get_jobs(): jobs.append({ "id": job.id, "next_run_time": job.next_run_time.isoformat() if job.next_run_time else None, "trigger": str(job.trigger), }) return jobs # 全局单例 scheduler = TaskScheduler()