初始化参数数据自动化管理系统
功能: - 文章内容库管理 - 待处理产品列表管理 - 自动处理流程 - 智能搜索和数据提取 - ParamHub API集成 - 定时任务调度 部署端口: 16043
This commit is contained in:
@@ -0,0 +1,127 @@
|
||||
"""
|
||||
定时任务调度器
|
||||
"""
|
||||
from apscheduler.schedulers.background import BackgroundScheduler
|
||||
from apscheduler.triggers.interval import IntervalTrigger
|
||||
from datetime import datetime
|
||||
from models.database import db
|
||||
from services.process_service import process_service
|
||||
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():
|
||||
"""
|
||||
自动处理产品的定时任务
|
||||
"""
|
||||
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:
|
||||
# 检查是否正在处理
|
||||
processing = db.get_processing_products()
|
||||
if any(p['product_name'] == product['product_name'] for p in processing):
|
||||
logger.warning(f"产品 {product['product_name']} 正在处理中,跳过")
|
||||
continue
|
||||
|
||||
# 添加到处理中列表
|
||||
db.start_processing(
|
||||
product_name=product['product_name'],
|
||||
category=product.get('category'),
|
||||
subcategory=product.get('subcategory')
|
||||
)
|
||||
|
||||
# 执行处理
|
||||
result = process_service.process_product(product)
|
||||
|
||||
# 从待处理列表移除
|
||||
db.remove_pending_product(product['product_name'])
|
||||
|
||||
logger.info(f"产品 {product['product_name']} 处理完成: {result['message']}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"处理产品 {product['product_name']} 时出错: {str(e)}")
|
||||
db.finish_processing(product['product_name'])
|
||||
finally:
|
||||
# 确保从处理中列表移除
|
||||
db.finish_processing(product['product_name'])
|
||||
|
||||
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'
|
||||
)
|
||||
Reference in New Issue
Block a user