""" 调度器模块 提供定时任务调度功能,支持cron表达式和间隔执行 """ import time import threading from datetime import datetime, timedelta from typing import Callable, Dict, Any, Optional, List import schedule from croniter import croniter try: from ..core.config_manager import get_config from ..core.logger import get_logger except ImportError: from core.config_manager import get_config from core.logger import get_logger logger = get_logger(__name__) class TaskScheduler: """任务调度器""" def __init__(self): self.config = get_config().scheduler self._running = False self._thread = None self._jobs = {} self._job_stats = {} # 初始化调度 if self.config.enabled: self._load_jobs_from_config() def _load_jobs_from_config(self): """从配置加载定时任务""" for job_config in self.config.jobs: if job_config.get("enabled", False): job_name = job_config["name"] cron_expr = job_config["cron"] # 根据任务名称创建对应的任务函数 if job_name == "daily_crawl": func = self._create_crawl_job() elif job_name == "data_cleanup": func = self._create_cleanup_job() else: logger.warning(f"未知的任务类型: {job_name}") continue self.add_cron_job(job_name, cron_expr, func) def _create_crawl_job(self) -> Callable: """创建爬取任务""" def crawl_job(): try: logger.info("开始执行定时爬取任务") # 导入这里避免循环导入 try: from ..crawler.spider import crawl_announcements from ..filters.filters import filter_from_config from ..storage.postgresql import save_announcements_to_storage from ..storage.md_generator import generate_onu_md from ..notification.wechat import send_announcements_notification except ImportError: from crawler.spider import crawl_announcements from filters.filters import filter_from_config from storage.postgresql import save_announcements_to_storage from storage.md_generator import generate_onu_md from notification.wechat import send_announcements_notification # 执行爬取 crawl_results = crawl_announcements() if not crawl_results: logger.info("定时爬取任务完成:无数据") return # 收集所有公告 all_announcements = [] for result in crawl_results: if result.announcements: all_announcements.extend(result.announcements) if not all_announcements: logger.info("定时爬取任务完成:无新公告") return # 筛选公告 filter_obj = filter_from_config() filtered_announcements, filter_stats = filter_obj.filter(all_announcements) logger.info(f"筛选结果: {len(all_announcements)} -> {len(filtered_announcements)}") # 保存到数据库 saved_count = save_announcements_to_storage(filtered_announcements) # 生成Markdown文件 generate_onu_md(filtered_announcements) # 发送通知 if filtered_announcements: send_announcements_notification(filtered_announcements) logger.info(f"定时爬取任务完成:处理 {len(filtered_announcements)} 条公告,保存 {saved_count} 条") except Exception as e: logger.error(f"定时爬取任务执行失败: {str(e)}") # 发送错误通知 from ..notification.wechat import send_error_alert send_error_alert("定时爬取任务失败", str(e)) return crawl_job def _create_cleanup_job(self) -> Callable: """创建数据清理任务""" def cleanup_job(): try: logger.info("开始执行数据清理任务") try: from ..storage.postgresql import cleanup_storage except ImportError: from storage.postgresql import cleanup_storage # 执行清理 deleted_count = cleanup_storage() logger.info(f"数据清理任务完成:删除 {deleted_count} 条过期数据") # 发送通知(如果删除的数据较多) if deleted_count > 0: try: from ..notification.wechat import send_system_notification except ImportError: from notification.wechat import send_system_notification send_system_notification( "数据清理完成", f"已清理 {deleted_count} 条过期数据" ) except Exception as e: logger.error(f"数据清理任务执行失败: {str(e)}") return cleanup_job def add_cron_job(self, name: str, cron_expr: str, func: Callable) -> bool: """ 添加cron定时任务 Args: name: 任务名称 cron_expr: cron表达式 func: 任务函数 Returns: bool: 添加是否成功 """ try: # 验证cron表达式 croniter(cron_expr) # 添加到schedule schedule.every().day.at("00:00").do(func) # 临时设置,会被替换 # 存储任务信息 self._jobs[name] = { "func": func, "cron": cron_expr, "next_run": None, "last_run": None, "run_count": 0, "error_count": 0 } logger.info(f"添加定时任务: {name} ({cron_expr})") return True except Exception as e: logger.error(f"添加定时任务失败 {name}: {str(e)}") return False def add_interval_job(self, name: str, interval_seconds: int, func: Callable) -> bool: """ 添加间隔执行任务 Args: name: 任务名称 interval_seconds: 执行间隔(秒) func: 任务函数 Returns: bool: 添加是否成功 """ try: schedule.every(interval_seconds).seconds.do(func) self._jobs[name] = { "func": func, "interval": interval_seconds, "next_run": None, "last_run": None, "run_count": 0, "error_count": 0 } logger.info(f"添加间隔任务: {name} ({interval_seconds}秒)") return True except Exception as e: logger.error(f"添加间隔任务失败 {name}: {str(e)}") return False def remove_job(self, name: str) -> bool: """ 移除任务 Args: name: 任务名称 Returns: bool: 移除是否成功 """ if name in self._jobs: # 注意:schedule库没有直接的移除方法 # 这里只是从我们的记录中移除 del self._jobs[name] logger.info(f"移除任务: {name}") return True return False def start(self): """启动调度器""" if self._running: logger.warning("调度器已经在运行中") return if not self.config.enabled: logger.info("调度器已禁用") return self._running = True self._thread = threading.Thread(target=self._run_scheduler, daemon=True) self._thread.start() logger.info("调度器已启动") def stop(self): """停止调度器""" if not self._running: return self._running = False if self._thread and self._thread.is_alive(): self._thread.join(timeout=5) logger.info("调度器已停止") def _run_scheduler(self): """运行调度器主循环""" logger.info("调度器主循环开始") while self._running: try: schedule.run_pending() time.sleep(1) except Exception as e: logger.error(f"调度器运行异常: {str(e)}") time.sleep(5) # 出错后等待5秒再继续 logger.info("调度器主循环结束") def run_once(self, job_name: Optional[str] = None): """ 手动执行任务一次 Args: job_name: 任务名称,如果为None则执行所有任务 """ if job_name: if job_name in self._jobs: job_info = self._jobs[job_name] logger.info(f"手动执行任务: {job_name}") try: job_info["func"]() job_info["run_count"] += 1 job_info["last_run"] = datetime.now() logger.info(f"任务 {job_name} 执行完成") except Exception as e: job_info["error_count"] += 1 logger.error(f"任务 {job_name} 执行失败: {str(e)}") else: logger.error(f"任务不存在: {job_name}") else: # 执行所有任务 for name in list(self._jobs.keys()): self.run_once(name) def get_status(self) -> Dict[str, Any]: """ 获取调度器状态 Returns: Dict[str, Any]: 状态信息 """ jobs_status = {} for name, job_info in self._jobs.items(): jobs_status[name] = { "enabled": True, "run_count": job_info.get("run_count", 0), "error_count": job_info.get("error_count", 0), "last_run": job_info.get("last_run").isoformat() if job_info.get("last_run") else None, "next_run": job_info.get("next_run").isoformat() if job_info.get("next_run") else None } return { "enabled": self.config.enabled, "running": self._running, "timezone": self.config.timezone, "jobs": jobs_status } def is_running(self) -> bool: """检查调度器是否正在运行""" return self._running # 全局调度器实例 _scheduler = None def get_scheduler() -> TaskScheduler: """ 获取调度器实例 Returns: TaskScheduler: 调度器实例 """ global _scheduler if _scheduler is None: _scheduler = TaskScheduler() return _scheduler def start_scheduler(): """启动调度器""" scheduler = get_scheduler() scheduler.start() def stop_scheduler(): """停止调度器""" scheduler = get_scheduler() scheduler.stop() def run_scheduled_jobs(job_name: Optional[str] = None): """ 手动执行定时任务 Args: job_name: 任务名称 """ scheduler = get_scheduler() scheduler.run_once(job_name)