Files
crawler-plugin/app/amazon/match_action.py
铭坤 d8098b0378 modified: app/amazon/__pycache__/del_brand.cpython-39.pyc
modified:   app/amazon/__pycache__/match_action.cpython-39.pyc
	modified:   app/amazon/del_brand.py
	modified:   app/amazon/match_action.py
	modified:   app/assets/delete-brand.js
	modified:   app/new_web_source/delete-brand.html
	new file:   "app/web_source/brand - \345\211\257\346\234\254.html"
2026-04-12 20:52:26 +08:00

717 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 re
import traceback
from DrissionPage._pages.chromium_tab import ChromiumTab
from config import runing_task, runing_shop
from datetime import datetime
from amazon.del_brand import AmamzonBase, kill_process
from amazon.tool import show_notification,get_shop_info
class AmzoneMatchAction(AmamzonBase):
mark_name = "匹配价格"
def SwitchPage(self):
"""
切换至 管理所有库存页面
1、等待 //navigation-favorites-bar[@class="hydrated"] 出现
"""
navigation = self.tab.ele('xpath://navigation-favorites-bar[@class="hydrated"]')
navigation.wait.displayed(raise_err=False)
page_btn = navigation.sr('xpath://internal-fav-bar-links[@data-internal="navigation"]').sr(
'xpath://a[@data-page-id="ezdpc-gui-inventory-mons"]')
page_btn.wait.displayed(raise_err=False)
page_btn.click(timeout=5)
self.tab.wait.doc_loaded()
# 等待搜索框出现
search_region = self.tab.ele('xpath://div[@id="searchBoxContainer"]//kat-input-group')
search_region.wait.displayed(raise_err=False)
def search(self,filter_type="ApprovalRequired"):
sku_ls = []
load_ele = self.tab.eles("xpath://div[contains(@class,'Loader-module__loader')]",timeout=5)
if len(load_ele) > 0:
load_ele[0].wait.deleted(timeout=3, raise_err=False)
time.sleep(0.5)
drop_down = self.tab.ele('xpath://div[contains(@class,"VolusListingStatusDropDown-module__verticalContainer")]//kat-dropdown')
drop_down.wait.displayed(raise_err=False)
drop_down.wait.enabled(raise_err=False)
time.sleep(0.6)
drop_down.click()
# //kat-option[@value="SearchSuppressed"]
xp = f'xpath://kat-option[@value="{filter_type}"]'
print(f"{self.mark_name}】正在寻找筛选条件 {filter_type}xpath: {xp}")
approval_required = self.tab.eles(xp,timeout=5)
if len(approval_required) == 0:
print(f"{self.mark_name}】没有需要{filter_type}选项】没有需要{filter_type}的商品了")
return sku_ls # "没有需要审批的商品了"
else:
approval_required = approval_required[0]
approval_required.wait.displayed(raise_err=False)
approval_required.click()
approval_required_text = approval_required.text
print(f"{self.mark_name}】已选择筛选条件: {approval_required_text}")
count = re.findall(r'\d+', approval_required_text)
if count:
count = int(count[0])
print(f"{self.mark_name}】待审批的商品数量: {count}")
if count <= 0:
print(f"{self.mark_name}】没有需要{filter_type}的商品了")
return sku_ls #"没有需要审批的商品了"
for _ in range(3):
# 等待加载完成
load_ele = self.tab.eles("xpath://div[contains(@class,'Loader-module__loader')]")
if len(load_ele) > 0:
load_ele[0].wait.deleted(timeout=3, raise_err=False)
time.sleep(0.5)
sku_ls = self.tab.eles("xpath://div[@data-sku]",timeout=3)
if len(sku_ls) > 0:
break
approval_required.click()
return sku_ls
def wait_loaded(self):
# 等待加载完成
try:
load_ele = self.tab.ele(
'xpath://kat-panel[@data-testid="kat-panel-ActionPanelContent"]//div[contains(@class,"Loader-module__loader")]|//kat-panel[@data-testid="kat-panel-ActionPanelContent"]/div[@data-f1-component]//div[contains(@class,"==")]/div[contains(@class,"==")]/span',
timeout=3)
load_ele.wait.deleted(timeout=5, raise_err=False)
except Exception as e:
print(f"{self.mark_name}】等待加载中消失出错", e)
def run_page_action(self):
print(f"{self.mark_name}】开始执行")
num = 0
retry_num = 0
already_asin = set()
while retry_num < 3: # 最多重试3次
# if num > 3: #测试
# return
# 等待加载完成
try:
load_ele = self.tab.eles("xpath://div[contains(@class,'Loader-module__loader')]",timeout=5)
if len(load_ele) > 0:
load_ele[0].wait.deleted(timeout=3, raise_err=False)
time.sleep(0.5)
# 获取当前页码
try:
page_pamel = self.tab.eles('xpath://kat-pagination',timeout=5)
if len(page_pamel) > 0:
current_page = page_pamel[0].sr('xpath:.//ul[@class="pages"]//li[@aria-current="true"]').text
# 总页数
total_page = page_pamel[0].sr.eles('xpath:.//ul[@class="pages"]//span[@class="page__inner"][last()]')
if len(total_page) > 0:
total_page = total_page[-1].text
else:
total_page = 0
print(f"{self.mark_name}】当前页码: {current_page} / 总页数: {total_page}")
except Exception as e:
print(f"{self.mark_name}】获取页码失败", e)
sku_ls = self.tab.eles("xpath://div[@data-sku]",timeout=10)
print(f"{self.mark_name}】获取到 {len(sku_ls)}")
# for sku_ele in sku_ls[0:2]:
for sku_ele in sku_ls:
# solve_problem = sku_ele.eles('xpath:.//kat-link[@label="解决商品信息问题"]')
asin = sku_ele.ele('xpath:.//div[contains(@class,"JanusSplitBox-module__container")]//div[contains(@class,"JanusSplitBox-module__panel--") and contains(string(.),"ASIN")]/..//div[last()]',timeout=10).text
print(f"{self.mark_name}】ASIN {asin} 找到....")
if asin in already_asin:
print(f"{self.mark_name}{asin} 已经处理过了,跳过")
continue
price_match = sku_ele.eles('xpath:.//div[@data-test-id="FeaturedOfferPrice"]//a[text()="匹配"]',timeout=5)
if len(price_match) == 0:
print(f"{self.mark_name}{asin},没有推荐价格匹配按钮")
yield (asin,"无需处理")
already_asin.add(asin)
continue
try:
price_match[0].click()
save_all_btn = self.tab.ele('xpath://kat-button[@label="保存所有"]',timeout=10)
save_all_btn.wait.displayed(timeout=10, raise_err=False)
save_all_btn.click()
save_all_btn.wait.deleted(timeout=10, raise_err=False)
yield (asin,"处理完成")
already_asin.add(asin)
except Exception as e:
print(f"{asin} 处理失败", e)
yield (asin,"处理失败")
already_asin.add(asin)
continue
# 判断是否存在需要翻页的情况
page_pamel = self.tab.eles('xpath://kat-pagination',timeout=5)
if len(page_pamel) == 0:
break
next_page_btn = page_pamel[0].sr('xpath:.//span[@part="pagination-nav-right"]')
class_str = next_page_btn.attr('class')
if "end" in class_str:
break
next_page_btn.click()
num += 1
print(f"{self.mark_name}】【程序计算】正在翻页,已翻 {num} 页...")
already_asin = set()
except Exception as e:
print(f"{self.mark_name}】处理匹配操作异常", e)
traceback.print_exc()
retry_num += 1
self.tab.refresh()
self.tab.wait.doc_loaded(raise_err=False,timeout=120)
class MatchTak:
country_info = {
"DE": "德国",
"FR": "法国",
"ES": "西班牙",
"IT": "意大利",
"UK": "英国"
}
def __init__(self, user_info: dict = None):
"""初始化审批任务处理器
Args:
user_info: 用户信息字典,包含 company, username, password
"""
self.user_info = user_info or {}
self.running = True
def log(self, message: str, level: str = "INFO"):
"""日志输出
Args:
message: 日志消息
level: 日志级别
"""
from datetime import datetime
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
if level == "ERROR":
show_notification(message, "error")
print(f"[{timestamp}] [MatchTak] [{level}] {message}")
def process_task(self, task_data: dict):
"""处理审批任务主入口
Args:
task_data: 任务数据
"""
try:
data = task_data.get("data", {})
task_id = data.get("taskId")
items = data.get("items", [])
country_codes = data.get("country_codes", [])
risk_listing_filter = data.get("risk_listing_filter", "All")
user_id = data.get("user_id")
stage_index = data.get("stage_index")
final_stage = bool(data.get("final_stage", True))
if not task_id:
self.log("任务ID为空跳过", "WARNING")
return
if not items:
self.log("店铺列表为空,跳过", "WARNING")
return
if not country_codes:
self.log("国家列表为空,跳过", "WARNING")
return
self.log(f"开始处理审批任务 {task_id},共 {len(items)} 个店铺,{len(country_codes)} 个国家")
from config import runing_task
runing_task[task_id] = {
"status": "running",
"start_time": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
"total_shops": len(items),
"processed_shops": 0,
"total_countries": len(country_codes) * len(items),
"processed_countries": 0,
"total_asins": 0,
"processed_asins": 0,
"success_count": 0,
"failed_count": 0,
"stop_requested": False
}
# 遍历处理每个店铺
for idx, shop_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")
runing_task[task_id]["status"] = "stopped"
return
shop_name = shop_item.get("shopName", "未知店铺")
self.log(f"[{idx}/{len(items)}] 开始处理店铺: {shop_name}")
show_notification(f"开始处理店铺: {shop_name}", "info")
try:
self.process_shop(shop_item, country_codes, task_id, risk_listing_filter, user_id, stage_index, final_stage)
# self.process_shop(shop_item, country_codes, task_id,risk_listing_filter)
# 更新已处理店铺数
if task_id in runing_task:
runing_task[task_id]["processed_shops"] += 1
except Exception as e:
import traceback
self.log(f"处理店铺 {shop_name} 失败: {str(e)}", "ERROR")
self.log(traceback.format_exc(), "ERROR")
# 更新任务状态
if task_id in runing_task:
if runing_task[task_id].get("stop_requested", False):
runing_task[task_id]["status"] = "stopped"
self.log(f"任务 {task_id} 已被暂停!")
else:
runing_task[task_id]["status"] = "completed"
self.log(f"任务 {task_id} 处理完成!")
except Exception as e:
import traceback
self.log(f"任务处理失败: {traceback.format_exc()}", "ERROR")
if task_id:
from config import runing_task
if task_id in runing_task:
runing_task[task_id]["status"] = "failed"
runing_task[task_id]["error"] = str(e)
# def process_shop(self, shop_item: dict, country_codes: list, task_id: int, risk_listing_filter: str):
def process_shop(self, shop_item: dict, country_codes: list, task_id: int, risk_listing_filter: str, user_id=None, stage_index=None, final_stage: bool = True):
"""处理单个店铺
Args:
shop_item: 店铺信息
country_codes: 国家代码列表
task_id: 任务ID
risk_listing_filter: 风险商品筛选条件
"""
shop_name = shop_item.get("shopName", "未知店铺")
company_name = shop_item.get("companyName", "")
if not company_name:
self.log(f"店铺 {shop_name} 的公司名称为空,跳过", "WARNING")
return
if task_id in runing_task:
runing_task[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"店铺 {shop_name} 已添加到执行列表,账号: {company_name},开始时间: {start_time}")
# 店铺打开重试最多3次
driver = None
max_retries = 3
error_info = ""
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 = AmzoneMatchAction(user_info)
browser = driver.open_shop(shop_name)
if browser and browser != "店铺不存在":
self.log(f"成功打开店铺 {shop_name}")
else:
self.log(f"打开店铺失败: {browser}", "WARNING")
driver = None
continue
# 判断是否需要登录
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
else:
break
except Exception as e:
import traceback
self.log(f"打开店铺异常: {traceback.format_exc()}", "INFO")
driver = None
error_info = str(e)
time.sleep(10)
# 如果还有重试机会,等待后继续
if retry < max_retries - 1:
time.sleep(3)
# 检查是否成功打开
if not driver or not browser or browser == "店铺不存在":
error_msg = f"店铺 {shop_name} 打开失败,已重试 {max_retries} 次,跳过该店铺,{error_info}"
self.log(error_msg, "ERROR")
# 从执行列表中移除
if shop_name in runing_shop:
del runing_shop[shop_name]
return
try:
# 处理每个国家
for country_code in country_codes:
# 检查是否收到暂停请求
if task_id in runing_task and runing_task[task_id].get("stop_requested", False):
self.log(f"检测到任务 {task_id} 的暂停请求,停止处理国家", "WARNING")
break
try:
self.process_country(driver, country_code, task_id, shop_name,risk_listing_filter)
except Exception as e:
import traceback
self.log(f"处理国家 {country_code} 失败: {str(e)}", "ERROR")
self.log(traceback.format_exc(), "ERROR")
# 最后回传,标记完成
try:
# self.post_result(task_id, shop_name, country_code, "", "", is_done=True)
if final_stage:
self.post_result(task_id, shop_name, country_code, "", "", is_done=True)
else:
self.post_stage_finished(task_id, user_id, stage_index)
except Exception as e:
self.log(f"回传结果失败: {str(e)}", "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: AmzoneMatchAction, country_code: str, task_id: int, shop_name: str, risk_listing_filter: str):
"""处理单个国家的审批任务
Args:
driver: AmzoneApprove驱动实例
country_code: 国家代码(如 UK, DE, FR 等)
task_id: 任务ID
shop_name: 店铺名称
risk_listing_filter: 风险商品筛选条件
"""
from config import runing_task
# 转换国家代码为中文名称
country_name = self.country_info.get(country_code, country_code)
info_mes = f"开始处理国家: {country_name} ({country_code})"
self.log(info_mes)
show_notification(info_mes, "info")
# 更新当前处理的国家
if task_id in runing_task:
runing_task[task_id]["current_country"] = country_name
# 切换国家最多重试3次
max_retries = 3
switch_success = False
for retry in range(max_retries):
try:
self.log(f"尝试切换到国家 {country_name} (第 {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_name)
if switch_success:
self.log(f"成功切换到国家 {country_name}")
break
else:
self.log(f"切换到国家 {country_name} 失败", "WARNING")
except Exception as e:
import traceback
self.log(f"切换国家 {country_name} 异常: {str(e)}", "ERROR")
self.log(traceback.format_exc(), "ERROR")
# 如果还有重试机会,等待后继续
if retry < max_retries - 1:
time.sleep(2)
# 如果切换失败,直接返回
if not switch_success:
error_message = f"切换到国家 {country_name} 失败,已重试 {max_retries} 次,跳过该国家"
self.log(error_message, "ERROR")
if task_id in runing_task:
runing_task[task_id]["processed_countries"] += 1
return
# 切换到库存管理页面
try:
driver.SwitchPage()
self.log(f"已切换到库存管理页面")
except Exception as e:
import traceback
self.log(f"切换页面失败: {str(e)}", "ERROR")
self.log(traceback.format_exc(), "ERROR")
if task_id in runing_task:
runing_task[task_id]["processed_countries"] += 1
return
# 搜索需要审批的商品最多重试3次
sku_ls = []
for retry in range(max_retries):
try:
self.log(f"尝试搜索匹配操作商品 (第 {retry + 1}/{max_retries} 次)")
sku_ls = driver.search(filter_type=risk_listing_filter)
break
except Exception as e:
self.log(f"搜索商品异常: {str(e)}", "ERROR")
if retry < max_retries - 1:
try:
driver.tab.refresh()
time.sleep(3)
except Exception as refresh_error:
self.log(f"刷新页面失败: {str(refresh_error)}", "WARNING")
# 如果没有需要审批的商品,直接返回
if len(sku_ls) == 0:
self.log(f"国家 {country_name} 没有搜索出的商品")
if task_id in runing_task:
runing_task[task_id]["processed_countries"] += 1
return
self.log(f"国家 {country_name} 搜索出 {len(sku_ls)} 商品,开始处理...")
# 处理所有需要审批的商品通过yield获取结果
try:
for asin, status in driver.run_page_action():
# 检查是否收到暂停请求
if task_id in runing_task and runing_task[task_id].get("stop_requested", False):
self.log(f"检测到任务 {task_id} 的暂停请求停止处理ASIN", "WARNING")
break
self.log(f"ASIN {asin} 处理结果: {status}")
# 更新任务状态
if task_id in runing_task:
runing_task[task_id]["current_asin"] = asin
runing_task[task_id]["processed_asins"] += 1
runing_task[task_id]["failed_count"] += 1
# 回传结果到API
try:
self.post_result(task_id, shop_name, country_code, asin, status)
except Exception as e:
self.log(f"回传结果失败: {str(e)}", "ERROR")
except Exception as e:
import traceback
self.log(f"处理审批商品异常: {str(e)}", "ERROR")
self.log(traceback.format_exc(), "ERROR")
# 更新已处理国家数
if task_id in runing_task:
runing_task[task_id]["processed_countries"] += 1
self.log(f"国家 {country_name} 处理完成")
def post_stage_finished(self, task_id: int, user_id, stage_index):
import requests
from config import DELETE_BRAND_API_BASE
if user_id in (None, "", 0):
raise ValueError("user_id is required for stage completion callback")
if stage_index is None:
raise ValueError("stage_index is required for stage completion callback")
url = f"{DELETE_BRAND_API_BASE}/api/shop-match/tasks/{task_id}/stage-finished"
payload = {"stage_index": stage_index}
params = {"user_id": user_id}
max_retries = 3
for retry in range(max_retries):
try:
self.log(f"Attempting stage completion callback ({retry + 1}/{max_retries})")
response = requests.post(
url,
params=params,
json=payload,
headers={"Content-Type": "application/json"},
timeout=30,
verify=False,
)
self.log(f"Stage completion callback response: {response.text}")
data = response.json() if response.text else {}
if response.status_code == 200 and isinstance(data, dict) and data.get("success"):
self.log(f"Stage completion callback succeeded: task={task_id}, stage={stage_index}")
return
self.log(f"Stage completion callback failed, status={response.status_code}", "WARNING")
except Exception as e:
self.log(f"Stage completion callback exception: {str(e)}", "ERROR")
if retry < max_retries - 1:
time.sleep(2)
raise RuntimeError(f"Stage completion callback failed after retries: task={task_id}, stage={stage_index}")
def post_result(self, task_id: int, shop_name: str, country_code: str, asin: str, status: str,is_done: bool = False):
"""回传处理结果到API
Args:
task_id: 任务ID
shop_name: 店铺名称
country_code: 国家代码
asin: ASIN
status: 处理状态
"""
import requests
from config import DELETE_BRAND_API_BASE
url = f"{DELETE_BRAND_API_BASE}/api/shop-match/tasks/{task_id}/result"
payload = {
"shops": [
{
"error": "",
"countries": {
country_code: [
{
"asin": asin,
"status": status,
"done": is_done
}
]
},
"shopName": shop_name
}
]
}
max_retries = 3
for retry in range(max_retries):
try:
print("================【匹配】=====================")
self.log(f"尝试回传结果 (第 {retry + 1}/{max_retries} 次)")
self.log(f"回传URL: {url}")
self.log(f"回传数据: {payload}")
response = requests.post(
url,
json=payload,
headers={"Content-Type": "application/json"},
timeout=30,
verify=False
)
self.log(f"回传结果: {response.text}")
data = response.json() if response.text else {}
if response.status_code == 200 and isinstance(data, dict) and data.get("success"):
self.log(f"结果回传成功: {asin} - {status}")
return
else:
self.log(f"结果回传失败,状态码: {response.status_code}", "WARNING")
print("=====================================")
except Exception as e:
self.log(f"调用API异常: {str(e)}", "ERROR")
print("=====================================")
# 如果还有重试机会,等待后继续
if retry < max_retries - 1:
time.sleep(2)
self.log(f"已达到最大重试次数,结果回传最终失败", "ERROR")
raise RuntimeError("已达到最大重试次数,结果回传最终失败")
if __name__ == "__main__":
user_info = {
"company": "rongchuang123",
"username": "自动化_Robot",
"password": "#20zsg25"
}
shop_name = "魏振峰"
country = "德国"
kill_process('v6')
driver = AmzoneMatchAction(user_info)
browser = driver.open_shop(shop_name)
sw_suc = driver.SwitchingCountries(country)
driver.SwitchPage()
risk_listing_filter = "All"
for _ in range(3):
try:
sku_ls = driver.search(filter_type=risk_listing_filter)
break
except Exception as e:
print(e)
driver.tab.refresh()
if len(sku_ls) > 0:
print("有数据,开始操作")
for asin, status in driver.run_page_action():
print(f"ASIN {asin} 的处理结果: {status}")
print("已完成操作")