modified: app/ali_oss.py

new file:   app/amazon/del_brand.py
	new file:   app/amazon/main.py
	modified:   app/blueprints/__pycache__/brand.cpython-39.pyc
	modified:   app/blueprints/__pycache__/communication.cpython-39.pyc
	modified:   app/blueprints/__pycache__/main.cpython-39.pyc
	modified:   app/blueprints/brand.py
	new file:   "app/blueprints/brand_\345\244\207\344\273\275.py"
	modified:   app/blueprints/communication.py
	modified:   app/blueprints/main.py
	modified:   app/config.py
	modified:   app/main.py
This commit is contained in:
铭坤
2026-04-01 11:53:12 +08:00
parent b6dd1d7931
commit 785d7563f9
12 changed files with 1779 additions and 81 deletions

519
app/amazon/del_brand.py Normal file
View File

@@ -0,0 +1,519 @@
import traceback
import winreg
import subprocess
import time
import uuid
import requests
import json
import os
from typing import Literal
from DrissionPage import Chromium
from DrissionPage.common import By
def kill_process(version: Literal["v5", "v6"]):
"""杀紫鸟客户端进程(独立函数版本)"""
driver = ZiniaoDriver({})
driver.kill_process(version)
class ZiniaoDriver:
"""紫鸟浏览器自动化驱动类"""
def __init__(self, user_info: dict, socket_port: int = 19890):
"""
初始化紫鸟浏览器驱动
Args:
user_info: 用户信息字典,包含 company, username, password
socket_port: 客户端通信端口,默认 19890
"""
self.user_info = user_info
self.socket_port = socket_port
self.client_path = None
self.browser = None
self.tab = None
self.store_id = None
def get_zinaio_exe(self, protocol_name: str = "superbrowser"):
"""
获取紫鸟安装目录
Args:
protocol_name: 协议名称,默认 "superbrowser"
Returns:
exe_path: 紫鸟浏览器可执行文件路径
"""
try:
key_path = rf"SOFTWARE\Classes\{protocol_name}\shell\open\command"
key = winreg.OpenKey(winreg.HKEY_CURRENT_USER, key_path)
command, _ = winreg.QueryValueEx(key, "")
winreg.CloseKey(key)
if isinstance(command, str):
exe_path = command.split()[0]
else:
exe_path = command[0]
return exe_path
except FileNotFoundError:
try:
key = winreg.OpenKey(winreg.HKEY_LOCAL_MACHINE, key_path)
command, _ = winreg.QueryValueEx(key, "")
winreg.CloseKey(key)
if isinstance(command, str):
exe_path = command.split()[0]
else:
exe_path = command[0]
return exe_path
except FileNotFoundError:
return None
def update_core(self):
"""
下载所有内核打开店铺前调用需客户端版本5.285.7以上
因为http有超时时间所以这个action适合循环调用直到返回成功
"""
data = {
"action": "updateCore",
"requestId": str(uuid.uuid4()),
}
data.update(self.user_info)
while True:
url = f'http://127.0.0.1:{self.socket_port}'
response = requests.post(url, json.dumps(data).encode('utf-8'), timeout=120)
result = response.json()
print(result)
if result is None:
print("等待客户端启动...")
time.sleep(2)
continue
if result.get("statusCode") is None or result.get("statusCode") == -10003:
print("当前版本不支持此接口,请升级客户端")
return
elif result.get("statusCode") == 0:
print("更新内核完成")
return
else:
print(f"等待更新内核: {json.dumps(result)}")
time.sleep(2)
def kill_process(self, version: Literal["v5", "v6"]):
"""
杀紫鸟客户端进程
Args:
version: 客户端版本
"""
if version == "v5":
process_name = 'SuperBrowser.exe'
else:
process_name = 'ziniao.exe'
os.system('taskkill /f /t /im ' + process_name)
time.sleep(3)
def get_browser_list(self) -> list:
"""
获取浏览器列表
Returns:
list: 浏览器列表
"""
request_id = str(uuid.uuid4())
data = {
"action": "getBrowserList",
"requestId": request_id
}
data.update(self.user_info)
url = f'http://127.0.0.1:{self.socket_port}'
response = requests.post(url, json.dumps(data).encode('utf-8'), timeout=120)
r = response.json()
if str(r.get("statusCode")) == "0":
print(r)
return r.get("browserList")
elif str(r.get("statusCode")) == "-10003":
print(f"login Err {json.dumps(r, ensure_ascii=False)}")
exit()
else:
print(f"Fail {json.dumps(r, ensure_ascii=False)} ")
exit()
def open_store(self, store_info, isWebDriverReadOnlyMode=0, isprivacy=0,
isHeadless=0, cookieTypeSave=0, jsInfo=""):
"""
打开店铺
Args:
store_info: 店铺信息browserId 或 browserOauth
isWebDriverReadOnlyMode: 是否只读模式,默认 0
isprivacy: 隐私模式,默认 0
isHeadless: 无头模式,默认 0
cookieTypeSave: cookie保存类型默认 0
jsInfo: 注入的JS信息默认 ""
Returns:
dict: 返回结果
"""
request_id = str(uuid.uuid4())
data = {
"action": "startBrowser",
"isWaitPluginUpdate": 0,
"isHeadless": isHeadless,
"requestId": request_id,
"isWebDriverReadOnlyMode": isWebDriverReadOnlyMode,
"cookieTypeLoad": 0,
"cookieTypeSave": cookieTypeSave,
"runMode": "1",
"isLoadUserPlugin": False,
"pluginIdType": 1,
"privacyMode": isprivacy
}
data.update(self.user_info)
if store_info.isdigit():
data["browserId"] = store_info
else:
data["browserOauth"] = store_info
if len(str(jsInfo)) > 2:
data["injectJsInfo"] = json.dumps(jsInfo)
url = f'http://127.0.0.1:{self.socket_port}'
response = requests.post(url, json.dumps(data).encode('utf-8'), timeout=120)
r = response.json()
if str(r.get("statusCode")) == "0":
return r
elif str(r.get("statusCode")) == "-10003":
print(f"login Err {json.dumps(r, ensure_ascii=False)}")
exit()
else:
print(f"Fail {json.dumps(r, ensure_ascii=False)} ")
exit()
def get_browser(self, port) -> Chromium:
"""
获取浏览器实例
Args:
port: 调试端口
Returns:
Chromium: DrissionPage浏览器实例
"""
browser = Chromium(port)
return browser
def start_client(self):
"""启动紫鸟客户端"""
self.client_path = self.get_zinaio_exe("superbrowser").strip('"')
print(self.client_path)
cmd = [self.client_path, '--run_type=web_driver', '--ipc_type=http',
'--port=' + str(self.socket_port)]
print(" ".join(cmd))
subprocess.Popen(cmd)
time.sleep(5)
def open_shop(self, shop_name: str):
"""
打开指定店铺
Args:
shop_name: 店铺名称
Returns:
Chromium or str: 成功返回浏览器实例,失败返回错误信息
"""
# 启动客户端
self.start_client()
# 更新内核
self.update_core()
# 获取店铺列表
shop_ls = self.get_browser_list()
self.store_id = None
# print(shop_ls)
for shop in shop_ls:
if shop.get("browserName") == shop_name:
self.store_id = shop.get('browserOauth')
break
if not self.store_id:
print("店铺不存在")
return "店铺不存在"
# 打开店铺
ret_json = self.open_store(self.store_id)
print(ret_json)
self.store_id = ret_json.get("browserOauth")
if self.store_id is None:
self.store_id = ret_json.get("browserId")
# 获取drissionpage浏览器会话
self.browser = self.get_browser(ret_json.get('debuggingPort'))
ip_check_url = ret_json.get("ipDetectionPage")
if not ip_check_url:
print("ip检测页地址为空请升级紫鸟浏览器到最新版")
print(f"=====关闭店铺:{shop_name}=====")
self.close_store(self.store_id)
exit()
ip_usable = self.open_ip_check(self.browser, ip_check_url)
if ip_usable:
print("ip检测通过打开店铺平台主页")
self.open_launcher_page(ret_json.get("launcherPage"), self.browser)
else:
print("IP检测不通过")
return self.browser
def close_store(self, browser_oauth=None):
"""
关闭店铺
Args:
browser_oauth: 店铺OAuth标识如果不提供则使用当前打开的店铺
Returns:
dict: 返回结果
"""
if browser_oauth is None:
browser_oauth = self.store_id
request_id = str(uuid.uuid4())
data = {
"action": "stopBrowser",
"requestId": request_id,
"duplicate": 0,
"browserOauth": browser_oauth
}
data.update(self.user_info)
url = f'http://127.0.0.1:{self.socket_port}'
response = requests.post(url, json.dumps(data).encode('utf-8'), timeout=120)
r = response.json()
if str(r.get("statusCode")) == "0":
return r
elif str(r.get("statusCode")) == "-10003":
print(f"login Err {json.dumps(r, ensure_ascii=False)}")
exit()
else:
print(f"Fail {json.dumps(r, ensure_ascii=False)} ")
exit()
def open_launcher_page(self, launcher_page: str, browser: Chromium = None):
"""
打开启动页面
Args:
launcher_page: 要打开的页面URL
browser: 浏览器实例,如果不提供则使用当前浏览器实例
"""
if browser is None:
browser = self.browser
tab = browser.new_tab(url=launcher_page)
self.tab = tab
return tab
def open_ip_check(self, browser: Chromium, ip_check_url: str):
"""
打开ip检测页检测ip是否正常
:param browser: drissionpage浏览器会话
:param ip_check_url ip检测页地址
:return 检测结果
"""
try:
tab = browser.latest_tab
tab.get(ip_check_url)
success_button = tab.ele((By.XPATH, '//button[contains(@class, "styles_btn--success")]'),
timeout=60) # 等待查找元素60秒
if success_button:
print("ip检测成功")
return True
else:
print("ip检测超时")
return False
except Exception as e:
print("ip检测异常:" + traceback.format_exc())
return False
class AmazoneDriver(ZiniaoDriver):
"""亚马逊专用驱动类继承自ZiniaoDriver"""
def SwitchingCountries(self, country_name: str):
"""
切换国家
操作:
1、//div[@class="dropdown-account-switcher-header-label"]/span[last()] 获取此元素文本,判断当前国家,如果与目标国家相同则不操作,否则执行下一步
2、点击 //div[@class="dropdown-account-switcher-header-label"] 打开下拉框
3、点击 //div[@class="dropdown-account-switcher-list-item"] 第一个展开国家列表
4、点击 //div[@class="dropdown-account-switcher-list-item dropdown-account-switcher-list-item-indented" and @title="国家名"] 切换到目标国家
5、等待页面加载完成判断国家是否切换成功成功则返回True否则返回False
Args:
country_name: 目标国家名称
Returns:
bool: 切换成功返回True失败返回False
"""
try:
if self.browser is None:
print("浏览器实例不存在,请先打开店铺")
return False
# 获取当前标签页
tab = self.tab
# 步骤1获取当前国家名称
print(f"正在检查当前国家...")
current_country_ele = tab.ele('xpath://div[@class="dropdown-account-switcher-header-label"]/span[last()]',
timeout=10)
if current_country_ele:
current_country = current_country_ele.text.strip()
print(f"当前国家:{current_country}")
# 判断是否与目标国家相同
if current_country == country_name:
print(f"当前已经是目标国家 {country_name},无需切换")
return True
else:
print("无法获取当前国家信息")
return False
# 步骤2点击打开下拉框
print(f"正在打开国家切换下拉框...")
dropdown_header = tab.ele('xpath://div[@class="dropdown-account-switcher-header-label"]', timeout=10)
if not dropdown_header:
print("找不到国家切换下拉框")
return False
dropdown_header.click()
time.sleep(1) # 等待下拉框展开
# 步骤3点击第一个展开国家列表
print(f"正在展开国家列表...")
first_item = tab.ele('xpath://div[@class="dropdown-account-switcher-list-item"]', timeout=10)
if not first_item:
print("找不到国家列表项")
return False
first_item.click()
time.sleep(1) # 等待国家列表展开
# 步骤4点击目标国家
print(f"正在切换到国家:{country_name}")
target_country_xpath = f'//div[@class="dropdown-account-switcher-list-item dropdown-account-switcher-list-item-indented" and @title="{country_name}"]'
target_country = tab.ele(f'xpath:{target_country_xpath}', timeout=10)
if not target_country:
print(f"找不到目标国家:{country_name}")
return False
target_country.click()
# 步骤5等待页面加载完成并验证切换结果
print(f"等待页面加载...")
# time.sleep(3) # 等待页面加载
self.tab.wait.doc_loaded()
# 再次检查当前国家
new_country_ele = tab.ele('xpath://div[@class="dropdown-account-switcher-header-label"]/span[last()]',
timeout=10)
if new_country_ele:
new_country = new_country_ele.text.strip()
if new_country == country_name:
print(f"国家切换成功:{new_country}")
return True
else:
print(f"国家切换失败,当前国家:{new_country},目标国家:{country_name}")
return False
else:
print("无法验证切换结果")
return False
except Exception as e:
print(f"切换国家时发生异常:{traceback.format_exc()}")
return False
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, asin):
search_region = self.tab.ele('xpath://div[@id="searchBoxContainer"]//kat-input-group')
search_region.wait.displayed(raise_err=False)
time.sleep(0.6)
search_input = self.tab.ele("xpath://kat-input[contains(@class,'SearchBox-module__searchInput')]").sr(
'xpath://span[@class="container"]//input[@part="input"]')
search_input.input(asin)
for _ in range(3):
search_btn = self.tab.ele("xpath://kat-icon[@name='search']")
search_btn.click()
load_ele = self.tab.ele("xpath://div[contains(@class,'Loader-module__loader')]")
load_ele.wait.hidden(timeout=20)
time.sleep(1)
sku_ls = self.tab.eles("//div[@data-sku]")
return sku_ls
def del_action(self, sku_ele):
dropdown = sku_ele.ele('xpath:.//kat-dropdown-button[@variant="secondary"]')
dropdown.click()
time.sleep(1)
del_btn = dropdown.ele("xpath:.//button[@role='menuitem' and text()='删除商品信息']")
del_btn.wait.displayed(raise_err=False)
del_btn.click()
confirm_btn = self.tab.ele('xpath://kat-modal[@data-testid="action-modal"]//kat-button[@variant="primary"]')
confirm_btn.wait.displayed(raise_err=False)
confirm_btn.wait.enabled()
confirm_btn.click()
suc_alert = self.tab.ele("xpath://kat-alert[@variant='success']")
suc = suc_alert.wait.displayed(raise_err=False, timeout=5)
return suc
if __name__ == '__main__':
# 使用示例
user_info = {
"company": "rongchuang123",
"username": "自动化_Robot",
"password": "#20zsg25"
}
shop_name = "郭亚芳"
kill_process("v6")
# 创建驱动实例
driver = AmazoneDriver(user_info)
# 打开店铺并获取浏览器实例
browser = driver.open_shop(shop_name)
country = "西班牙"
asin = "Voanos"
sw_suc = driver.SwitchingCountries(country)
driver.SwitchPage()
sku_ls = driver.search(asin=asin)
print("查询结果有:", len(sku_ls))
for sku in sku_ls:
suc = driver.del_action(sku)
print(sku, "删除结果", suc)

