This commit is contained in:
super
2026-04-23 15:25:49 +08:00
8 changed files with 544 additions and 378 deletions

View File

@@ -923,46 +923,18 @@ class ApproveTask:
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):
"""处理单个店铺
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
def open_shop(self,max_retries,company_name,shop_name,iskill=False):
error_info = ""
driver = None
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)
if iskill:
self.log("重试前先杀掉浏览器进程...")
kill_process("v6")
kill_process("v5")
time.sleep(2)
# 组装用户信息并创建驱动
user_info = {
@@ -974,7 +946,6 @@ class ApproveTask:
if browser and browser != "店铺不存在":
self.log(f"成功打开店铺 {shop_name}")
# break
else:
self.log(f"打开店铺失败: {browser}", "WARNING")
driver = None
@@ -1012,6 +983,7 @@ class ApproveTask:
driver = None
else:
break
except Exception as e:
import traceback
self.log(f"打开店铺异常: {traceback.format_exc()}", "INFO")
@@ -1030,9 +1002,40 @@ class ApproveTask:
# 从执行列表中移除
if shop_name in runing_shop:
del runing_shop[shop_name]
return driver
return driver
def process_shop(self, shop_item: dict, country_codes: list, task_id: int, risk_listing_filter: str):
"""处理单个店铺
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
try:
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
iskill = False
# 处理每个国家
for country_code in country_codes:
# 检查是否收到暂停请求
@@ -1040,26 +1043,32 @@ class ApproveTask:
self.log(f"检测到任务 {task_id} 的暂停请求,停止处理国家", "WARNING")
break
# 打开店铺
driver = self.open_shop(max_retries=max_retries, company_name=company_name,
shop_name=shop_name, iskill=iskill)
if driver is None:
self.log(f"任务 {task_id} 启动店铺失败,结束任务", "ERROR")
for country_code in country_codes:
self.post_result(task_id, shop_name, country_code, "", "", is_done=True)
return
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")
if "与页面的连接已断开" in e:
iskill = True
# 更新已处理国家数
if task_id in runing_task:
runing_task[task_id]["processed_countries"] += 1
# 最后回传,标记完成
try:
self.post_result(task_id, shop_name, country_code, "", "", is_done=True)
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:
@@ -1176,7 +1185,6 @@ class ApproveTask:
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):
@@ -1204,14 +1212,6 @@ class ApproveTask:
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} 处理完成")

View File

@@ -169,7 +169,7 @@ class ZiniaoDriver:
"cookieTypeLoad": 0,
"cookieTypeSave": cookieTypeSave,
"runMode": "1",
"isLoadUserPlugin": False,
"isLoadUserPlugin": True,
"pluginIdType": 1,
"privacyMode": isprivacy
}

View File

@@ -309,48 +309,18 @@ class MatchTak:
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,limit:str=None):
"""处理单个店铺
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
def open_shop(self,max_retries,company_name,shop_name,iskill=False):
error_info = ""
driver = None
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)
if iskill:
self.log("重试前先杀掉浏览器进程...")
kill_process("v6")
kill_process("v5")
time.sleep(2)
# 组装用户信息并创建驱动
user_info = {
@@ -418,9 +388,41 @@ class MatchTak:
# 从执行列表中移除
if shop_name in runing_shop:
del runing_shop[shop_name]
return driver
return driver
# 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,limit:str=None):
"""处理单个店铺
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
try:
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
iskill = False
# 处理每个国家
for country_code in country_codes:
# 检查是否收到暂停请求
@@ -428,12 +430,28 @@ class MatchTak:
self.log(f"检测到任务 {task_id} 的暂停请求,停止处理国家", "WARNING")
break
# 打开店铺
driver = self.open_shop(max_retries=max_retries, company_name=company_name,
shop_name=shop_name, iskill=iskill)
if driver is None:
self.log(f"任务 {task_id} 启动店铺失败,结束任务", "ERROR")
for country_code in country_codes:
self.post_result(task_id, shop_name, country_code, "", "", is_done=True)
return
try:
self.process_country(driver, country_code, task_id, shop_name,risk_listing_filter,limit)
except Exception as e:
import traceback
self.log(f"处理国家 {country_code} 失败: {str(e)}", "ERROR")
self.log(traceback.format_exc(), "ERROR")
if "与页面的连接已断开" in e:
iskill = True
# 更新已处理国家数
if task_id in runing_task:
runing_task[task_id]["processed_countries"] += 1
# 最后回传,标记完成
try:
# self.post_result(task_id, shop_name, country_code, "", "", is_done=True)
@@ -443,15 +461,7 @@ class MatchTak:
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:
@@ -487,6 +497,15 @@ class MatchTak:
for retry in range(max_retries):
try:
self.log(f"尝试切换到国家 {country_name} (第 {retry + 1}/{max_retries} 次)")
if retry > 1:
# 刷新不行就重新打开店铺
self.log("重试前重新打开店铺...")
try:
driver.close_store()
time.sleep(3)
driver.open_shop(shop_name)
except Exception as e:
self.log(f"关闭重新打开店铺: {str(e)}", "WARNING")
# 如果不是第一次尝试,先刷新页面
if retry > 0:
@@ -559,7 +578,6 @@ class MatchTak:
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):
@@ -581,14 +599,7 @@ class MatchTak:
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} 处理完成")

