Files
GX-gp-notify/gx_gp_monitor/scheduler/scheduler.py
T
2026-01-07 17:37:09 +08:00

373 lines
11 KiB
Python

"""
调度器模块
提供定时任务调度功能,支持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)