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:
421
app/amazon/main.py
Normal file
421
app/amazon/main.py
Normal 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()
|
||||
Reference in New Issue
Block a user