Files
crawler-plugin/app/amazon/main.py
koko 9146d17625 Align patrol delete module naming
Rename the patrol deletion implementation and browser helper scripts away from the generic product name so the filesystem matches the queued patrol-delete task domain. Update imports, script loading, and tests to use the patrol_delete module path.

Constraint: Existing task class remains ProductTask to avoid broad API churn beyond the requested file and script naming

Rejected: Rename ProductTask class now | would expand the change into a larger API update across tests and task dispatch

Confidence: high

Scope-risk: narrow

Directive: Keep patrol-delete script paths under app/amazon/scripts/patrol_delete when adding new helper scripts

Tested: uv run pytest tests/test_product.py

Not-tested: Remote push before this commit due previous SSH publickey failure
2026-04-27 17:58:40 +08:00

697 lines
30 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 time
import traceback
import requests
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, kill_process
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 ProductTask
from amazon.detail_spider import SpiderTask
from amazon.tool import get_shop_info,show_notification
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 # 线程池执行器
# 在提交新任务前杀掉旧进程(确保环境干净)
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 = {
"product-risk-resolve-run" : "产品风险审批",
"shop-match-run" : "匹配价格",
"price-track-run" : "跟价",
"query-asin-run" : "状态查询",
"patrol-delete-run" : "巡店删除",
"appearance-patent-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 == "delete-brand-run":
# 提交删除品牌任务到线程池
self.log(f"接收到【删除品牌】任务,提交到线程池处理...")
future = self.executor.submit(self._process_task_wrapper, task_data)
futures.append(future)
elif 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):
"""审批任务处理包装器(用于线程池调用)
Args:
task_data: 任务数据
"""
try:
TASK_INFO = {
"产品风险审批" : ApproveTask,
"匹配价格" : MatchTak,
"跟价" : PriceTask,
"状态查询" : StatusTask,
"巡店删除" : ProductTask,
"亚马逊采集" : 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)} 产品风险审批任务处理完成")
except Exception as e:
self.log(f"线程 {id(task_data)} 产品风险审批任务处理异常: {traceback.format_exc()}", "ERROR")
def process_task(self, task_data: Dict[str, Any]):
"""处理单个任务
Args:
task_data: 任务数据
"""
try:
data = task_data.get("data", {})
task_id = data.get("taskId")
items = data.get("items", [])
if not task_id:
self.log("任务ID为空跳过", "WARNING")
return
# 初始化任务状态
self.update_task_status(
task_id,
status="running",
start_time=datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
total_shops=len(items),
processed_shops=0,
total_asins=sum(shop.get("totalRows", 0) for shop in items),
processed_asins=0,
success_count=0,
failed_count=0,
stop_requested=False # 暂停请求标志
)
self.log(f"开始处理任务 {task_id},共 {len(items)} 个店铺")
# 遍历处理每个店铺
for idx, shop_data in enumerate(items, 1):
shop_name = shop_data.get("shopName", "未知店铺")
self.log(f"[{idx}/{len(items)}] 开始处理店铺: {shop_name}")
try:
self.process_shop(shop_data, task_id)
# 更新已处理店铺数
if task_id in runing_task:
runing_task[task_id]["processed_shops"] += 1
except Exception as e:
self.log(f"处理店铺 {shop_name} 失败: {str(e)}", "ERROR")
self.log(traceback.format_exc(), "ERROR")
# 检查任务最终状态
if task_id in runing_task and runing_task[task_id].get("stop_requested", False):
self.update_task_status(task_id, status="stopped")
self.log(f"任务 {task_id} 已被暂停!")
else:
self.update_task_status(task_id, status="completed")
self.log(f"任务 {task_id} 处理完成!")
except Exception as e:
self.log(f"任务处理失败: {traceback.format_exc()}", "ERROR")
if task_id:
self.update_task_status(task_id, status="failed", error=str(e))
def process_shop(self, shop_data: Dict[str, Any], task_id: int):
"""处理单个店铺(包含重试逻辑)
Args:
shop_data: 店铺数据
task_id: 任务ID
"""
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")
return
# 更新当前处理的店铺
self.update_task_status(task_id, current_shop=shop_name)
# 将店铺添加到正在执行中的店铺列表
start_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
runing_shop[shop_name] = start_time
self.log(f"账号:{company_name},店铺 {shop_name} 已添加到执行列表,开始时间: {start_time}")
# 店铺打开重试最多3次
driver = None
browser = None
max_retries = 3
for retry in range(max_retries):
try:
self.log(f"尝试打开店铺 {shop_name} (第 {retry + 1}/{max_retries} 次)")
# 如果不是第一次尝试,先杀进程
# if retry > 0:
# self.log("重试前先杀掉浏览器进程...")
# kill_process("v6")
# kill_process("v5")
# time.sleep(2)
# 创建驱动并打开店铺
user_info = {
**self.user_info,
"company": company_name
}
driver = AmazoneDriver(user_info)
browser = driver.open_shop(shop_name)
if browser and browser != "店铺不存在":
self.log(f"成功打开店铺 {shop_name}")
break
else:
self.log(f"打开店铺失败: {browser}", "WARNING")
driver = None
# 判断是否需要登录
driver.tab.wait.doc_loaded(timeout=30, raise_err=False) # 等待页面加载,避免过早判断登录状态
need_login = driver.need_login()
print("【是否需要登录】:",need_login)
if need_login:
self.log(f"店铺 {shop_name} 需要登录,正在登录...")
# 获取店铺凭证
response = get_shop_info(shop_name)
print("【获取店铺凭证返回】:",response.text)
shop_data = response.json()
if not shop_data:
mes = f"获取店铺凭证失败,响应数据: {shop_data.get('message', '未知错误')}"
self.log(mes, "ERROR")
show_notification(mes, "ERROR")
continue
password = shop_data["data"]["password"]
login_success = driver.login(password)
if login_success:
self.log(f"店铺 {shop_name} 登录成功,正在重新打开店铺...")
browser = driver.open_shop(shop_name)
if browser and browser != "店铺不存在":
self.log(f"成功打开店铺 {shop_name} 登录后")
break
else:
self.log(f"登录后打开店铺失败: {browser}", "WARNING")
driver = None
else:
self.log(f"店铺 {shop_name} 登录失败", "WARNING")
driver = None
except Exception as e:
self.log(f"打开店铺异常: {str(e)}", "ERROR")
driver = None
# 如果还有重试机会,等待后继续
if retry < max_retries - 1:
time.sleep(3)
# 检查是否成功打开
if not driver or not browser or browser == "店铺不存在":
self.log(f"店铺 {shop_name} 打开失败,已重试 {max_retries} 次,跳过该店铺", "ERROR")
return
try:
# 处理每个国家
chunk_index =1
for country_data in countries:
# 检查是否收到暂停请求
if task_id in runing_task and runing_task[task_id].get("stop_requested", False):
self.log(f"检测到任务 {task_id} 的暂停请求,停止处理国家", "WARNING")
break # 跳出循环进入finally关闭店铺
try:
chunk_index = self.process_country(driver, country_data, task_id, result_id, shop_data,chunk_index)
except Exception as e:
country_name = country_data.get("country", "未知")
self.log(f"处理国家 {country_name} 失败: {str(e)}", "ERROR")
self.log(traceback.format_exc(), "ERROR")
finally:
# 关闭店铺
try:
if driver:
self.log(f"关闭店铺 {shop_name}")
driver.close_store()
time.sleep(2)
except Exception as e:
self.log(f"关闭店铺失败: {str(e)}", "WARNING")
# 从正在执行中的店铺列表中移除
if shop_name in runing_shop:
del runing_shop[shop_name]
self.log(f"店铺 {shop_name} 已从执行列表中移除")
def process_country(self, driver: AmazoneDriver, country_data: Dict[str, Any],
task_id: int, result_id: int, shop_data: Dict[str, Any],chunk_index:int):
"""处理单个国家的所有ASIN
Args:
driver: 亚马逊驱动实例
country_data: 国家数据
task_id: 任务ID
result_id: 结果ID
shop_data: 店铺数据
"""
country = country_data.get("country", "未知")
items = country_data.get("items", [])
self.log(f"开始处理国家: {country},共 {len(items)} 个ASIN")
self.update_task_status(task_id, current_country=country)
# 切换国家最多重试3次
max_retries = 3
switch_success = False
for retry in range(max_retries):
try:
self.log(f"尝试切换到国家 {country} (第 {retry + 1}/{max_retries} 次)")
# 如果不是第一次尝试,先刷新页面
if retry > 0:
self.log("重试前刷新页面...")
try:
driver.tab.refresh()
time.sleep(3)
except Exception as e:
self.log(f"刷新页面失败: {str(e)}", "WARNING")
switch_success = driver.SwitchingCountries(country)
if switch_success:
self.log(f"成功切换到国家 {country}")
break
else:
self.log(f"切换到国家 {country} 失败", "WARNING")
except Exception as e:
self.log(f"切换国家 {country} 异常: {str(e)}", "ERROR")
# 如果还有重试机会,等待后继续
if retry < max_retries - 1:
time.sleep(2)
# 如果切换失败回传该国家所有ASIN为失败状态
if not switch_success:
self.log(f"切换到国家 {country} 失败,已重试 {max_retries}将所有ASIN标记为失败", "ERROR")
chunk_index = self._report_all_asins_failed(country, items, task_id, shop_data,chunk_index)
return chunk_index
# 切换到库存管理页面
try:
driver.SwitchPage()
self.log(f"已切换到库存管理页面")
except Exception as e:
self.log(f"切换页面失败: {str(e)}", "ERROR")
# 切换页面失败也回传所有ASIN为失败
chunk_index = self._report_all_asins_failed(country, items, task_id, shop_data,chunk_index)
return chunk_index
# 处理每个ASIN
file_key = shop_data.get("fileKey", "")
source_filename = shop_data.get("sourceFilename", "")
total_rows = shop_data.get("totalRows", 0)
for idx, asin_item in enumerate(items, 1):
# 检查是否收到暂停请求
if task_id in runing_task and runing_task[task_id].get("stop_requested", False):
self.log(f"检测到任务 {task_id} 的暂停请求,停止处理", "WARNING")
return # 返回到上层,会触发店铺关闭
try:
self.log(f"[{idx}/{len(items)}] 处理ASIN: {asin_item.get('asin', '')}")
self.process_asin(driver, asin_item, country, task_id, result_id,
file_key, source_filename, total_rows,chunk_index)
except Exception as e:
asin = asin_item.get("asin", "未知")
self.log(f"处理ASIN {asin} 失败: {str(e)}", "ERROR")
# 继续处理下一个ASIN
chunk_index += 1
return chunk_index
def process_asin(self, driver: AmazoneDriver, asin_item: Dict[str, Any],
country: str, task_id: int, result_id: int, file_key: str,
source_filename: str, total_rows: int,chunk_index:int):
"""处理单个ASIN并回传结果
Args:
driver: 亚马逊驱动实例
asin_item: ASIN数据项
country: 国家名称
task_id: 任务ID
result_id: 结果ID
file_key: 文件KEY
source_filename: 源文件名
total_rows: 总行数
chunk_index: 当前处理索引
"""
asin = asin_item.get("asin", "")
row_index = asin_item.get("rowIndex", 0)
# 更新当前处理的ASIN
self.update_task_status(task_id, current_asin=asin)
status = "失败"
max_retries = 3 # 最多重试3次
for retry in range(max_retries):
try:
self.log(f"处理ASIN {asin} (第 {retry + 1}/{max_retries} 次)")
# 如果不是第一次尝试,先刷新页面
if retry > 0:
self.log("重试前刷新页面...")
try:
driver.tab.refresh()
time.sleep(3)
except Exception as e:
self.log(f"刷新页面失败: {str(e)}", "WARNING")
# 搜索ASIN
sku_ls = driver.search(asin=asin)
self.log(f"搜索到 {len(sku_ls)} 个SKU")
if len(sku_ls) == 0:
status = "查询不到"
self.log(f"ASIN {asin} 未找到商品", "WARNING")
break # 查询不到商品,无需重试
else:
# 删除所有找到的SKU如果任何一个失败则重新开始整个流程
success_count = 0
all_success = True # 标记是否所有SKU都删除成功
total_sku_count = len(sku_ls)
for sku in sku_ls:
try:
suc = driver.del_action(sku)
if suc:
success_count += 1
self.log(f"SKU 删除成功 ({success_count}/{total_sku_count})")
else:
self.log(f"SKU 删除失败", "WARNING")
all_success = False
break # 任何一个失败,退出循环,准备重试整个流程
except Exception as e:
self.log(f"删除SKU异常: {str(e)}", "ERROR")
all_success = False
break # 发生异常,退出循环,准备重试整个流程
# 判断删除结果
if all_success and success_count == total_sku_count:
status = "成功"
self.log(f"ASIN {asin} 所有SKU删除成功 ({success_count}/{total_sku_count})")
# 更新成功计数
if task_id in runing_task:
runing_task[task_id]["success_count"] += 1
break # 全部成功,跳出重试循环
else:
# 有失败的SKU
if retry < max_retries - 1:
self.log(f"有SKU删除失败准备重试整个流程... ({retry + 1}/{max_retries})")
time.sleep(2)
continue # 继续下一次重试
else:
# 所有重试都用完了
if success_count > 0:
status = "部分成功"
self.log(f"ASIN {asin} 部分SKU删除成功 ({success_count}/{total_sku_count})", "WARNING")
if task_id in runing_task:
runing_task[task_id]["success_count"] += 1
else:
status = "失败"
self.log(f"ASIN {asin} 所有SKU删除失败", "ERROR")
if task_id in runing_task:
runing_task[task_id]["failed_count"] += 1
except Exception as e:
status = "删除异常"
self.log(f"处理ASIN {asin} 异常: {str(e)}", "ERROR")
# 如果还有重试机会,继续重试
if retry < max_retries - 1:
self.log(f"发生异常,准备重试... ({retry + 1}/{max_retries})")
time.sleep(2)
else:
# 所有重试都失败了,更新失败计数
if task_id in runing_task:
runing_task[task_id]["failed_count"] += 1
# 更新已处理ASIN计数
if task_id in runing_task:
runing_task[task_id]["processed_asins"] += 1
# 回传结果到API
try:
payload = {
"submissionId": "",
"files": [{
"fileKey": file_key,
"sourceFilename": source_filename,
"chunkIndex": chunk_index,
"chunkTotal": total_rows,
"processedRows": chunk_index,
"totalRows": total_rows,
"currentCountry": country,
"currentAsin": asin,
"countries": [{
"country": country,
"items": [{
"asin": asin,
"status": status
}]
}]
}]
}
self.post_result(task_id, payload)
self.log(f"ASIN {asin} 结果已回传,状态: {status}")
except Exception as e:
self.log(f"回传结果失败: {str(e)}", "ERROR")
def _report_all_asins_failed(self, country: str, items: List[Dict[str, Any]],
task_id: int, shop_data: Dict[str, Any],chunk_index:int):
"""将国家下所有ASIN标记为失败并回传
Args:
country: 国家名称
items: ASIN列表
task_id: 任务ID
shop_data: 店铺数据
"""
self.log(f"开始回传国家 {country} 下的 {len(items)} 个ASIN失败状态")
file_key = shop_data.get("fileKey", "")
source_filename = shop_data.get("sourceFilename", "")
total_rows = shop_data.get("totalRows", 0)
for asin_item in items:
asin = asin_item.get("asin", "")
try:
# 更新失败计数
if task_id in runing_task:
runing_task[task_id]["failed_count"] += 1
runing_task[task_id]["processed_asins"] += 1
# 回传失败状态
payload = {
"submissionId": "",
"files": [{
"fileKey": file_key,
"sourceFilename": source_filename,
"chunkIndex": chunk_index,
"chunkTotal": total_rows,
"processedRows": chunk_index,
"totalRows": total_rows,
"currentCountry": country,
"currentAsin": asin,
"countries": [{
"country": country,
"items": [{
"asin": asin,
"status": "失败"
}]
}]
}]
}
self.post_result(task_id, payload)
chunk_index += 1
self.log(f"ASIN {asin} 失败状态已回传")
except Exception as e:
self.log(f"回传ASIN {asin} 失败状态时出错: {str(e)}", "ERROR")
return chunk_index
def post_result(self, task_id: int, payload: Dict[str, Any]):
"""回传结果到API带重试机制
Args:
task_id: 任务ID
payload: 结果数据
"""
url = f"{DELETE_BRAND_API_BASE}/api/delete-brand/tasks/{task_id}/result"
max_retries = 3 # 最多重试3次
for retry in range(max_retries):
try:
self.log(f"尝试回传结果 (第 {retry + 1}/{max_retries} 次)")
response = requests.post(
url,
json=payload,
headers={"Content-Type": "application/json"},
timeout=30,
verify=False # 忽略SSL证书验证
)
print("【结果提交】:",payload)
print("【结果提交返回】:",response.text)
if response.status_code == 200:
self.log(f"结果回传成功: {url}")
return # 成功后直接返回,不再重试
else:
self.log(f"结果回传失败,状态码: {response.status_code}", "WARNING")
# 如果还有重试机会,继续重试
if retry < max_retries - 1:
self.log(f"准备重试... ({retry + 1}/{max_retries})")
time.sleep(2) # 等待2秒后重试
else:
self.log(f"已达到最大重试次数,结果回传最终失败", "ERROR")
except Exception as e:
self.log(f"调用API异常: {str(e)}", "ERROR")
# 如果还有重试机会,继续重试
if retry < max_retries - 1:
self.log(f"发生异常,准备重试... ({retry + 1}/{max_retries})")
time.sleep(2) # 等待2秒后重试
else:
self.log(f"已达到最大重试次数,结果回传最终失败", "ERROR")
def update_task_status(self, task_id: int, **kwargs):
"""更新任务状态
Args:
task_id: 任务ID
**kwargs: 要更新的字段
"""
if task_id not in runing_task:
runing_task[task_id] = {}
runing_task[task_id].update(kwargs)
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()