Files
crawler-plugin/app/amazon/main.py
铭坤 3fbbf87c10 modified: .env
modified:   amazon/approve.py
	modified:   amazon/chrome_base.py
	modified:   amazon/detail_spider.py
	modified:   amazon/main.py
	modified:   amazon/price_match.py
	modified:   amazon/similar_asin.py
	modified:   amazon/tool.py
	new file:   amazon/user_data/chrome_data/BrowserMetrics-spare.pma
	new file:   amazon/user_data/chrome_data/BrowserMetrics/BrowserMetrics-69FAEF07-558.pma
2026-05-06 22:30:44 +08:00

176 lines
7.2 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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.similar_asin import SimilarAsinTask
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" : "相似ASIN"
}
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,
"相似ASIN" : SimilarAsinTask
}
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()