421
app/amazon/main.py Normal file
View File

@@ -0,0 +1,421 @@
import time
import traceback
import requests
from datetime import datetime
from typing import Dict, Any, List
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
class TaskMonitor:
"""任务监控器:负责监控队列并执行品牌删除任务"""
def __init__(self):
"""初始化任务监控器"""
self.running = True
self.user_info = {
"company": ZN_COMPANY,
"username": ZN_USERNAME,
"password": ZN_PASSWORD
}
self.chunk_index = 1 # 当前处理的分块索引
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("任务监控器启动,开始监听队列...")
while self.running:
try:
# 阻塞式获取任务避免CPU空转
task_data = JSON_TASK_QUEUE.get(block=True, timeout=1)
# 检查任务类型
task_type = task_data.get("type", "")
if task_type != "delete-brand-run":
self.log(f"未知任务类型: {task_type},跳过", "WARNING")
continue
# 处理任务
self.log(f"接收到删除品牌任务,开始处理...")
self.process_task(task_data)
except Exception as e:
if "Empty" not in str(e): # 忽略队列为空的超时异常
self.log(f"任务处理异常: {traceback.format_exc()}", "ERROR")
time.sleep(0.1) # 短暂等待避免异常循环
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", [])
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"店铺 {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")
time.sleep(2)
# 创建驱动并打开店铺
driver = AmazoneDriver(self.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
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:
# 处理每个国家
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:
self.process_country(driver, country_data, task_id, result_id, shop_data)
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]):
"""处理单个国家的所有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)
# 切换国家
try:
switch_success = driver.SwitchingCountries(country)
if not switch_success:
self.log(f"切换到国家 {country} 失败,跳过该国家", "ERROR")
return
except Exception as e:
self.log(f"切换国家 {country} 异常: {str(e)}", "ERROR")
return
# 切换到库存管理页面
try:
driver.SwitchPage()
self.log(f"已切换到库存管理页面")
except Exception as e:
self.log(f"切换页面失败: {str(e)}", "ERROR")
return
# 处理每个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, idx)
except Exception as e:
asin = asin_item.get("asin", "未知")
self.log(f"处理ASIN {asin} 失败: {str(e)}", "ERROR")
# 继续处理下一个ASIN
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):
"""处理单个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 = "失败"
try:
# 搜索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")
else:
# 删除所有找到的SKU
success_count = 0
for sku in sku_ls:
try:
suc = driver.del_action(sku)
if suc:
success_count += 1
self.log(f"SKU 删除成功")
else:
self.log(f"SKU 删除失败", "WARNING")
except Exception as e:
self.log(f"删除SKU异常: {str(e)}", "ERROR")
if success_count > 0:
status = "成功"
self.log(f"ASIN {asin} 删除成功 ({success_count}/{len(sku_ls)})")
# 更新成功计数
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 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": self.chunk_index,
"chunkTotal": total_rows,
"processedRows": self.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 post_result(self, task_id: int, payload: Dict[str, Any]):
"""回传结果到API
Args:
task_id: 任务ID
payload: 结果数据
"""
url = f"{DELETE_BRAND_API_BASE}/newApi/api/delete-brand/tasks/{task_id}/result"
try:
response = requests.post(
url,
json=payload,
headers={"Content-Type": "application/json"},
timeout=30,
verify=False # 忽略SSL证书验证
)
if response.status_code == 200:
self.log(f"结果回传成功: {url}")
else:
self.log(f"结果回传失败,状态码: {response.status_code}", "WARNING")
except Exception as e:
self.log(f"调用API异常: {str(e)}", "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 as e:
monitor.log(f"监控器异常退出: {traceback.format_exc()}", "ERROR")
if __name__ == "__main__":
main()