new file: amazon/__pycache__/approve.cpython-39.pyc

modified:   amazon/__pycache__/del_brand.cpython-39.pyc
	modified:   amazon/__pycache__/main.cpython-39.pyc
	new file:   amazon/approve.py
	modified:   amazon/del_brand.py
	modified:   amazon/main.py
	modified:   main.py
This commit is contained in:
铭坤
2026-04-07 14:09:42 +08:00
parent 1b02160b20
commit 489e7e8269
7 changed files with 905 additions and 18 deletions

View File

@@ -6,6 +6,7 @@ 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, kill_process
from amazon.approve import ApproveTask
class TaskMonitor:
@@ -26,6 +27,7 @@ class TaskMonitor:
# 在提交新任务前杀掉旧进程(确保环境干净)
kill_process("v6")
kill_process("v5")
def log(self, message: str, level: str = "INFO"):
"""日志输出
@@ -54,15 +56,23 @@ class TaskMonitor:
# 检查任务类型
task_type = task_data.get("type", "")
if task_type != "delete-brand-run":
if task_type == "delete-brand-run":
# 提交删除品牌任务到线程池
self.log(f"接收到删除品牌任务,提交到线程池处理...")
future = self.executor.submit(self._process_task_wrapper, task_data)
futures.append(future)
elif task_type == "product-risk-resolve-run":
# 提交产品风险审批任务到线程池
self.log(f"接收到产品风险审批任务,提交到线程池处理...")
future = self.executor.submit(self._process_approve_task_wrapper, task_data)
futures.append(future)
else:
self.log(f"未知任务类型: {task_type},跳过", "WARNING")
continue
# 提交任务到线程池
self.log(f"接收到删除品牌任务,提交到线程池处理...")
future = self.executor.submit(self._process_task_wrapper, task_data)
futures.append(future)
# 清理已完成的future对象避免内存累积
futures = [f for f in futures if not f.done()]
@@ -86,11 +96,26 @@ class TaskMonitor:
task_data: 任务数据
"""
try:
self.log(f"线程 {id(task_data)} 开始处理任务...")
self.log(f"线程 {id(task_data)} 开始处理删除品牌任务...")
self.process_task(task_data)
self.log(f"线程 {id(task_data)} 任务处理完成")
self.log(f"线程 {id(task_data)} 删除品牌任务处理完成")
except Exception as e:
self.log(f"线程 {id(task_data)} 任务处理异常: {traceback.format_exc()}", "ERROR")
self.log(f"线程 {id(task_data)} 删除品牌任务处理异常: {traceback.format_exc()}", "ERROR")
def _process_approve_task_wrapper(self, task_data: Dict[str, Any]):
"""审批任务处理包装器(用于线程池调用)
Args:
task_data: 任务数据
"""
try:
self.log(f"线程 {id(task_data)} 开始处理产品风险审批任务...")
# 创建ApproveTask实例并处理任务
approve_task = ApproveTask(user_info=self.user_info)
approve_task.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_task(self, task_data: Dict[str, Any]):
@@ -161,6 +186,7 @@ class TaskMonitor:
shop_name = shop_data.get("shopName", "未知店铺")
result_id = shop_data.get("resultId")
countries = shop_data.get("countries", [])
company_name = shop_data.get("companyName", "未知公司")
if not countries:
self.log(f"店铺 {shop_name} 没有国家数据,跳过", "WARNING")
@@ -172,7 +198,7 @@ class TaskMonitor:
# 将店铺添加到正在执行中的店铺列表
start_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
runing_shop[shop_name] = start_time
self.log(f"店铺 {shop_name} 已添加到执行列表,开始时间: {start_time}")
self.log(f"账号:{company_name},店铺 {shop_name} 已添加到执行列表,开始时间: {start_time}")
# 店铺打开重试最多3次
driver = None
@@ -187,10 +213,15 @@ class TaskMonitor:
if retry > 0:
self.log("重试前先杀掉浏览器进程...")
kill_process("v6")
kill_process("v5")
time.sleep(2)
# 创建驱动并打开店铺
driver = AmazoneDriver(self.user_info)
user_info = {
**self.user_info,
"company": company_name
}
driver = AmazoneDriver(user_info)
browser = driver.open_shop(shop_name)
if browser and browser != "店铺不存在":