""" 定时任务调度器 """ from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.interval import IntervalTrigger from datetime import datetime from models.database import db from config import Config import logging # 配置日志 logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s' ) logger = logging.getLogger('scheduler') class TaskScheduler: def __init__(self): self.scheduler = BackgroundScheduler() self.running = False def start(self): """启动调度器""" if not self.running: self.scheduler.start() self.running = True logger.info("定时任务调度器已启动") def stop(self): """停止调度器""" if self.running: self.scheduler.shutdown() self.running = False logger.info("定时任务调度器已停止") def add_interval_job(self, func, interval_seconds, job_id, **kwargs): """添加定时任务""" self.scheduler.add_job( func, trigger=IntervalTrigger(seconds=interval_seconds), id=job_id, replace_existing=True, **kwargs ) logger.info(f"添加定时任务: {job_id}, 间隔: {interval_seconds}秒") def remove_job(self, job_id): """移除定时任务""" try: self.scheduler.remove_job(job_id) logger.info(f"移除定时任务: {job_id}") except Exception as e: logger.error(f"移除任务失败: {job_id}, 错误: {str(e)}") def get_jobs(self): """获取所有任务""" return self.scheduler.get_jobs() # 创建全局调度器实例 task_scheduler = TaskScheduler() def auto_process_task(): """ 自动处理产品的定时任务(大模型驱动) 通过 process_monitor 启动异步处理会话,避免与手动处理并发冲突 """ from services.process_monitor import process_monitor try: # 检查是否启用自动处理 enabled = db.get_system_config('auto_process_enabled', 'true') if enabled.lower() != 'true': logger.info("自动处理已禁用,跳过本次执行") return # 获取批量处理数量 batch_size = int(db.get_system_config('batch_size', '5')) # 获取待处理产品 products = db.get_pending_products(limit=batch_size) if not products: logger.info("没有待处理的产品") return logger.info(f"开始处理 {len(products)} 个产品") for product in products: try: # 启动大模型处理会话(内部有防重,已有活跃会话会自动跳过) session_id, started = process_monitor.start_process( product_name=product['product_name'], category=product.get('category'), subcategory=product.get('subcategory') ) if started: logger.info(f"产品 {product['product_name']} 处理会话已启动: {session_id}") else: logger.info(f"产品 {product['product_name']} 已有处理会话: {session_id}") except Exception as e: logger.error(f"启动产品 {product['product_name']} 处理会话时出错: {str(e)}") except Exception as e: logger.error(f"自动处理任务执行失败: {str(e)}") def setup_auto_process_job(): """设置自动处理定时任务""" interval = int(db.get_system_config('process_interval', str(Config.PROCESS_INTERVAL))) task_scheduler.add_interval_job( auto_process_task, interval_seconds=interval, job_id='auto_process' )