import traceback import requests import threading from datetime import datetime from typing import Dict, Any, List from concurrent.futures import ThreadPoolExecutor, as_completed from config import JSON_TASK_QUEUE, runing_task, runing_shop, DELETE_BRAND_API_BASE, ZN_COMPANY, ZN_USERNAME, ZN_PASSWORD from amazon.del_brand import AmazoneDriver, DelbramdTask from amazon.approve import ApproveTask from amazon.match_action import MatchTak from amazon.price_match import PriceTask from amazon.asin_status import StatusTask from amazon.patrol_delete import PatrolDeleteTask from amazon.detail_spider import SpiderTask from amazon.amazon_base import kill_process class TaskMonitor: """任务监控器:负责监控队列并执行品牌删除任务""" def __init__(self): """初始化任务监控器""" self.running = True self.user_info = { "company": ZN_COMPANY, "username": ZN_USERNAME, "password": ZN_PASSWORD } self.chunk_index = 1 # 当前处理的分块索引 self.max_workers = 5 # 线程池最大线程数 self.executor = None # 线程池执行器 self.serial_task_locks = { "product-risk-resolve-run": threading.Lock(), } # 在提交新任务前杀掉旧进程(确保环境干净) kill_process("v6") kill_process("v5") def log(self, message: str, level: str = "INFO"): """日志输出 Args: message: 日志消息 level: 日志级别 """ timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") print(f"[{timestamp}] [{level}] {message}") def start(self): """启动任务监控(使用线程池处理任务)""" self.log(f"任务监控器启动,开始监听队列... (线程池大小: {self.max_workers})") # 创建线程池 self.executor = ThreadPoolExecutor(max_workers=self.max_workers) futures = [] # 保存所有提交的任务Future对象 task_type_info = { "delete-brand-run" : "删除品牌", "product-risk-resolve-run" : "产品风险审批", "shop-match-run" : "匹配价格", "price-track-run" : "跟价", "query-asin-run" : "状态查询", "patrol-delete-run" : "巡店删除", "appearance-patent-run" : "亚马逊采集", "similar-asin-run" : "亚马逊采集", } try: while self.running: try: # 超时后继续循环检查self.running状态,避免卡死 task_data = JSON_TASK_QUEUE.get(block=True, timeout=30) # 检查任务类型 task_type = task_data.get("type", "") if task_type in task_type_info: # 提交产品风险审批任务到线程池 self.log(f"接收到任务,提交到线程池处理...") future = self.executor.submit(self._process_approve_task_wrapper, task_data, task_type_info[task_type]) futures.append(future) else: self.log(f"未知任务类型: {task_type},跳过", "WARNING") continue # 清理已完成的future对象,避免内存累积 futures = [f for f in futures if not f.done()] except Exception as e: # 静默处理队列为空的超时,只记录真正的异常 if str(e) and "Empty" not in str(e): self.log(f"任务监控异常: {str(e)}", "ERROR") # 队列为空时不需要额外等待,直接继续循环 finally: # 关闭监控时,等待所有任务完成 self.log("正在关闭任务监控器,等待所有任务完成...") if self.executor: self.executor.shutdown(wait=True) self.log("所有任务已完成,监控器已关闭") def _process_task_wrapper(self, task_data: Dict[str, Any]): """任务处理包装器(用于线程池调用) Args: task_data: 任务数据 """ try: self.log(f"线程 {id(task_data)} 开始处理任务...") self.process_task(task_data) self.log(f"线程 {id(task_data)} 任务处理完成") except Exception as e: self.log(f"线程 {id(task_data)} 删除品牌任务处理异常: {traceback.format_exc()}", "ERROR") def _process_approve_task_wrapper(self, task_data: Dict[str, Any],TASK_TYPE:str): task_type = task_data.get("type", "") task_lock = self.serial_task_locks.get(task_type) if task_lock is not None: self.log(f"线程 {id(task_data)} 等待 {task_type} 串行锁...") with task_lock: self.log(f"线程 {id(task_data)} 获取 {task_type} 串行锁,开始处理...") return self._process_approve_task_unlocked(task_data, TASK_TYPE) return self._process_approve_task_unlocked(task_data, TASK_TYPE) def _process_approve_task_unlocked(self, task_data: Dict[str, Any],TASK_TYPE:str): """审批任务处理包装器(用于线程池调用) Args: task_data: 任务数据 """ try: TASK_INFO = { "删除品牌": DelbramdTask , "产品风险审批" : ApproveTask, "匹配价格" : MatchTak, "跟价" : PriceTask, "状态查询" : StatusTask, "巡店删除" : PatrolDeleteTask, "亚马逊采集" : SpiderTask, } self.log(f"线程 {id(task_data)} 开始处理产品风险审批任务...") # 创建ApproveTask实例并处理任务 TASK_CLS = TASK_INFO[TASK_TYPE] # 根据任务类型选择处理类,默认为ApproveTask approve_task = TASK_CLS(user_info=self.user_info) approve_task.process_task(task_data) self.log(f"线程 {id(task_data)} {TASK_TYPE} 任务处理完成") except Exception as e: self.log(f"线程 {id(task_data)} {TASK_TYPE} 任务处理异常: {traceback.format_exc()}", "ERROR") def stop(self): """停止监控""" self.running = False self.log("任务监控器已停止") def main(): """主函数:启动任务监控""" # 创建并启动任务监控器 monitor = TaskMonitor() try: monitor.start() except KeyboardInterrupt: monitor.log("接收到中断信号,正在停止...") monitor.stop() except Exception: monitor.log(f"监控器异常退出: {traceback.format_exc()}", "ERROR") if __name__ == "__main__": main()