View File

@@ -114,6 +114,30 @@ def calculate_standard_competitor_pricing(my_current_price, target_competitor_pr
class ChromeAmzone:
mark_name = "亚马逊详情采集"
country_info = {
"英国": {
"url": "https://www.amazon.co.uk/dp/B0CJ8SNXXV",
"zip_code": "SW1A 1AA",
"mark": "SW1A 1"
},
"德国": {
"url": "https://www.amazon.de/dp/B0CC8CW9G2?th=1",
"zip_code": "10115"
},
"法国": {
"url": "https://www.amazon.fr/dp/B0FRG1MJ8H?th=1",
"zip_code": "75001"
},
"西班牙": {
"url": "https://www.amazon.es/dp/B08ZXVNYNN",
"zip_code": "28001"
},
"意大利": {
"url": "https://www.amazon.it/dp/B0D1P17T2Q",
"zip_code": "20121"
}
}
def __init__(self,tab):
self.tab = tab
@@ -136,8 +160,92 @@ class ChromeAmzone:
if len(accept_btn) > 0:
accept_btn[0].click()
def _set_zip_code(self, zip_code, mark=None):
"""
设置邮编
def run(self):
Args:
zip_code: 目标邮编
"""
try:
# 检查当前邮编
zip_display = self.tab.ele('xpath://div[@id="glow-ingress-block"]', timeout=10)
if zip_display:
current_text = zip_display.text
# print(f"当前地址信息: {current_text}")
# 先检查标识
if mark is not None and mark in current_text:
print(f"邮编检测到标识: {mark},无需修改")
return True
# 检查是否已经包含目标邮编
if zip_code in current_text:
print(f"邮编已经设置为: {zip_code},无需修改")
return True
# 需要设置邮编
print(f"正在设置邮编为: {zip_code}")
# 点击地址选择按钮
location_link = self.tab.ele('xpath://a[@id="nav-global-location-popover-link"]', timeout=10)
if not location_link:
print("找不到地址设置按钮")
return False
location_link.click()
time.sleep(1)
# 等待邮编输入框出现
zip_input = self.tab.ele('xpath://input[@id="GLUXZipUpdateInput"]', timeout=10)
if not zip_input:
print("找不到邮编输入框")
return False
# 输入邮编
zip_input.input(zip_code, clear=True)
time.sleep(0.5)
# 点击提交按钮
submit_btn = self.tab.ele('xpath://input[@aria-labelledby="GLUXZipUpdate-announce"]', timeout=10)
if not submit_btn:
print("找不到提交按钮")
return False
submit_btn.click()
continue_btn = self.tab.eles('xpath://div[@class="a-popover-footer"]//input[@id="GLUXConfirmClose"]',
timeout=10)
if len(continue_btn) > 0:
continue_btn[0].click()
# 等待提交按钮消失(表示请求已发送)
print("等待邮编更新...")
time.sleep(2)
# 等待页面加载完成
self.tab.wait.doc_loaded(timeout=30, raise_err=False)
time.sleep(2)
# 验证邮编是否设置成功
zip_display_after = self.tab.ele('xpath://div[@id="glow-ingress-block"]', timeout=10)
if zip_display_after:
updated_text = zip_display_after.text
# print(f"更新后的地址信息: {updated_text}")
if zip_code in updated_text:
print(f"邮编设置成功: {zip_code}")
return True
else:
# print(f"邮编设置可能失败,当前显示: {updated_text}")
return False
return True
except Exception as e:
print(f"设置邮编时出错: {traceback.format_exc()}")
return False
def run(self,country):
"""
运行亚马逊详情采集任务
@@ -149,6 +257,20 @@ class ChromeAmzone:
dict: 包含采集到的数据
"""
try:
if country not in self.country_info:
error_msg = f"不支持的国家: {country},支持的国家有: {list(self.country_info.keys())}"
print(error_msg)
show_notification(error_msg, "error")
return None
# 获取国家配置
country_config = self.country_info[country]
zip_code = country_config["zip_code"]
mark = country_config.get("mark")
self._set_zip_code(zip_code,mark)
self.tab.wait.doc_loaded(timeout=5,raise_err=False)
self.tab.wait.doc_loaded(timeout=30, raise_err=False)
@@ -379,7 +501,15 @@ class AmzonePriceMatch(AmamzonBase):
print("当前操作的tab_id",self.tab.tab_id)
self.browser.close_tabs(close_tab)
def run_page_action(self,current_shop_name:str,appoint_asin:str=None,skip_asin:list=[]):
def run_page_action(self,current_shop_name:str,mode:str,
appoint_asin:str=None,skip_asin:list=[],miniprice_info:dict={}):
"""
miniprice_info :
{
"B0DSJGJBWV" : "19.81" # asin : 最低价
}
"""
print(f"{self.mark_name}】,开始执行")
num = 0
retry_num = 0
@@ -453,6 +583,24 @@ class AmzonePriceMatch(AmamzonBase):
if len(current_price_ele) > 0:
current_price = current_price_ele[0].attr("value")
print("获取到当前价格为 ->",current_price)
# 如果紫鸟里当前价格低于最低价则删除,否则跳过
if mode == "status" and miniprice_info.get(asin):
backend_price = float(miniprice_info.get(asin))
if float(current_price) < backend_price:
yield (asin, {
"statu": "删除(当前价格低于最低价)",
"deleteSkipAsin": True,
"removeAsin": asin,
})
continue
else:
yield (asin, {
"statu": "跳过(当前价格不低于最低价)",
"deleteSkipAsin": True,
"removeAsin": asin,
})
continue
else:
print(f"{self.mark_name} 没有获取到当前价格")
yield (asin, {
@@ -476,10 +624,16 @@ class AmzonePriceMatch(AmamzonBase):
time.sleep(1)
pass
chrome = ChromeAmzone(tab=new_tab)
front_end_data = chrome.run()
front_end_data = chrome.run(country=current_country)
chrome.close()
print(self.mark_name,"亚马逊前台抓取到数据",front_end_data)
#
if front_end_data is None or len(front_end_data.get("top_sellers")) == 0:
yield (asin, {
"statu": "失败(第一名获取失败,请检查卖家精灵是否启用)",
"currentPrice": current_price,
})
continue
cart_seller = front_end_data.get("cart_seller")
if cart_seller == current_shop_name and len(front_end_data.get("top_sellers"))< 2:
yield (asin, {
@@ -769,6 +923,8 @@ class PriceTask:
shopMallName = shop_item.get("shopMallName","")
skip_asins_by_country = shop_item.get("skip_asins_by_country",{})
asin_rows_by_country = shop_item.get("asin_rows_by_country",{})
skip_asin_details_by_country = shop_item.get("skip_asin_details_by_country",{})
mode = shop_item.get("mode")
if not company_name:
self.log(f"店铺 {shop_name} 的公司名称为空,跳过", "WARNING")
@@ -785,30 +941,42 @@ class PriceTask:
# 店铺打开重试最多3次
driver = None
max_retries = 3
iskill = False
# 处理每个国家
for country_code in country_codes:
# 打开店铺
driver = self.open_shop(max_retries=max_retries,company_name=company_name,shop_name=shop_name)
driver = self.open_shop(max_retries=max_retries, company_name=company_name,
shop_name=shop_name,iskill=iskill)
if driver is None:
self.log(f"任务 {task_id} 启动店铺失败,结束任务", "ERROR")
for country_code in country_codes:
self.post_result(task_id, shop_name, country_code, "", {}, shopMallName, is_done=True)
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
skip_asin = skip_asins_by_country.get(country_code,[]) #需要跳过的asin
appoint_asin = asin_rows_by_country.get(country_code,[])
miniprice_info = {i.get('asin'):i.get('minimumPrice') for i in skip_asin_details_by_country.get(country_code,[])}
try:
self.process_country(driver, country_code, task_id, shop_name, risk_listing_filter,
shopMallName=shopMallName,skip_asin=skip_asin, limit=limit,appoint_asin=appoint_asin)
shopMallName=shopMallName,skip_asin=skip_asin, limit=limit,
appoint_asin=appoint_asin,mode=mode,miniprice_info=miniprice_info)
except Exception as e:
import traceback
self.log(f"处理国家 {country_code} 失败: {str(e)}", "ERROR")
self.log(traceback.format_exc(), "ERROR")
if "与页面的连接已断开" in e:
iskill = True
# 更新已处理国家数
if task_id in runing_task:
runing_task[task_id]["processed_countries"] += 1
# 关闭店铺
driver.close_store()
# 最后回传,标记完成
try:
# self.post_result(task_id, shop_name, country_code, "", "", is_done=True)
@@ -818,15 +986,6 @@ class PriceTask:
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:
@@ -834,7 +993,9 @@ class PriceTask:
self.log(f"店铺 {shop_name} 已从执行列表中移除")
def process_country(self, driver: AmzonePriceMatch, country_code: str, task_id: int, shop_name: str,
risk_listing_filter: str,shopMallName:str,skip_asin:list,appoint_asin:list,limit: str = None):
risk_listing_filter: str,shopMallName:str,skip_asin:list,appoint_asin:list,mode:str,
miniprice_info:dict={},
limit: str = None):
"""处理单个国家的审批任务
Args:
@@ -864,7 +1025,7 @@ class PriceTask:
self.log(f"尝试切换到国家 {country_name} (第 {retry + 1}/{max_retries} 次)")
if retry > 1:
# 刷新不行就重新打开店铺
self.log("重试前刷新页面...")
self.log("重试前重新打开店铺...")
try:
driver.close_store()
time.sleep(3)
@@ -946,7 +1107,7 @@ class PriceTask:
self.log(f"国家 {country_name} 搜索出 {len(sku_ls)} 商品,开始处理...")
# 处理所有需要审批的商品通过yield获取结果
try:
# 指定 asin
max_range = max(1,len(appoint_asin))
for i in range(max_range):
@@ -959,7 +1120,8 @@ class PriceTask:
for asin, status in driver.run_page_action(
current_shop_name=_shopMallName,
appoint_asin=ap_asin,skip_asin=skip_asin
appoint_asin=ap_asin,skip_asin=skip_asin,mode=mode,
miniprice_info=miniprice_info
):
# 检查是否收到暂停请求
if task_id in runing_task and runing_task[task_id].get("stop_requested", False):
@@ -981,15 +1143,6 @@ class PriceTask:
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} 处理完成")
@@ -1073,6 +1226,8 @@ class PriceTask:
"secondPlace": status.get("secondPlace") if status.get("secondPlace") else "",
"cartShopName": status.get("cartShopName") if status.get("cartShopName") else "",
"priceChangeStatus": "UPDATED",
"deleteSkipAsin": status.get("deleteSkipAsin",False),
"removeAsin": status.get("removeAsin") if status.get("removeAsin") else "",
# "modifyCount": "2",
"status": status.get("statu")
}