Compare commits
11 Commits
855f1affd9
...
46d91fd8ca
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
46d91fd8ca | ||
|
|
a24db2f73c | ||
|
|
ca9537f862 | ||
|
|
94ae7e68ff | ||
|
|
94e1930538 | ||
|
|
5bbe5a3077 | ||
|
|
8cead01270 | ||
|
|
99fc14b9f9 | ||
|
|
d3671206e1 | ||
|
|
0a8202788b | ||
|
|
5a294e5c85 |
3
app/.env
3
app/.env
@@ -6,6 +6,9 @@ mysql_user=aiimage
|
||||
proxy_url=https://api.jikip.com/ip-get?num=1&minute=1&format=json&area=all&protocol=1&mode=2&key=t24g6gi44ubufd8
|
||||
proxy_mode=2
|
||||
|
||||
zn_company=rongchuang123
|
||||
zn_username=%E8%87%AA%E5%8A%A8%E5%8C%96_Robot
|
||||
|
||||
client_name=ShuFuAI
|
||||
|
||||
|
||||
|
||||
15
app/.env copy
Normal file
15
app/.env copy
Normal file
@@ -0,0 +1,15 @@
|
||||
base_url=http://8.136.19.173:15124
|
||||
workflow_id=7608812635877900322
|
||||
mysql_host=8.136.19.173
|
||||
mysql_user=aiimage
|
||||
|
||||
proxy_url=https://api.jikip.com/ip-get?num=1&minute=1&format=json&area=all&protocol=1&mode=2&key=t24g6gi44ubufd8
|
||||
proxy_mode=2
|
||||
|
||||
client_name=ShuFuAI
|
||||
|
||||
|
||||
java_api_base=http://127.0.0.1:18080
|
||||
# java_api_base=http://8.136.19.173:18080
|
||||
|
||||
|
||||
Binary file not shown.
BIN
app/amazon/__pycache__/del_brand.cpython-312.pyc
Normal file
BIN
app/amazon/__pycache__/del_brand.cpython-312.pyc
Normal file
Binary file not shown.
BIN
app/amazon/__pycache__/del_brand.cpython-39.pyc
Normal file
BIN
app/amazon/__pycache__/del_brand.cpython-39.pyc
Normal file
Binary file not shown.
BIN
app/amazon/__pycache__/main.cpython-312.pyc
Normal file
BIN
app/amazon/__pycache__/main.cpython-312.pyc
Normal file
Binary file not shown.
BIN
app/amazon/__pycache__/main.cpython-39.pyc
Normal file
BIN
app/amazon/__pycache__/main.cpython-39.pyc
Normal file
Binary file not shown.
@@ -233,7 +233,7 @@ class ZiniaoDriver:
|
||||
# 获取店铺列表
|
||||
shop_ls = self.get_browser_list()
|
||||
self.store_id = None
|
||||
# print(shop_ls)
|
||||
print(shop_ls)
|
||||
for shop in shop_ls:
|
||||
if shop.get("browserName") == shop_name:
|
||||
self.store_id = shop.get('browserOauth')
|
||||
@@ -456,16 +456,20 @@ class AmazoneDriver(ZiniaoDriver):
|
||||
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)
|
||||
|
||||
search_input.input(asin,clear=True)
|
||||
sku_ls = []
|
||||
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]")
|
||||
# load_ele.wait.hidden(timeout=3, raise_err=False)
|
||||
load_ele.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
|
||||
return sku_ls
|
||||
|
||||
def del_action(self, sku_ele):
|
||||
@@ -473,7 +477,7 @@ class AmazoneDriver(ZiniaoDriver):
|
||||
dropdown.click()
|
||||
time.sleep(1)
|
||||
|
||||
del_btn = dropdown.ele("xpath:.//button[@role='menuitem' and text()='删除商品信息']")
|
||||
del_btn = dropdown.sr("xpath:.//button[@role='menuitem' and @data-action='DeleteListing']")
|
||||
del_btn.wait.displayed(raise_err=False)
|
||||
del_btn.click()
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@ 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
|
||||
|
||||
@@ -20,6 +21,11 @@ class TaskMonitor:
|
||||
}
|
||||
|
||||
self.chunk_index = 1 # 当前处理的分块索引
|
||||
self.max_workers = 5 # 线程池最大线程数
|
||||
self.executor = None # 线程池执行器
|
||||
|
||||
# 在提交新任务前杀掉旧进程(确保环境干净)
|
||||
kill_process("v6")
|
||||
|
||||
def log(self, message: str, level: str = "INFO"):
|
||||
"""日志输出
|
||||
@@ -32,28 +38,60 @@ class TaskMonitor:
|
||||
print(f"[{timestamp}] [{level}] {message}")
|
||||
|
||||
def start(self):
|
||||
"""启动任务监控(阻塞式循环)"""
|
||||
self.log("任务监控器启动,开始监听队列...")
|
||||
"""启动任务监控(使用线程池处理任务)"""
|
||||
self.log(f"任务监控器启动,开始监听队列... (线程池大小: {self.max_workers})")
|
||||
|
||||
while self.running:
|
||||
try:
|
||||
# 阻塞式获取任务,避免CPU空转
|
||||
task_data = JSON_TASK_QUEUE.get(block=True, timeout=1)
|
||||
# 创建线程池
|
||||
self.executor = ThreadPoolExecutor(max_workers=self.max_workers)
|
||||
futures = [] # 保存所有提交的任务Future对象
|
||||
|
||||
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"未知任务类型: {task_type},跳过", "WARNING")
|
||||
continue
|
||||
|
||||
# 提交任务到线程池
|
||||
self.log(f"接收到删除品牌任务,提交到线程池处理...")
|
||||
future = self.executor.submit(self._process_task_wrapper, task_data)
|
||||
futures.append(future)
|
||||
|
||||
# 清理已完成的future对象,避免内存累积
|
||||
futures = [f for f in futures if not f.done()]
|
||||
|
||||
except Exception as e:
|
||||
# 静默处理队列为空的超时,只记录真正的异常
|
||||
if "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")
|
||||
|
||||
# 检查任务类型
|
||||
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]):
|
||||
"""处理单个任务
|
||||
@@ -177,6 +215,7 @@ class TaskMonitor:
|
||||
|
||||
try:
|
||||
# 处理每个国家
|
||||
chunk_index =1
|
||||
for country_data in countries:
|
||||
# 检查是否收到暂停请求
|
||||
if task_id in runing_task and runing_task[task_id].get("stop_requested", False):
|
||||
@@ -184,7 +223,7 @@ class TaskMonitor:
|
||||
break # 跳出循环,进入finally关闭店铺
|
||||
|
||||
try:
|
||||
self.process_country(driver, country_data, task_id, result_id, shop_data)
|
||||
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")
|
||||
@@ -206,7 +245,7 @@ class TaskMonitor:
|
||||
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]):
|
||||
task_id: int, result_id: int, shop_data: Dict[str, Any],chunk_index:int):
|
||||
"""处理单个国家的所有ASIN
|
||||
|
||||
Args:
|
||||
@@ -222,15 +261,42 @@ class TaskMonitor:
|
||||
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
|
||||
# 切换国家,最多重试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:
|
||||
@@ -238,7 +304,9 @@ class TaskMonitor:
|
||||
self.log(f"已切换到库存管理页面")
|
||||
except Exception as e:
|
||||
self.log(f"切换页面失败: {str(e)}", "ERROR")
|
||||
return
|
||||
# 切换页面失败也回传所有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", "")
|
||||
@@ -254,15 +322,18 @@ class TaskMonitor:
|
||||
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)
|
||||
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):
|
||||
source_filename: str, total_rows: int,chunk_index:int):
|
||||
"""处理单个ASIN并回传结果
|
||||
|
||||
Args:
|
||||
@@ -283,48 +354,89 @@ class TaskMonitor:
|
||||
self.update_task_status(task_id, current_asin=asin)
|
||||
|
||||
status = "失败"
|
||||
max_retries = 3 # 最多重试3次
|
||||
|
||||
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:
|
||||
for retry in range(max_retries):
|
||||
try:
|
||||
self.log(f"处理ASIN {asin} (第 {retry + 1}/{max_retries} 次)")
|
||||
|
||||
# 如果不是第一次尝试,先刷新页面
|
||||
if retry > 0:
|
||||
self.log("重试前刷新页面...")
|
||||
try:
|
||||
suc = driver.del_action(sku)
|
||||
if suc:
|
||||
success_count += 1
|
||||
self.log(f"SKU 删除成功")
|
||||
else:
|
||||
self.log(f"SKU 删除失败", "WARNING")
|
||||
driver.tab.refresh()
|
||||
time.sleep(3)
|
||||
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
|
||||
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:
|
||||
status = "失败"
|
||||
self.log(f"ASIN {asin} 所有SKU删除失败", "ERROR")
|
||||
# 更新失败计数
|
||||
# 删除所有找到的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
|
||||
|
||||
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:
|
||||
@@ -337,9 +449,9 @@ class TaskMonitor:
|
||||
"files": [{
|
||||
"fileKey": file_key,
|
||||
"sourceFilename": source_filename,
|
||||
"chunkIndex": self.chunk_index,
|
||||
"chunkIndex": chunk_index,
|
||||
"chunkTotal": total_rows,
|
||||
"processedRows": self.chunk_index,
|
||||
"processedRows": chunk_index,
|
||||
"totalRows": total_rows,
|
||||
"currentCountry": country,
|
||||
"currentAsin": asin,
|
||||
@@ -359,31 +471,105 @@ class TaskMonitor:
|
||||
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
|
||||
"""回传结果到API(带重试机制)
|
||||
|
||||
Args:
|
||||
task_id: 任务ID
|
||||
payload: 结果数据
|
||||
"""
|
||||
url = f"{DELETE_BRAND_API_BASE}/newApi/api/delete-brand/tasks/{task_id}/result"
|
||||
url = f"{DELETE_BRAND_API_BASE}/api/delete-brand/tasks/{task_id}/result"
|
||||
max_retries = 3 # 最多重试3次
|
||||
|
||||
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")
|
||||
for retry in range(max_retries):
|
||||
try:
|
||||
self.log(f"尝试回传结果 (第 {retry + 1}/{max_retries} 次)")
|
||||
|
||||
except Exception as e:
|
||||
self.log(f"调用API异常: {str(e)}", "ERROR")
|
||||
response = requests.post(
|
||||
url,
|
||||
json=payload,
|
||||
headers={"Content-Type": "application/json"},
|
||||
timeout=30,
|
||||
verify=False # 忽略SSL证书验证
|
||||
)
|
||||
print("【结果提交】:",payload)
|
||||
|
||||
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):
|
||||
"""更新任务状态
|
||||
@@ -404,7 +590,7 @@ class TaskMonitor:
|
||||
|
||||
|
||||
def main():
|
||||
"""主函数:启动任务监控"""
|
||||
"""主函数:启动任务监控"""
|
||||
# 创建并启动任务监控器
|
||||
monitor = TaskMonitor()
|
||||
|
||||
@@ -413,7 +599,7 @@ def main():
|
||||
except KeyboardInterrupt:
|
||||
monitor.log("接收到中断信号,正在停止...")
|
||||
monitor.stop()
|
||||
except Exception as e:
|
||||
except Exception:
|
||||
monitor.log(f"监控器异常退出: {traceback.format_exc()}", "ERROR")
|
||||
|
||||
|
||||
|
||||
File diff suppressed because one or more lines are too long
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
@@ -110,14 +110,16 @@ def proxy(path):
|
||||
if key.lower() in ['host', 'content-length', 'connection']:
|
||||
continue
|
||||
headers[key] = value
|
||||
try:
|
||||
print("=============================")
|
||||
print("target_url:",target_url)
|
||||
print("params:",params)
|
||||
print("data:",request.get_data())
|
||||
print("=============================")
|
||||
except Exception as e:
|
||||
print("打印失败",e)
|
||||
ignore_url = [f"{JAVA_API_BASE}/api/delete-brand/tasks/batch",f"{JAVA_API_BASE}/api/delete-brand/history"]
|
||||
if target_url not in ignore_url:
|
||||
try:
|
||||
print("=============================")
|
||||
print("target_url:",target_url)
|
||||
print("params:",params)
|
||||
print("data:",request.get_data())
|
||||
print("=============================")
|
||||
except Exception as e:
|
||||
print("打印失败",e)
|
||||
|
||||
try:
|
||||
# 使用流式请求
|
||||
|
||||
@@ -21,9 +21,11 @@ client_name=os.getenv("client_name") + ".exe"
|
||||
JAVA_API_BASE = os.getenv("java_api_base", "http://127.0.0.1:18080")
|
||||
|
||||
# 紫鸟浏览器配置
|
||||
from urllib.parse import unquote
|
||||
ZN_COMPANY = os.getenv("zn_company", "")
|
||||
ZN_USERNAME = os.getenv("zn_username", "")
|
||||
ZN_PASSWORD = os.getenv("zn_password", "")
|
||||
ZN_USERNAME = unquote(ZN_USERNAME)
|
||||
ZN_PASSWORD = os.getenv("zn_password", "#20zsg25")
|
||||
|
||||
# 删除品牌API配置
|
||||
DELETE_BRAND_API_BASE = os.getenv("DELETE_BRAND_API_BASE", JAVA_API_BASE)
|
||||
|
||||
23
app/main.py
23
app/main.py
@@ -7,6 +7,7 @@ import json
|
||||
import sys
|
||||
import shutil
|
||||
import os
|
||||
from amazon.main import TaskMonitor
|
||||
os.makedirs(cache_path,exist_ok=True)
|
||||
if not debug:
|
||||
today = datetime.datetime.now().strftime("%Y_%m_%d")
|
||||
@@ -267,8 +268,8 @@ class WindowAPI:
|
||||
payload = json.loads(data)
|
||||
if not isinstance(payload, (dict, list)):
|
||||
return {'success': False, 'error': '仅支持 JSON 对象或数组'}
|
||||
if payload.get("type") == "delete-brand-run" and payload.get("data").get("taskId") in runing_task:
|
||||
return {'success': False, 'error': '已存在正在执行的删除品牌任务'}
|
||||
# if payload.get("type") == "delete-brand-run" and payload.get("data").get("taskId") in runing_task:
|
||||
# return {'success': False, 'error': '已存在正在执行的删除品牌任务'}
|
||||
# 判断当前店铺是否正在运行中
|
||||
if payload.get("data").get("items"):
|
||||
shop_name = payload["data"]["items"][0]["shopName"]
|
||||
@@ -301,11 +302,29 @@ def start_flask():
|
||||
run_app(host='127.0.0.1', port=PORT)
|
||||
|
||||
|
||||
def start_task_monitor():
|
||||
"""启动任务监控器(在独立线程中运行)"""
|
||||
monitor = TaskMonitor()
|
||||
try:
|
||||
print(f"[{datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 任务监控线程已启动")
|
||||
monitor.start()
|
||||
except Exception:
|
||||
import traceback
|
||||
print(f"[{datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 任务监控线程异常退出: {traceback.format_exc()}")
|
||||
|
||||
|
||||
def main():
|
||||
if not os.path.exists(cache_path):
|
||||
os.makedirs(cache_path,exist_ok=True)
|
||||
|
||||
# 启动 Flask 服务
|
||||
t = threading.Thread(target=start_flask, daemon=True)
|
||||
t.start()
|
||||
|
||||
# 启动任务监控线程
|
||||
monitor_thread = threading.Thread(target=start_task_monitor, daemon=True)
|
||||
monitor_thread.start()
|
||||
|
||||
# 等待 Flask 启动
|
||||
time.sleep(2)
|
||||
|
||||
|
||||
@@ -7,7 +7,7 @@
|
||||
<script type="module" crossorigin src="/assets/delete-brand.js"></script>
|
||||
<link rel="modulepreload" crossorigin href="/assets/pywebview-Dee3nQjE.js">
|
||||
<link rel="stylesheet" crossorigin href="/assets/pywebview-MkjZQlBu.css">
|
||||
<link rel="stylesheet" crossorigin href="/assets/delete-brand-CS44Bk-y.css">
|
||||
<link rel="stylesheet" crossorigin href="/assets/delete-brand-_jGFWTAn.css">
|
||||
</head>
|
||||
<body>
|
||||
<div id="app"></div>
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"version": "1.0.22"}
|
||||
{"version": "1.0.13"}
|
||||
1001
app/web_source/brand-旧.html
Normal file
1001
app/web_source/brand-旧.html
Normal file
File diff suppressed because it is too large
Load Diff
@@ -392,7 +392,7 @@
|
||||
<body>
|
||||
<header class="top-bar">
|
||||
<div class="logo-area">
|
||||
<span class="app-name">南日AI-亚马逊</span>
|
||||
<span class="app-name">数富AI-亚马逊</span>
|
||||
<a href="/home" class="btn-home">返回首页</a>
|
||||
</div>
|
||||
<nav class="nav-tabs">
|
||||
|
||||
BIN
backend-java/aiimage-backend.log.2026-04-02.0.gz
Normal file
BIN
backend-java/aiimage-backend.log.2026-04-02.0.gz
Normal file
Binary file not shown.
@@ -1053,6 +1053,7 @@ public class BrandTaskService {
|
||||
for (BrandCrawlTaskEntity task : runningTasks) {
|
||||
Map<Object, Object> progress = brandTaskProgressCacheService.getProgress(task.getId());
|
||||
if (progress.isEmpty()) {
|
||||
failStaleRunningTask(task.getId(), "任务进度缓存已过期,任务已自动失败");
|
||||
continue;
|
||||
}
|
||||
long lastHeartbeatAt = 0L;
|
||||
@@ -1061,21 +1062,26 @@ public class BrandTaskService {
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
if (lastHeartbeatAt <= 0) {
|
||||
failStaleRunningTask(task.getId(), "任务心跳信息缺失,任务已自动失败");
|
||||
continue;
|
||||
}
|
||||
LocalDateTime lastHeartbeat = LocalDateTime.ofInstant(Instant.ofEpochMilli(lastHeartbeatAt), ZoneId.systemDefault());
|
||||
if (lastHeartbeat.isAfter(threshold)) {
|
||||
continue;
|
||||
}
|
||||
brandCrawlTaskMapper.update(null, new LambdaUpdateWrapper<BrandCrawlTaskEntity>()
|
||||
.eq(BrandCrawlTaskEntity::getId, task.getId())
|
||||
.eq(BrandCrawlTaskEntity::getStatus, STATUS_RUNNING)
|
||||
.set(BrandCrawlTaskEntity::getStatus, STATUS_FAILED)
|
||||
.set(BrandCrawlTaskEntity::getErrorMessage, "前端长时间无响应,任务已自动失败"));
|
||||
brandTaskProgressCacheService.markFailed(task.getId(), "前端长时间无响应,任务已自动失败");
|
||||
failStaleRunningTask(task.getId(), "前端长时间无响应,任务已自动失败");
|
||||
}
|
||||
}
|
||||
|
||||
private void failStaleRunningTask(Long taskId, String message) {
|
||||
brandCrawlTaskMapper.update(null, new LambdaUpdateWrapper<BrandCrawlTaskEntity>()
|
||||
.eq(BrandCrawlTaskEntity::getId, taskId)
|
||||
.eq(BrandCrawlTaskEntity::getStatus, STATUS_RUNNING)
|
||||
.set(BrandCrawlTaskEntity::getStatus, STATUS_FAILED)
|
||||
.set(BrandCrawlTaskEntity::getErrorMessage, message));
|
||||
brandTaskProgressCacheService.markFailed(taskId, message);
|
||||
}
|
||||
|
||||
private record ParsedBrandFile(String sheetName, List<String> columns, List<Map<String, Object>> rows,
|
||||
List<String> uniqueBrands) {
|
||||
}
|
||||
|
||||
@@ -12,6 +12,9 @@ public class DeleteBrandParsedFileCacheDto {
|
||||
private String fileKey;
|
||||
private String sourceFilename;
|
||||
private String shopName;
|
||||
private String shopId;
|
||||
private String platform;
|
||||
private String openStoreUrl;
|
||||
private Integer totalRows;
|
||||
private List<DeleteBrandCountryGroupVo> countries = new ArrayList<>();
|
||||
private List<DeleteBrandPreviewRowVo> previewRows = new ArrayList<>();
|
||||
|
||||
@@ -17,4 +17,7 @@ public class DeleteBrandTaskDetailVo {
|
||||
|
||||
@Schema(description = "任务结果项快照")
|
||||
private List<DeleteBrandResultItemVo> items = new ArrayList<>();
|
||||
|
||||
@Schema(description = "文件级进度快照")
|
||||
private List<DeleteBrandTaskFileProgressVo> fileProgress = new ArrayList<>();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
package com.nanri.aiimage.modules.deletebrand.model.vo;
|
||||
|
||||
import io.swagger.v3.oas.annotations.media.Schema;
|
||||
import lombok.Data;
|
||||
|
||||
@Data
|
||||
@Schema(description = "删除品牌任务中的文件级进度")
|
||||
public class DeleteBrandTaskFileProgressVo {
|
||||
|
||||
@Schema(description = "文件标识(优先 fileKey)")
|
||||
private String fileKey;
|
||||
|
||||
@Schema(description = "源文件名")
|
||||
private String sourceFilename;
|
||||
|
||||
@Schema(description = "已处理行数")
|
||||
private Integer processedRows;
|
||||
|
||||
@Schema(description = "总行数")
|
||||
private Integer totalRows;
|
||||
|
||||
@Schema(description = "百分比(0-100)")
|
||||
private Integer percent;
|
||||
|
||||
@Schema(description = "文件状态:PENDING/RUNNING/COMPLETED/FAILED")
|
||||
private String status;
|
||||
}
|
||||
@@ -25,6 +25,7 @@ import com.nanri.aiimage.modules.deletebrand.model.vo.DeleteBrandResultItemVo;
|
||||
import com.nanri.aiimage.modules.deletebrand.model.vo.DeleteBrandRunVo;
|
||||
import com.nanri.aiimage.modules.deletebrand.model.vo.DeleteBrandTaskDeletionStatusVo;
|
||||
import com.nanri.aiimage.modules.deletebrand.model.vo.DeleteBrandTaskDetailVo;
|
||||
import com.nanri.aiimage.modules.deletebrand.model.vo.DeleteBrandTaskFileProgressVo;
|
||||
import com.nanri.aiimage.modules.deletebrand.model.vo.DeleteBrandTaskItemVo;
|
||||
import com.nanri.aiimage.modules.file.service.LocalFileStorageService;
|
||||
import com.nanri.aiimage.modules.file.service.oss.OssStorageService;
|
||||
@@ -151,6 +152,9 @@ public class DeleteBrandRunService {
|
||||
cacheDto.setFileKey(sourceFile.getFileKey());
|
||||
cacheDto.setSourceFilename(sourceFile.getOriginalFilename());
|
||||
cacheDto.setShopName(item.getShopName());
|
||||
cacheDto.setShopId(item.getShopId());
|
||||
cacheDto.setPlatform(item.getPlatform());
|
||||
cacheDto.setOpenStoreUrl(item.getOpenStoreUrl());
|
||||
cacheDto.setTotalRows(parsed.totalRows());
|
||||
cacheDto.setCountries(parsed.countries());
|
||||
cacheDto.setPreviewRows(parsed.previewRows());
|
||||
@@ -284,9 +288,9 @@ public class DeleteBrandRunService {
|
||||
item.setMatchStatus(candidate.getMatchStatus());
|
||||
item.setMatchMessage(candidate.getMatchMessage());
|
||||
// 如果物理结果已经是成功了,说明已经处理过了,强制置为匹配成功
|
||||
if (item.isSuccess()) {
|
||||
if (item.isSuccess() && candidate.getShopId() != null && !candidate.getShopId().isBlank()) {
|
||||
item.setMatched(true);
|
||||
item.setMatchStatus("MATCHED");
|
||||
item.setMatchStatus(ZiniaoShopIndexService.MATCH_STATUS_MATCHED);
|
||||
item.setMatchMessage(null);
|
||||
}
|
||||
|
||||
@@ -566,9 +570,11 @@ public class DeleteBrandRunService {
|
||||
DeleteBrandLineProgressVo progressVo = buildLineProgress(task.getId(), progress);
|
||||
|
||||
DeleteBrandTaskDetailVo detail = new DeleteBrandTaskDetailVo();
|
||||
List<DeleteBrandResultItemVo> taskItems = resolveTaskResultItems(task);
|
||||
detail.setTask(taskVo);
|
||||
detail.setLine_progress(progressVo);
|
||||
detail.setItems(resolveTaskResultItems(task));
|
||||
detail.setItems(taskItems);
|
||||
detail.setFileProgress(buildFileProgress(task, taskItems, progressVo));
|
||||
return detail;
|
||||
}
|
||||
|
||||
@@ -588,7 +594,7 @@ public class DeleteBrandRunService {
|
||||
}
|
||||
item.setTaskId(task.getId());
|
||||
item.setTaskStatus(task.getStatus());
|
||||
if ("SUCCESS".equals(task.getStatus()) && item.isSuccess()) {
|
||||
if ("SUCCESS".equals(task.getStatus()) && item.isSuccess() && item.getShopId() != null && !item.getShopId().isBlank()) {
|
||||
item.setMatched(true);
|
||||
item.setMatchStatus(ZiniaoShopIndexService.MATCH_STATUS_MATCHED);
|
||||
item.setMatchMessage(null);
|
||||
@@ -599,6 +605,73 @@ public class DeleteBrandRunService {
|
||||
}
|
||||
}
|
||||
|
||||
private List<DeleteBrandTaskFileProgressVo> buildFileProgress(FileTaskEntity task,
|
||||
List<DeleteBrandResultItemVo> taskItems,
|
||||
DeleteBrandLineProgressVo progressVo) {
|
||||
if (taskItems == null || taskItems.isEmpty()) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
String taskStatus = task == null ? null : task.getStatus();
|
||||
DeleteBrandLineProgressInfoVo info = progressVo == null ? null : progressVo.getInfo();
|
||||
Integer finishedFiles = info == null ? null : info.getFinished_files();
|
||||
String runningFileName = info == null ? null : info.getFile_name();
|
||||
Integer currentLine = info == null ? null : info.getCurrent_line();
|
||||
Integer totalLines = info == null ? null : info.getTotal_lines();
|
||||
|
||||
List<DeleteBrandTaskFileProgressVo> fileProgress = new ArrayList<>();
|
||||
int completedAssigned = 0;
|
||||
int normalizedFinished = finishedFiles == null ? 0 : Math.max(0, finishedFiles);
|
||||
|
||||
for (DeleteBrandResultItemVo item : taskItems) {
|
||||
if (item == null) {
|
||||
continue;
|
||||
}
|
||||
|
||||
DeleteBrandTaskFileProgressVo fp = new DeleteBrandTaskFileProgressVo();
|
||||
fp.setFileKey(item.getFileKey());
|
||||
fp.setSourceFilename(item.getSourceFilename());
|
||||
|
||||
Integer itemTotalRows = item.getTotalRows() == null ? 0 : Math.max(item.getTotalRows(), 0);
|
||||
fp.setTotalRows(itemTotalRows);
|
||||
|
||||
String status = "PENDING";
|
||||
int processedRows = 0;
|
||||
int percent = 0;
|
||||
|
||||
if ("SUCCESS".equals(taskStatus)) {
|
||||
status = "COMPLETED";
|
||||
processedRows = itemTotalRows;
|
||||
percent = 100;
|
||||
} else if ("FAILED".equals(taskStatus)) {
|
||||
status = "FAILED";
|
||||
processedRows = Math.min(itemTotalRows, Math.max(0, currentLine == null ? 0 : currentLine));
|
||||
percent = itemTotalRows > 0 ? Math.min(100, (int) Math.round(processedRows * 100.0 / itemTotalRows)) : 0;
|
||||
} else {
|
||||
if (completedAssigned < normalizedFinished) {
|
||||
status = "COMPLETED";
|
||||
processedRows = itemTotalRows;
|
||||
percent = 100;
|
||||
completedAssigned++;
|
||||
} else if (runningFileName != null && runningFileName.equals(item.getSourceFilename())) {
|
||||
status = "RUNNING";
|
||||
processedRows = Math.min(itemTotalRows, Math.max(0, currentLine == null ? 0 : currentLine));
|
||||
if (totalLines != null && totalLines > 0) {
|
||||
percent = Math.min(100, Math.max(0, (int) Math.round(processedRows * 100.0 / totalLines)));
|
||||
} else if (itemTotalRows > 0) {
|
||||
percent = Math.min(100, Math.max(0, (int) Math.round(processedRows * 100.0 / itemTotalRows)));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fp.setProcessedRows(processedRows);
|
||||
fp.setPercent(percent);
|
||||
fp.setStatus(status);
|
||||
fileProgress.add(fp);
|
||||
}
|
||||
return fileProgress;
|
||||
}
|
||||
|
||||
private FileTaskEntity refreshIndexedMatchesIfNeeded(FileTaskEntity task) {
|
||||
if (task == null || task.getResultJson() == null || task.getResultJson().isBlank()) {
|
||||
return task;
|
||||
@@ -1054,9 +1127,12 @@ public class DeleteBrandRunService {
|
||||
item.setCountries(parsedFile.getCountries());
|
||||
item.setPreviewRows(parsedFile.getPreviewRows());
|
||||
item.setTruncated(false);
|
||||
item.setMatched(true);
|
||||
item.setMatchStatus("MATCHED");
|
||||
item.setMatchMessage(null);
|
||||
item.setShopId(parsedFile.getShopId());
|
||||
item.setPlatform(parsedFile.getPlatform());
|
||||
item.setOpenStoreUrl(parsedFile.getOpenStoreUrl());
|
||||
item.setMatched(parsedFile.getShopId() != null && !parsedFile.getShopId().isBlank());
|
||||
item.setMatchStatus(item.isMatched() ? ZiniaoShopIndexService.MATCH_STATUS_MATCHED : ZiniaoShopIndexService.MATCH_STATUS_PENDING);
|
||||
item.setMatchMessage(item.isMatched() ? null : "店铺索引信息缺失,请重新匹配索引");
|
||||
finalItems.add(item);
|
||||
successCount++;
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@ import com.nanri.aiimage.modules.ziniao.memory.model.entity.ZiniaoMemoryStoreEnt
|
||||
import com.fasterxml.jackson.databind.JavaType;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
@@ -18,8 +19,12 @@ import java.util.Optional;
|
||||
|
||||
@Service
|
||||
@RequiredArgsConstructor
|
||||
@Slf4j
|
||||
public class ZiniaoMemoryStoreService {
|
||||
|
||||
/** 本表仅用于持久化店铺索引条目(见 ZiniaoShopIndexService)。 */
|
||||
public static final String CACHE_TYPE_SHOP_INDEX_ENTRY = "SHOP_INDEX_ENTRY";
|
||||
|
||||
private final ZiniaoMemoryStoreMapper ziniaoMemoryStoreMapper;
|
||||
private final ObjectMapper objectMapper;
|
||||
|
||||
@@ -80,6 +85,21 @@ public class ZiniaoMemoryStoreService {
|
||||
.toList();
|
||||
}
|
||||
|
||||
/**
|
||||
* 未过期行(含 cache_key),用于索引巡检、修复与去重后的 stale 标记。
|
||||
*/
|
||||
public List<ZiniaoMemoryStoreEntity> listAliveEntitiesByType(String cacheType, int limit) {
|
||||
String normalizedType = normalizeRequired(cacheType, "cacheType 不能为空");
|
||||
int safeLimit = Math.max(limit, 1);
|
||||
LocalDateTime now = LocalDateTime.now();
|
||||
List<ZiniaoMemoryStoreEntity> entities = ziniaoMemoryStoreMapper.selectList(new LambdaQueryWrapper<ZiniaoMemoryStoreEntity>()
|
||||
.eq(ZiniaoMemoryStoreEntity::getCacheType, normalizedType)
|
||||
.gt(ZiniaoMemoryStoreEntity::getExpiresAt, now)
|
||||
.orderByAsc(ZiniaoMemoryStoreEntity::getCacheKey)
|
||||
.last("LIMIT " + safeLimit));
|
||||
return entities == null ? List.of() : entities;
|
||||
}
|
||||
|
||||
@Transactional
|
||||
public void put(String cacheType, String cacheKey, Object payload, Duration ttl) {
|
||||
String normalizedType = normalizeRequired(cacheType, "cacheType 不能为空");
|
||||
@@ -105,19 +125,32 @@ public class ZiniaoMemoryStoreService {
|
||||
entity.setCreatedAt(now);
|
||||
entity.setUpdatedAt(now);
|
||||
ziniaoMemoryStoreMapper.insert(entity);
|
||||
if (CACHE_TYPE_SHOP_INDEX_ENTRY.equals(normalizedType)) {
|
||||
log.info("[ziniao-shop-index-db] insert cacheKey={} expiresAt={} bytes={}",
|
||||
normalizedKey, expiresAt, payloadJson.length());
|
||||
}
|
||||
return;
|
||||
}
|
||||
entity.setPayloadJson(payloadJson);
|
||||
entity.setExpiresAt(expiresAt);
|
||||
entity.setUpdatedAt(now);
|
||||
ziniaoMemoryStoreMapper.updateById(entity);
|
||||
if (CACHE_TYPE_SHOP_INDEX_ENTRY.equals(normalizedType)) {
|
||||
log.info("[ziniao-shop-index-db] update cacheKey={} id={} expiresAt={} bytes={}",
|
||||
normalizedKey, entity.getId(), expiresAt, payloadJson.length());
|
||||
}
|
||||
}
|
||||
|
||||
@Transactional
|
||||
public void delete(String cacheType, String cacheKey) {
|
||||
ZiniaoMemoryStoreEntity entity = findOne(cacheType, cacheKey);
|
||||
String normalizedType = normalizeRequired(cacheType, "cacheType 不能为空");
|
||||
String normalizedKey = normalizeRequired(cacheKey, "cacheKey 不能为空");
|
||||
ZiniaoMemoryStoreEntity entity = findOne(normalizedType, normalizedKey);
|
||||
if (entity != null) {
|
||||
ziniaoMemoryStoreMapper.deleteById(entity.getId());
|
||||
if (CACHE_TYPE_SHOP_INDEX_ENTRY.equals(normalizedType)) {
|
||||
log.info("[ziniao-shop-index-db] delete cacheKey={} id={}", entity.getCacheKey(), entity.getId());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,114 @@
|
||||
package com.nanri.aiimage.modules.ziniao.memory.service;
|
||||
|
||||
import com.fasterxml.jackson.databind.JavaType;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.nanri.aiimage.common.exception.BusinessException;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
/**
|
||||
* 进程内缓存:紫鸟 token/员工/店铺列表等短期数据。
|
||||
* 不写入 biz_ziniao_memory_store,避免与「店铺索引落库」混用同一张表。
|
||||
*/
|
||||
@Service
|
||||
@RequiredArgsConstructor
|
||||
@Slf4j
|
||||
public class ZiniaoTransientCacheService {
|
||||
|
||||
private static final String SEP = "::";
|
||||
|
||||
private final ObjectMapper objectMapper;
|
||||
|
||||
private final ConcurrentHashMap<String, Holder> map = new ConcurrentHashMap<>();
|
||||
|
||||
public <T> Optional<T> get(String cacheType, String cacheKey, Class<T> valueType) {
|
||||
Holder holder = getHolder(cacheType, cacheKey);
|
||||
if (holder == null) {
|
||||
return Optional.empty();
|
||||
}
|
||||
try {
|
||||
return Optional.ofNullable(objectMapper.readValue(holder.json, valueType));
|
||||
} catch (Exception ex) {
|
||||
throw new BusinessException("读取紫鸟进程缓存失败");
|
||||
}
|
||||
}
|
||||
|
||||
public <T> Optional<T> get(String cacheType, String cacheKey, JavaType javaType) {
|
||||
Holder holder = getHolder(cacheType, cacheKey);
|
||||
if (holder == null) {
|
||||
return Optional.empty();
|
||||
}
|
||||
try {
|
||||
@SuppressWarnings("unchecked")
|
||||
T value = (T) objectMapper.readValue(holder.json, javaType);
|
||||
return Optional.ofNullable(value);
|
||||
} catch (Exception ex) {
|
||||
throw new BusinessException("读取紫鸟进程缓存失败");
|
||||
}
|
||||
}
|
||||
|
||||
public <E> Optional<List<E>> getList(String cacheType, String cacheKey, Class<E> elementType) {
|
||||
JavaType type = objectMapper.getTypeFactory().constructCollectionType(List.class, elementType);
|
||||
return get(cacheType, cacheKey, type);
|
||||
}
|
||||
|
||||
public void put(String cacheType, String cacheKey, Object payload, Duration ttl) {
|
||||
String normalizedType = normalizeRequired(cacheType, "cacheType 不能为空");
|
||||
String normalizedKey = normalizeRequired(cacheKey, "cacheKey 不能为空");
|
||||
if (ttl == null || ttl.isZero() || ttl.isNegative()) {
|
||||
throw new BusinessException("ttl 不合法");
|
||||
}
|
||||
try {
|
||||
String json = objectMapper.writeValueAsString(payload);
|
||||
LocalDateTime expiresAt = LocalDateTime.now().plusSeconds(ttl.getSeconds());
|
||||
map.put(compoundKey(normalizedType, normalizedKey), new Holder(json, expiresAt));
|
||||
log.trace("[ziniao-transient] put type={} keyLen={} ttlSec={}", normalizedType, normalizedKey.length(), ttl.getSeconds());
|
||||
} catch (Exception ex) {
|
||||
throw new BusinessException("写入紫鸟进程缓存失败");
|
||||
}
|
||||
}
|
||||
|
||||
public void delete(String cacheType, String cacheKey) {
|
||||
String normalizedType = normalizeRequired(cacheType, "cacheType 不能为空");
|
||||
String normalizedKey = normalizeRequired(cacheKey, "cacheKey 不能为空");
|
||||
map.remove(compoundKey(normalizedType, normalizedKey));
|
||||
log.trace("[ziniao-transient] delete type={}", normalizedType);
|
||||
}
|
||||
|
||||
private Holder getHolder(String cacheType, String cacheKey) {
|
||||
String normalizedType = normalizeRequired(cacheType, "cacheType 不能为空");
|
||||
String normalizedKey = normalizeRequired(cacheKey, "cacheKey 不能为空");
|
||||
Holder holder = map.get(compoundKey(normalizedType, normalizedKey));
|
||||
if (holder == null) {
|
||||
return null;
|
||||
}
|
||||
if (holder.expiresAt == null || !holder.expiresAt.isAfter(LocalDateTime.now())) {
|
||||
map.remove(compoundKey(normalizedType, normalizedKey));
|
||||
return null;
|
||||
}
|
||||
return holder;
|
||||
}
|
||||
|
||||
private static String compoundKey(String cacheType, String cacheKey) {
|
||||
return cacheType + SEP + cacheKey;
|
||||
}
|
||||
|
||||
private static String normalizeRequired(String value, String message) {
|
||||
String normalized = Objects.toString(value, "").trim();
|
||||
if (normalized.isEmpty()) {
|
||||
throw new BusinessException(message);
|
||||
}
|
||||
return normalized;
|
||||
}
|
||||
|
||||
private record Holder(String json, LocalDateTime expiresAt) {
|
||||
}
|
||||
}
|
||||
@@ -3,7 +3,7 @@ package com.nanri.aiimage.modules.ziniao.service;
|
||||
import com.nanri.aiimage.common.exception.BusinessException;
|
||||
import com.nanri.aiimage.config.ZiniaoProperties;
|
||||
import com.nanri.aiimage.modules.ziniao.client.ZiniaoClient;
|
||||
import com.nanri.aiimage.modules.ziniao.memory.service.ZiniaoMemoryStoreService;
|
||||
import com.nanri.aiimage.modules.ziniao.memory.service.ZiniaoTransientCacheService;
|
||||
import com.nanri.aiimage.modules.ziniao.model.cache.ZiniaoSessionCacheDto;
|
||||
import com.nanri.aiimage.modules.ziniao.model.cache.ZiniaoShopCacheDto;
|
||||
import com.nanri.aiimage.modules.ziniao.service.ZiniaoApiKeyProvider;
|
||||
@@ -151,23 +151,30 @@ public class ZiniaoAuthService {
|
||||
|| message.contains("userId存在无效的参数值");
|
||||
}
|
||||
|
||||
public void evictUserStoresCacheForIndex(String apiKey, Long companyId, Long userId) {
|
||||
if (apiKey == null || apiKey.isBlank() || companyId == null || companyId <= 0 || userId == null || userId <= 0) {
|
||||
return;
|
||||
}
|
||||
ziniaoTransientCacheService.delete(CACHE_TYPE_USER_STORES, buildApiKeyHash(apiKey) + ":" + companyId + ":" + userId);
|
||||
}
|
||||
|
||||
public void evictInvalidUserForIndex(String apiKey, Long companyId, Long userId) {
|
||||
if (apiKey == null || apiKey.isBlank() || companyId == null || companyId <= 0 || userId == null || userId <= 0) {
|
||||
return;
|
||||
}
|
||||
String staffCacheKey = buildApiKeyHash(apiKey) + ":" + companyId;
|
||||
ziniaoMemoryStoreService.getList(CACHE_TYPE_STAFF_LIST, staffCacheKey, ZiniaoStaffItemVo.class)
|
||||
ziniaoTransientCacheService.getList(CACHE_TYPE_STAFF_LIST, staffCacheKey, ZiniaoStaffItemVo.class)
|
||||
.ifPresent(staff -> {
|
||||
List<ZiniaoStaffItemVo> filtered = staff.stream()
|
||||
.filter(item -> item != null && !userId.equals(item.getUserId()))
|
||||
.toList();
|
||||
if (filtered.isEmpty()) {
|
||||
ziniaoMemoryStoreService.delete(CACHE_TYPE_STAFF_LIST, staffCacheKey);
|
||||
ziniaoTransientCacheService.delete(CACHE_TYPE_STAFF_LIST, staffCacheKey);
|
||||
} else {
|
||||
ziniaoMemoryStoreService.put(CACHE_TYPE_STAFF_LIST, staffCacheKey, filtered, STAFF_LIST_CACHE_TTL);
|
||||
ziniaoTransientCacheService.put(CACHE_TYPE_STAFF_LIST, staffCacheKey, filtered, STAFF_LIST_CACHE_TTL);
|
||||
}
|
||||
});
|
||||
ziniaoMemoryStoreService.delete(CACHE_TYPE_USER_STORES, buildApiKeyHash(apiKey) + ":" + companyId + ":" + userId);
|
||||
ziniaoTransientCacheService.delete(CACHE_TYPE_USER_STORES, buildApiKeyHash(apiKey) + ":" + companyId + ":" + userId);
|
||||
}
|
||||
|
||||
private boolean isSkippableUserStoresError(BusinessException ex) {
|
||||
@@ -196,7 +203,7 @@ public class ZiniaoAuthService {
|
||||
private final ZiniaoProperties ziniaoProperties;
|
||||
private final ZiniaoClient ziniaoClient;
|
||||
private final ZiniaoSessionCacheService ziniaoSessionCacheService;
|
||||
private final ZiniaoMemoryStoreService ziniaoMemoryStoreService;
|
||||
private final ZiniaoTransientCacheService ziniaoTransientCacheService;
|
||||
private final ZiniaoApiKeyProvider ziniaoApiKeyProvider;
|
||||
|
||||
public ZiniaoSessionVo getSession(String sessionId, Long userId) {
|
||||
@@ -393,14 +400,14 @@ public class ZiniaoAuthService {
|
||||
|
||||
String normalizedApiKey = requireText(apiKey, "紫鸟 apiKey 未配置");
|
||||
String cacheKey = buildApiKeyHash(normalizedApiKey);
|
||||
Long cached = ziniaoMemoryStoreService.get(CACHE_TYPE_COMPANY_ID, cacheKey, Long.class).orElse(null);
|
||||
Long cached = ziniaoTransientCacheService.get(CACHE_TYPE_COMPANY_ID, cacheKey, Long.class).orElse(null);
|
||||
if (cached != null && cached > 0) {
|
||||
log.info("[ziniao-company] hit cache, keyHash={}, companyId={}", shortKeyHash(apiKey), cached);
|
||||
return cached;
|
||||
}
|
||||
log.info("[ziniao-company] cache miss, requesting upstream, keyHash={}", shortKeyHash(apiKey));
|
||||
Long companyId = ziniaoClient.getCompanyIdByApiKey(normalizedApiKey);
|
||||
ziniaoMemoryStoreService.put(CACHE_TYPE_COMPANY_ID, cacheKey, companyId, COMPANY_ID_CACHE_TTL);
|
||||
ziniaoTransientCacheService.put(CACHE_TYPE_COMPANY_ID, cacheKey, companyId, COMPANY_ID_CACHE_TTL);
|
||||
log.info("[ziniao-company] cached upstream result, keyHash={}, companyId={}", shortKeyHash(apiKey), companyId);
|
||||
return companyId;
|
||||
}
|
||||
@@ -408,7 +415,7 @@ public class ZiniaoAuthService {
|
||||
|
||||
private List<ZiniaoStaffItemVo> getOrLoadStaff(String apiKey, Long companyId) {
|
||||
String cacheKey = buildApiKeyHash(apiKey) + ":" + companyId;
|
||||
List<ZiniaoStaffItemVo> cached = ziniaoMemoryStoreService.getList(CACHE_TYPE_STAFF_LIST, cacheKey, ZiniaoStaffItemVo.class)
|
||||
List<ZiniaoStaffItemVo> cached = ziniaoTransientCacheService.getList(CACHE_TYPE_STAFF_LIST, cacheKey, ZiniaoStaffItemVo.class)
|
||||
.orElse(null);
|
||||
if (cached != null) {
|
||||
log.info("[ziniao-staff] hit cache, keyHash={}, companyId={}, size={}", shortKeyHash(apiKey), companyId, cached.size());
|
||||
@@ -416,28 +423,43 @@ public class ZiniaoAuthService {
|
||||
}
|
||||
log.info("[ziniao-staff] cache miss, requesting upstream, keyHash={}, companyId={}", shortKeyHash(apiKey), companyId);
|
||||
List<ZiniaoStaffItemVo> staff = ziniaoClient.listStaff(apiKey, companyId);
|
||||
ziniaoMemoryStoreService.put(CACHE_TYPE_STAFF_LIST, cacheKey, staff, STAFF_LIST_CACHE_TTL);
|
||||
ziniaoTransientCacheService.put(CACHE_TYPE_STAFF_LIST, cacheKey, staff, STAFF_LIST_CACHE_TTL);
|
||||
log.info("[ziniao-staff] cached upstream result, keyHash={}, companyId={}, size={}", shortKeyHash(apiKey), companyId, staff.size());
|
||||
return staff;
|
||||
}
|
||||
|
||||
private List<ZiniaoShopCacheDto> getOrLoadUserStores(String apiKey, Long companyId, Long userId) {
|
||||
String cacheKey = buildApiKeyHash(apiKey) + ":" + companyId + ":" + userId;
|
||||
List<ZiniaoShopCacheDto> cached = ziniaoMemoryStoreService.getList(CACHE_TYPE_USER_STORES, cacheKey, ZiniaoShopCacheDto.class)
|
||||
List<ZiniaoShopCacheDto> cached = ziniaoTransientCacheService.getList(CACHE_TYPE_USER_STORES, cacheKey, ZiniaoShopCacheDto.class)
|
||||
.orElse(null);
|
||||
if (cached != null) {
|
||||
log.info("[ziniao-stores] hit cache, keyHash={}, companyId={}, userId={}, size={}", shortKeyHash(apiKey), companyId, userId, cached.size());
|
||||
return cached;
|
||||
}
|
||||
log.info("[ziniao-stores] cache miss, requesting upstream, keyHash={}, companyId={}, userId={}", shortKeyHash(apiKey), companyId, userId);
|
||||
List<ZiniaoShopCacheDto> stores = ziniaoClient.listUserStores(apiKey, companyId, userId);
|
||||
ziniaoMemoryStoreService.put(CACHE_TYPE_USER_STORES, cacheKey, stores, USER_STORES_CACHE_TTL);
|
||||
List<ZiniaoShopCacheDto> stores = ziniaoClient.listUserStores(apiKey, companyId, userId).stream()
|
||||
.filter(item -> item != null && item.getShopId() != null && !item.getShopId().isBlank())
|
||||
.collect(java.util.stream.Collectors.collectingAndThen(
|
||||
java.util.stream.Collectors.toMap(
|
||||
ZiniaoShopCacheDto::getShopId,
|
||||
item -> item,
|
||||
(left, right) -> left,
|
||||
java.util.LinkedHashMap::new
|
||||
),
|
||||
map -> new java.util.ArrayList<>(map.values())
|
||||
));
|
||||
if (stores.isEmpty()) {
|
||||
ziniaoTransientCacheService.delete(CACHE_TYPE_USER_STORES, cacheKey);
|
||||
log.info("[ziniao-stores] upstream returned empty stores, skip cache, keyHash={}, companyId={}, userId={}", shortKeyHash(apiKey), companyId, userId);
|
||||
return stores;
|
||||
}
|
||||
ziniaoTransientCacheService.put(CACHE_TYPE_USER_STORES, cacheKey, stores, USER_STORES_CACHE_TTL);
|
||||
log.info("[ziniao-stores] cached upstream result, keyHash={}, companyId={}, userId={}, size={}", shortKeyHash(apiKey), companyId, userId, stores.size());
|
||||
return stores;
|
||||
}
|
||||
|
||||
private StoreMatchResult getCachedShopMatch(String apiKey, Long companyId, String normalizedShopName) {
|
||||
StoreMatchResult cached = ziniaoMemoryStoreService.get(CACHE_TYPE_SHOP_MATCH, buildShopMatchCacheKey(apiKey, companyId, normalizedShopName), StoreMatchResult.class)
|
||||
StoreMatchResult cached = ziniaoTransientCacheService.get(CACHE_TYPE_SHOP_MATCH, buildShopMatchCacheKey(apiKey, companyId, normalizedShopName), StoreMatchResult.class)
|
||||
.orElse(null);
|
||||
if (cached == null || !cached.matched() || cached.shopId() == null || cached.matchedUserId() == null) {
|
||||
return null;
|
||||
@@ -446,7 +468,7 @@ public class ZiniaoAuthService {
|
||||
}
|
||||
|
||||
private void cacheShopMatch(String apiKey, Long companyId, String normalizedShopName, StoreMatchResult result, Duration ttl) {
|
||||
ziniaoMemoryStoreService.put(CACHE_TYPE_SHOP_MATCH, buildShopMatchCacheKey(apiKey, companyId, normalizedShopName), result, ttl);
|
||||
ziniaoTransientCacheService.put(CACHE_TYPE_SHOP_MATCH, buildShopMatchCacheKey(apiKey, companyId, normalizedShopName), result, ttl);
|
||||
}
|
||||
|
||||
private String buildShopMatchCacheKey(String apiKey, Long companyId, String normalizedShopName) {
|
||||
|
||||
@@ -1,8 +1,11 @@
|
||||
package com.nanri.aiimage.modules.ziniao.service;
|
||||
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.nanri.aiimage.common.exception.BusinessException;
|
||||
import com.nanri.aiimage.config.ZiniaoProperties;
|
||||
import com.nanri.aiimage.modules.ziniao.memory.model.entity.ZiniaoMemoryStoreEntity;
|
||||
import com.nanri.aiimage.modules.ziniao.memory.service.ZiniaoMemoryStoreService;
|
||||
import com.nanri.aiimage.modules.ziniao.memory.service.ZiniaoTransientCacheService;
|
||||
import com.nanri.aiimage.modules.ziniao.model.cache.ZiniaoShopIndexEntryDto;
|
||||
import com.nanri.aiimage.modules.ziniao.model.cache.ZiniaoShopIndexRefreshCursorDto;
|
||||
import com.nanri.aiimage.modules.ziniao.model.cache.ZiniaoShopCacheDto;
|
||||
@@ -17,6 +20,7 @@ import java.security.MessageDigest;
|
||||
import java.time.Duration;
|
||||
import java.time.Instant;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashSet;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
@@ -38,21 +42,26 @@ public class ZiniaoShopIndexService {
|
||||
public static final String MATCH_STATUS_CONFLICT = "CONFLICT";
|
||||
public static final String MATCH_STATUS_STALE = "INDEX_STALE";
|
||||
|
||||
private static final String CACHE_TYPE_SHOP_INDEX_ENTRY = "SHOP_INDEX_ENTRY";
|
||||
private static final String CACHE_TYPE_SHOP_INDEX_SCOPE_SNAPSHOT = "SHOP_INDEX_SCOPE_SNAPSHOT";
|
||||
private static final String CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR = "SHOP_INDEX_REFRESH_CURSOR";
|
||||
/** 有 shopId 时唯一键,避免同一店铺因不同员工/哈希产生多行。 */
|
||||
private static final String SHOP_ENTRY_KEY_SHOP_PREFIX = "s:";
|
||||
/** 无 shopId(冲突占位等)时仍按规范化店名存一行。 */
|
||||
private static final String SHOP_ENTRY_KEY_NAME_PREFIX = "n:";
|
||||
private static final Duration DEFAULT_ENTRY_TTL = Duration.ofHours(12);
|
||||
private static final Duration DEFAULT_CURSOR_TTL = Duration.ofHours(24);
|
||||
private static final int DEFAULT_REFRESH_BATCH_SIZE = 100;
|
||||
private static final int DEFAULT_LIST_BY_TYPE_LIMIT = 5000;
|
||||
private static final int SHOP_INDEX_LIST_LIMIT = 10000;
|
||||
|
||||
private final ZiniaoMemoryStoreService ziniaoMemoryStoreService;
|
||||
private final ZiniaoTransientCacheService ziniaoTransientCacheService;
|
||||
private final ZiniaoApiKeyProvider ziniaoApiKeyProvider;
|
||||
private final ZiniaoAuthService ziniaoAuthService;
|
||||
private final ZiniaoProperties ziniaoProperties;
|
||||
private final ObjectMapper objectMapper;
|
||||
|
||||
public ZiniaoShopIndexRefreshCursorDto getRefreshCursor() {
|
||||
ZiniaoShopIndexRefreshCursorDto cursor = ziniaoMemoryStoreService
|
||||
ZiniaoShopIndexRefreshCursorDto cursor = ziniaoTransientCacheService
|
||||
.get(CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR, "global", ZiniaoShopIndexRefreshCursorDto.class)
|
||||
.orElse(null);
|
||||
if (cursor == null) {
|
||||
@@ -70,9 +79,7 @@ public class ZiniaoShopIndexService {
|
||||
return pendingResult("店铺名为空,无法匹配索引");
|
||||
}
|
||||
|
||||
ZiniaoShopIndexEntryDto entry = ziniaoMemoryStoreService
|
||||
.get(CACHE_TYPE_SHOP_INDEX_ENTRY, normalizedShopName, ZiniaoShopIndexEntryDto.class)
|
||||
.orElse(null);
|
||||
ZiniaoShopIndexEntryDto entry = findShopIndexEntryForLookup(normalizedShopName);
|
||||
if (entry == null) {
|
||||
return pendingResult("店铺索引未命中,请等待后台刷新");
|
||||
}
|
||||
@@ -119,9 +126,11 @@ public class ZiniaoShopIndexService {
|
||||
cursor.setScopeKey("global");
|
||||
cursor.setStatus("RUNNING");
|
||||
cursor.setLastStartedAt(now);
|
||||
ziniaoMemoryStoreService.put(CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR, "global", cursor, DEFAULT_CURSOR_TTL);
|
||||
ziniaoTransientCacheService.put(CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR, "global", cursor, DEFAULT_CURSOR_TTL);
|
||||
log.info("[ziniao-index] refresh started");
|
||||
|
||||
Map<String, List<ZiniaoShopIndexEntryDto>> grouped = new LinkedHashMap<>();
|
||||
Map<String, Long> storesFingerprintToUserId = new LinkedHashMap<>();
|
||||
Set<Long> allInvalidUserIds = new LinkedHashSet<>();
|
||||
int refreshBatchSize = resolveRefreshBatchSize();
|
||||
try {
|
||||
@@ -156,7 +165,25 @@ public class ZiniaoShopIndexService {
|
||||
log.warn("[ziniao-index] skip user stores, companyId={}, userId={}, msg={}", companyId, userId, ex.getMessage());
|
||||
continue;
|
||||
}
|
||||
ziniaoMemoryStoreService.put(
|
||||
if (stores == null || stores.isEmpty()) {
|
||||
ziniaoTransientCacheService.delete(
|
||||
CACHE_TYPE_SHOP_INDEX_SCOPE_SNAPSHOT,
|
||||
buildScopeSnapshotKey(apiKey, companyId, userId)
|
||||
);
|
||||
continue;
|
||||
}
|
||||
String storesFingerprint = buildStoresFingerprint(stores);
|
||||
Long duplicatedUserId = storesFingerprintToUserId.putIfAbsent(storesFingerprint, userId);
|
||||
if (duplicatedUserId != null) {
|
||||
ziniaoTransientCacheService.delete(
|
||||
CACHE_TYPE_SHOP_INDEX_SCOPE_SNAPSHOT,
|
||||
buildScopeSnapshotKey(apiKey, companyId, userId)
|
||||
);
|
||||
ziniaoAuthService.evictUserStoresCacheForIndex(apiKey, companyId, userId);
|
||||
log.info("[ziniao-index] skip duplicated user stores snapshot, companyId={}, userId={}, duplicatedUserId={}", companyId, userId, duplicatedUserId);
|
||||
continue;
|
||||
}
|
||||
ziniaoTransientCacheService.put(
|
||||
CACHE_TYPE_SHOP_INDEX_SCOPE_SNAPSHOT,
|
||||
buildScopeSnapshotKey(apiKey, companyId, userId),
|
||||
stores,
|
||||
@@ -186,7 +213,7 @@ public class ZiniaoShopIndexService {
|
||||
}
|
||||
if (!invalidUserIds.isEmpty()) {
|
||||
for (Long invalidUserId : invalidUserIds) {
|
||||
ziniaoMemoryStoreService.delete(
|
||||
ziniaoTransientCacheService.delete(
|
||||
CACHE_TYPE_SHOP_INDEX_SCOPE_SNAPSHOT,
|
||||
buildScopeSnapshotKey(apiKey, companyId, invalidUserId)
|
||||
);
|
||||
@@ -194,6 +221,7 @@ public class ZiniaoShopIndexService {
|
||||
}
|
||||
}
|
||||
|
||||
List<ZiniaoShopIndexEntryDto> roundEntries = new ArrayList<>();
|
||||
for (Map.Entry<String, List<ZiniaoShopIndexEntryDto>> groupedEntry : grouped.entrySet()) {
|
||||
String normalizedShopName = groupedEntry.getKey();
|
||||
List<ZiniaoShopIndexEntryDto> candidates = groupedEntry.getValue();
|
||||
@@ -201,9 +229,32 @@ public class ZiniaoShopIndexService {
|
||||
if (candidates.size() > 1) {
|
||||
entryToStore = buildConflictEntry(normalizedShopName, candidates, now);
|
||||
}
|
||||
ziniaoMemoryStoreService.put(CACHE_TYPE_SHOP_INDEX_ENTRY, normalizedShopName, entryToStore, resolveEntryTtl());
|
||||
roundEntries.add(entryToStore);
|
||||
}
|
||||
markMissingEntriesAsStale(grouped.keySet(), now);
|
||||
Map<String, ZiniaoShopIndexEntryDto> byShopId = new LinkedHashMap<>();
|
||||
Map<String, ZiniaoShopIndexEntryDto> byNameOnly = new LinkedHashMap<>();
|
||||
for (ZiniaoShopIndexEntryDto e : roundEntries) {
|
||||
String sid = e.getShopId();
|
||||
if (sid != null && !sid.isBlank()) {
|
||||
byShopId.merge(sid, e, this::mergeDuplicateShopIndexEntries);
|
||||
} else if (e.getNormalizedShopName() != null && !e.getNormalizedShopName().isBlank()) {
|
||||
byNameOnly.merge(e.getNormalizedShopName(), e, this::mergeDuplicateShopIndexEntries);
|
||||
}
|
||||
}
|
||||
Set<String> activeCacheKeys = new HashSet<>();
|
||||
for (ZiniaoShopIndexEntryDto e : byShopId.values()) {
|
||||
String key = SHOP_ENTRY_KEY_SHOP_PREFIX + e.getShopId();
|
||||
activeCacheKeys.add(key);
|
||||
ziniaoMemoryStoreService.put(ZiniaoMemoryStoreService.CACHE_TYPE_SHOP_INDEX_ENTRY, key, e, resolveEntryTtl());
|
||||
}
|
||||
for (ZiniaoShopIndexEntryDto e : byNameOnly.values()) {
|
||||
String key = SHOP_ENTRY_KEY_NAME_PREFIX + e.getNormalizedShopName();
|
||||
activeCacheKeys.add(key);
|
||||
ziniaoMemoryStoreService.put(ZiniaoMemoryStoreService.CACHE_TYPE_SHOP_INDEX_ENTRY, key, e, resolveEntryTtl());
|
||||
}
|
||||
log.info("[ziniao-index] refresh persist shopIndex uniqueShopId={} nameOnly={} groupedNames={} activeRows={}",
|
||||
byShopId.size(), byNameOnly.size(), grouped.size(), activeCacheKeys.size());
|
||||
markMissingEntriesAsStale(activeCacheKeys, now);
|
||||
|
||||
cursor.setStatus("SUCCESS");
|
||||
cursor.setMessage(allInvalidUserIds.isEmpty()
|
||||
@@ -213,12 +264,14 @@ public class ZiniaoShopIndexService {
|
||||
cursor.setSampleInvalidUserIds(allInvalidUserIds.stream().limit(20).toList());
|
||||
cursor.setLastFinishedAt(now);
|
||||
cursor.setLastSuccessAt(now);
|
||||
ziniaoMemoryStoreService.put(CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR, "global", cursor, DEFAULT_CURSOR_TTL);
|
||||
ziniaoTransientCacheService.put(CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR, "global", cursor, DEFAULT_CURSOR_TTL);
|
||||
log.info("[ziniao-index] refresh success invalidUsersSkipped={}", allInvalidUserIds.size());
|
||||
} catch (Exception ex) {
|
||||
cursor.setStatus("FAILED");
|
||||
cursor.setMessage(ex.getMessage());
|
||||
cursor.setLastFinishedAt(now);
|
||||
ziniaoMemoryStoreService.put(CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR, "global", cursor, DEFAULT_CURSOR_TTL);
|
||||
ziniaoTransientCacheService.put(CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR, "global", cursor, DEFAULT_CURSOR_TTL);
|
||||
log.warn("[ziniao-index] refresh failed status=FAILED msg={}", ex.getMessage());
|
||||
if (ex instanceof BusinessException businessException) {
|
||||
throw businessException;
|
||||
}
|
||||
@@ -227,7 +280,8 @@ public class ZiniaoShopIndexService {
|
||||
}
|
||||
|
||||
public void invalidateIndex() {
|
||||
ziniaoMemoryStoreService.delete(CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR, "global");
|
||||
ziniaoTransientCacheService.delete(CACHE_TYPE_SHOP_INDEX_REFRESH_CURSOR, "global");
|
||||
log.info("[ziniao-index] cursor invalidated (transient only; shop rows unchanged)");
|
||||
}
|
||||
|
||||
public String normalizeShopName(String value) {
|
||||
@@ -292,20 +346,85 @@ public class ZiniaoShopIndexService {
|
||||
return pendingResult(defaultMessage);
|
||||
}
|
||||
|
||||
private void markMissingEntriesAsStale(Set<String> activeNames, long now) {
|
||||
List<ZiniaoShopIndexEntryDto> existingEntries = ziniaoMemoryStoreService.listByType(
|
||||
CACHE_TYPE_SHOP_INDEX_ENTRY,
|
||||
ZiniaoShopIndexEntryDto.class,
|
||||
DEFAULT_LIST_BY_TYPE_LIMIT
|
||||
);
|
||||
if (existingEntries.isEmpty()) {
|
||||
return;
|
||||
private ZiniaoShopIndexEntryDto findShopIndexEntryForLookup(String normalizedShopName) {
|
||||
String prefixedKey = SHOP_ENTRY_KEY_NAME_PREFIX + normalizedShopName;
|
||||
ZiniaoShopIndexEntryDto entry = ziniaoMemoryStoreService
|
||||
.get(ZiniaoMemoryStoreService.CACHE_TYPE_SHOP_INDEX_ENTRY, prefixedKey, ZiniaoShopIndexEntryDto.class)
|
||||
.orElse(null);
|
||||
if (entry != null) {
|
||||
log.debug("[ziniao-shop-index] lookup hit=prefixedName cacheKey={}", prefixedKey);
|
||||
return entry;
|
||||
}
|
||||
for (ZiniaoShopIndexEntryDto existingEntry : existingEntries) {
|
||||
if (existingEntry == null || existingEntry.getNormalizedShopName() == null || existingEntry.getNormalizedShopName().isBlank()) {
|
||||
entry = ziniaoMemoryStoreService
|
||||
.get(ZiniaoMemoryStoreService.CACHE_TYPE_SHOP_INDEX_ENTRY, normalizedShopName, ZiniaoShopIndexEntryDto.class)
|
||||
.orElse(null);
|
||||
if (entry != null) {
|
||||
log.info("[ziniao-shop-index] lookup hit=legacyKey please migrate cacheKey={}", normalizedShopName);
|
||||
return entry;
|
||||
}
|
||||
List<ZiniaoShopIndexEntryDto> rows = ziniaoMemoryStoreService.listByType(
|
||||
ZiniaoMemoryStoreService.CACHE_TYPE_SHOP_INDEX_ENTRY,
|
||||
ZiniaoShopIndexEntryDto.class,
|
||||
SHOP_INDEX_LIST_LIMIT
|
||||
);
|
||||
for (ZiniaoShopIndexEntryDto row : rows) {
|
||||
if (row == null) {
|
||||
continue;
|
||||
}
|
||||
if (activeNames.contains(existingEntry.getNormalizedShopName())) {
|
||||
if (normalizedShopName.equals(row.getNormalizedShopName())
|
||||
|| normalizedShopName.equals(normalizeShopName(row.getShopName()))) {
|
||||
log.info("[ziniao-shop-index] lookup hit=scanByShopName shopId={} normalizedInRow={}",
|
||||
row.getShopId(), row.getNormalizedShopName());
|
||||
return row;
|
||||
}
|
||||
}
|
||||
log.info("[ziniao-shop-index] lookup miss normalizedName={}", normalizedShopName);
|
||||
return null;
|
||||
}
|
||||
|
||||
private ZiniaoShopIndexEntryDto mergeDuplicateShopIndexEntries(ZiniaoShopIndexEntryDto a, ZiniaoShopIndexEntryDto b) {
|
||||
if (a == null) {
|
||||
return b;
|
||||
}
|
||||
if (b == null) {
|
||||
return a;
|
||||
}
|
||||
if (STATUS_CONFLICT.equals(a.getStatus()) != STATUS_CONFLICT.equals(b.getStatus())) {
|
||||
return STATUS_CONFLICT.equals(a.getStatus()) ? a : b;
|
||||
}
|
||||
if (STATUS_CONFLICT.equals(a.getStatus())) {
|
||||
int ca = a.getCandidateCount() == null ? 0 : a.getCandidateCount();
|
||||
int cb = b.getCandidateCount() == null ? 0 : b.getCandidateCount();
|
||||
return cb > ca ? b : a;
|
||||
}
|
||||
long ta = a.getLastRefreshedAt() != null ? a.getLastRefreshedAt() : (a.getLastSeenAt() != null ? a.getLastSeenAt() : 0L);
|
||||
long tb = b.getLastRefreshedAt() != null ? b.getLastRefreshedAt() : (b.getLastSeenAt() != null ? b.getLastSeenAt() : 0L);
|
||||
return tb >= ta ? b : a;
|
||||
}
|
||||
|
||||
private void markMissingEntriesAsStale(Set<String> activeCacheKeys, long now) {
|
||||
List<ZiniaoMemoryStoreEntity> entities = ziniaoMemoryStoreService.listAliveEntitiesByType(
|
||||
ZiniaoMemoryStoreService.CACHE_TYPE_SHOP_INDEX_ENTRY,
|
||||
SHOP_INDEX_LIST_LIMIT
|
||||
);
|
||||
if (entities.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
int staleCount = 0;
|
||||
for (ZiniaoMemoryStoreEntity entity : entities) {
|
||||
String rowKey = entity.getCacheKey();
|
||||
if (activeCacheKeys.contains(rowKey)) {
|
||||
continue;
|
||||
}
|
||||
ZiniaoShopIndexEntryDto existingEntry;
|
||||
try {
|
||||
existingEntry = objectMapper.readValue(entity.getPayloadJson(), ZiniaoShopIndexEntryDto.class);
|
||||
} catch (Exception ex) {
|
||||
log.warn("[ziniao-index] skip corrupt shop_index row cacheKey={}", rowKey);
|
||||
continue;
|
||||
}
|
||||
if (existingEntry == null || existingEntry.getNormalizedShopName() == null
|
||||
|| existingEntry.getNormalizedShopName().isBlank()) {
|
||||
continue;
|
||||
}
|
||||
if (STATUS_CONFLICT.equals(existingEntry.getStatus())
|
||||
@@ -313,18 +432,18 @@ public class ZiniaoShopIndexService {
|
||||
&& (existingEntry.getShopName() == null || existingEntry.getShopName().isBlank())
|
||||
&& (existingEntry.getApiKeyHash() == null || existingEntry.getApiKeyHash().isBlank())
|
||||
&& (existingEntry.getCompanyId() == null || existingEntry.getCompanyId() <= 0)) {
|
||||
ziniaoMemoryStoreService.delete(CACHE_TYPE_SHOP_INDEX_ENTRY, existingEntry.getNormalizedShopName());
|
||||
ziniaoMemoryStoreService.delete(ZiniaoMemoryStoreService.CACHE_TYPE_SHOP_INDEX_ENTRY, rowKey);
|
||||
log.info("[ziniao-index] removed ghost conflict shop_index row cacheKey={}", rowKey);
|
||||
continue;
|
||||
}
|
||||
existingEntry.setStatus(STATUS_STALE);
|
||||
existingEntry.setMessage("店铺索引已过期,请等待后台刷新");
|
||||
existingEntry.setLastRefreshedAt(now);
|
||||
ziniaoMemoryStoreService.put(
|
||||
CACHE_TYPE_SHOP_INDEX_ENTRY,
|
||||
existingEntry.getNormalizedShopName(),
|
||||
existingEntry,
|
||||
resolveEntryTtl()
|
||||
);
|
||||
ziniaoMemoryStoreService.put(ZiniaoMemoryStoreService.CACHE_TYPE_SHOP_INDEX_ENTRY, rowKey, existingEntry, resolveEntryTtl());
|
||||
staleCount++;
|
||||
}
|
||||
if (staleCount > 0) {
|
||||
log.info("[ziniao-index] marked stale shop_index rows count={}", staleCount);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -370,9 +489,10 @@ public class ZiniaoShopIndexService {
|
||||
if (left == null || right == null) {
|
||||
return false;
|
||||
}
|
||||
if (left.getShopId() != null && right.getShopId() != null) {
|
||||
if (left.getShopId() != null && !left.getShopId().isBlank()
|
||||
&& right.getShopId() != null && !right.getShopId().isBlank()) {
|
||||
// 同一 apiKey + 公司 + shopId:多员工可见同店,紫鸟 shopId 相同,保留任一员工即可
|
||||
return Objects.equals(left.getShopId(), right.getShopId())
|
||||
&& Objects.equals(left.getMatchedUserId(), right.getMatchedUserId())
|
||||
&& Objects.equals(left.getCompanyId(), right.getCompanyId())
|
||||
&& Objects.equals(left.getApiKeyHash(), right.getApiKeyHash());
|
||||
}
|
||||
@@ -382,6 +502,23 @@ public class ZiniaoShopIndexService {
|
||||
&& Objects.equals(left.getApiKeyHash(), right.getApiKeyHash());
|
||||
}
|
||||
|
||||
private String buildStoresFingerprint(List<ZiniaoShopCacheDto> stores) {
|
||||
if (stores == null || stores.isEmpty()) {
|
||||
return "";
|
||||
}
|
||||
StringBuilder sb = new StringBuilder();
|
||||
stores.stream()
|
||||
.filter(Objects::nonNull)
|
||||
.sorted((left, right) -> Objects.toString(left.getShopId(), "").compareTo(Objects.toString(right.getShopId(), "")))
|
||||
.forEach(item -> sb.append(Objects.toString(item.getShopId(), ""))
|
||||
.append('|')
|
||||
.append(Objects.toString(item.getShopName(), ""))
|
||||
.append('|')
|
||||
.append(Objects.toString(item.getPlatform(), ""))
|
||||
.append(';'));
|
||||
return sb.toString();
|
||||
}
|
||||
|
||||
private String buildScopeSnapshotKey(String apiKey, Long companyId, Long userId) {
|
||||
return buildApiKeyHash(apiKey) + ":" + companyId + ":" + userId;
|
||||
}
|
||||
|
||||
@@ -14,6 +14,7 @@ from blueprints.auth import auth
|
||||
from blueprints.main import main
|
||||
from blueprints.admin_api import admin_api
|
||||
from blueprints.version import version_bp
|
||||
from blueprints.get_resource import get_resource
|
||||
|
||||
BASE_DIR = os.path.dirname(os.path.abspath(__file__))
|
||||
app = Flask(__name__, template_folder=BASE_DIR, static_folder=BASE_DIR)
|
||||
@@ -27,6 +28,7 @@ app.register_blueprint(auth)
|
||||
app.register_blueprint(main)
|
||||
app.register_blueprint(admin_api)
|
||||
app.register_blueprint(version_bp)
|
||||
app.register_blueprint(get_resource)
|
||||
|
||||
|
||||
def run_app(host='0.0.0.0', port=15124):
|
||||
|
||||
@@ -115,13 +115,13 @@
|
||||
<div v-if="shouldShowProgress(item)" class="delete-brand-progress-block">
|
||||
<div class="delete-brand-progress-header">
|
||||
<span>任务进度</span>
|
||||
<span>{{ formatProgressPercent(item.taskId!) }}%</span>
|
||||
<span>{{ formatProgressPercent(item) }}%</span>
|
||||
</div>
|
||||
<div class="delete-brand-progress-bar">
|
||||
<div class="delete-brand-progress-bar-fill"
|
||||
:style="{ width: `${formatProgressPercent(item.taskId!)}%` }"></div>
|
||||
:style="{ width: `${formatProgressPercent(item)}%` }"></div>
|
||||
</div>
|
||||
<div class="files delete-brand-progress">{{ formatProgress(item.taskId!) }}</div>
|
||||
<div class="files delete-brand-progress">{{ formatProgress(item) }}</div>
|
||||
</div>
|
||||
<div v-if="item.previewRows?.length" class="files split-entry-list delete-brand-preview">
|
||||
{{ formatPreview(item.previewRows) }}
|
||||
@@ -167,13 +167,13 @@
|
||||
<div v-if="shouldShowProgress(item)" class="delete-brand-progress-block">
|
||||
<div class="delete-brand-progress-header">
|
||||
<span>任务进度</span>
|
||||
<span>{{ formatProgressPercent(item.taskId!) }}%</span>
|
||||
<span>{{ formatProgressPercent(item) }}%</span>
|
||||
</div>
|
||||
<div class="delete-brand-progress-bar">
|
||||
<div class="delete-brand-progress-bar-fill"
|
||||
:style="{ width: `${formatProgressPercent(item.taskId!)}%` }"></div>
|
||||
:style="{ width: `${formatProgressPercent(item)}%` }"></div>
|
||||
</div>
|
||||
<div class="files delete-brand-progress">{{ formatProgress(item.taskId!) }}</div>
|
||||
<div class="files delete-brand-progress">{{ formatProgress(item) }}</div>
|
||||
</div>
|
||||
<div v-if="item.previewRows?.length" class="files split-entry-list delete-brand-preview">
|
||||
{{ formatPreview(item.previewRows) }}
|
||||
@@ -224,9 +224,14 @@ import {
|
||||
} from '@/shared/api/java-modules'
|
||||
import { getPywebviewApi, type UploadedJavaFile } from '@/shared/bridges/pywebview'
|
||||
|
||||
interface SessionDeleteBrandItem extends DeleteBrandResultItem {
|
||||
_pushed?: boolean
|
||||
_completed?: boolean
|
||||
}
|
||||
|
||||
interface StoredCurrentTask {
|
||||
taskId: number
|
||||
items: DeleteBrandResultItem[]
|
||||
items: SessionDeleteBrandItem[]
|
||||
createdAt: number
|
||||
queuePushResult?: string
|
||||
queuePayloadText?: string
|
||||
@@ -318,6 +323,84 @@ const historySectionItems = computed(() =>
|
||||
);
|
||||
const hasVisibleItems = computed(() => currentSectionItems.value.length > 0 || historySectionItems.value.length > 0)
|
||||
|
||||
function persistSessionTasks() {
|
||||
if (typeof window !== 'undefined') {
|
||||
window.localStorage.setItem(getStorageKey(), JSON.stringify(sessionTasks.value))
|
||||
}
|
||||
}
|
||||
|
||||
function isMultiFileTask(taskId?: number) {
|
||||
if (!taskId) return false
|
||||
const sessionTask = sessionTasks.value.find(t => t.taskId === taskId)
|
||||
return (sessionTask?.items.length || 0) > 1
|
||||
}
|
||||
|
||||
function findSessionTask(taskId?: number) {
|
||||
if (!taskId) return null
|
||||
return sessionTasks.value.find(t => t.taskId === taskId) || null
|
||||
}
|
||||
|
||||
function isSameDeleteBrandItem(left: Pick<DeleteBrandResultItem, 'resultId' | 'fileKey' | 'sourceFilename'> | null | undefined,
|
||||
right: Pick<DeleteBrandResultItem, 'resultId' | 'fileKey' | 'sourceFilename'> | null | undefined) {
|
||||
if (!left || !right) return false
|
||||
if (left.resultId != null && right.resultId != null) {
|
||||
return left.resultId === right.resultId
|
||||
}
|
||||
if (left.fileKey && right.fileKey) {
|
||||
return left.fileKey === right.fileKey
|
||||
}
|
||||
return left.sourceFilename === right.sourceFilename
|
||||
}
|
||||
|
||||
function findSessionItem(taskId: number | undefined, item: DeleteBrandResultItem) {
|
||||
const sessionTask = findSessionTask(taskId)
|
||||
if (!sessionTask) return null
|
||||
return sessionTask.items.find(i => isSameDeleteBrandItem(i, item)) || null
|
||||
}
|
||||
|
||||
function updateCompletedItemsFromProgress(taskId: number | undefined) {
|
||||
if (!taskId) return false
|
||||
|
||||
const sessionTask = findSessionTask(taskId)
|
||||
if (!sessionTask) return false
|
||||
|
||||
let changed = false
|
||||
const info = taskDetails.value[taskId]?.line_progress?.info
|
||||
|
||||
// 多文件任务优先依据后端 finished_files 同步已完成文件数量。
|
||||
// 规则:按 sessionTask.items 的顺序,标记「已推送」中的前 N 个为 completed。
|
||||
if (isMultiFileTask(taskId) && info?.finished_files != null && info.finished_files >= 0) {
|
||||
const pushedItems = sessionTask.items.filter(i => i._pushed)
|
||||
const shouldCompleted = Math.max(0, Math.min(info.finished_files, pushedItems.length))
|
||||
|
||||
for (let i = 0; i < shouldCompleted; i += 1) {
|
||||
if (!pushedItems[i]._completed) {
|
||||
pushedItems[i]._completed = true
|
||||
changed = true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 兼容单文件/兜底:当当前文件行进度到达末行时,标记活动文件完成。
|
||||
if (info?.file_name && info.current_line != null && info.total_lines != null && info.total_lines > 0) {
|
||||
if (info.current_line >= info.total_lines) {
|
||||
const activeItem = currentSectionItems.value.find(item => getItemKey(item) === activeItemKey.value)
|
||||
const sessionItem = activeItem && activeItem.taskId === taskId
|
||||
? findSessionItem(taskId, activeItem)
|
||||
: sessionTask.items.find(i => i.sourceFilename === info.file_name)
|
||||
if (sessionItem && !sessionItem._completed) {
|
||||
sessionItem._completed = true
|
||||
changed = true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (changed) {
|
||||
persistSessionTasks()
|
||||
}
|
||||
return changed
|
||||
}
|
||||
|
||||
function getTaskStatus(taskId?: number) {
|
||||
if (!taskId) return ''
|
||||
return taskDetails.value[taskId]?.task?.status || ''
|
||||
@@ -334,7 +417,40 @@ function getDisplayError(item: DeleteBrandResultItem) {
|
||||
}
|
||||
|
||||
function shouldShowProgress(item: DeleteBrandResultItem) {
|
||||
const taskId = item.taskId
|
||||
if (!taskId) return false
|
||||
|
||||
const progress = taskDetails.value[taskId]?.line_progress
|
||||
const status = getTaskStatus(taskId)
|
||||
const multiFileTask = isMultiFileTask(taskId)
|
||||
const qs = getQueueStatus(item)
|
||||
|
||||
if (multiFileTask) {
|
||||
const sessionItem = findSessionItem(taskId, item)
|
||||
const isActiveItem = getItemKey(item) === activeItemKey.value
|
||||
const isProgressTarget = Boolean(progress?.has_progress && progress?.info?.file_name === item.sourceFilename)
|
||||
|
||||
// 多文件任务下,展示规则更宽松:
|
||||
// 1) 当前活跃文件;2) 后端进度指向的文件;3) 已推送文件;4) 已完成文件(任务未终态时保留)。
|
||||
if (isActiveItem || isProgressTarget || sessionItem?._pushed) {
|
||||
return true
|
||||
}
|
||||
if ((status === 'RUNNING' || status === 'PENDING' || !status) && sessionItem?._completed) {
|
||||
return true
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
// 单文件任务:完成/终态后隐藏。
|
||||
if (qs === '本文件解析完毕' || qs === '已被全局判定为终态') {
|
||||
return false
|
||||
}
|
||||
|
||||
if (progress?.has_progress || progress?.info) {
|
||||
return true
|
||||
}
|
||||
|
||||
return qs === '已入队' || qs === '处理中'
|
||||
}
|
||||
|
||||
@@ -342,14 +458,20 @@ function shouldShowProgress(item: DeleteBrandResultItem) {
|
||||
function normalizeDeleteBrandItems(items: DeleteBrandResultItem[]) {
|
||||
return items.map((item) => {
|
||||
if (item.matchStatus === 'MATCHED') {
|
||||
return item
|
||||
return {
|
||||
...item,
|
||||
_pushed: false,
|
||||
_completed: false,
|
||||
} as SessionDeleteBrandItem
|
||||
}
|
||||
return {
|
||||
...item,
|
||||
matched: false,
|
||||
openStoreUrl: undefined,
|
||||
error: item.error || undefined,
|
||||
}
|
||||
_pushed: false,
|
||||
_completed: false,
|
||||
} as SessionDeleteBrandItem
|
||||
})
|
||||
}
|
||||
|
||||
@@ -390,10 +512,40 @@ function updateSummary(items: DeleteBrandResultItem[]) {
|
||||
}
|
||||
}
|
||||
|
||||
function formatProgress(taskId: number) {
|
||||
function getFileProgress(item: DeleteBrandResultItem) {
|
||||
const taskId = item.taskId
|
||||
if (!taskId) return null
|
||||
const detail = taskDetails.value[taskId]
|
||||
const fileProgress = detail?.fileProgress || []
|
||||
return fileProgress.find((fp) => {
|
||||
if (fp.fileKey && item.fileKey) return fp.fileKey === item.fileKey
|
||||
if (fp.sourceFilename && item.sourceFilename) return fp.sourceFilename === item.sourceFilename
|
||||
return false
|
||||
}) || null
|
||||
}
|
||||
|
||||
function formatProgress(item: DeleteBrandResultItem) {
|
||||
const taskId = item.taskId
|
||||
if (!taskId) return ''
|
||||
|
||||
const detail = taskDetails.value[taskId]
|
||||
const info = detail?.line_progress?.info
|
||||
const fp = getFileProgress(item)
|
||||
|
||||
if (fp) {
|
||||
const parts: string[] = []
|
||||
const processedRows = fp.processedRows ?? 0
|
||||
const totalRows = fp.totalRows ?? 0
|
||||
if (fp.status) parts.push(`文件状态:${fp.status}`)
|
||||
parts.push(`文件进度:${processedRows}/${totalRows}`)
|
||||
if (info?.finished_files != null) {
|
||||
const total = info.file_total && info.file_total > 0 ? `/${info.file_total}` : ''
|
||||
parts.push(`已完成文件:${info.finished_files}${total}`)
|
||||
}
|
||||
if (info?.phase) parts.push(`阶段:${info.phase}`)
|
||||
return parts.join('|')
|
||||
}
|
||||
|
||||
if (!detail?.line_progress?.has_progress || !info) {
|
||||
const status = detail?.task?.status
|
||||
return status ? `任务状态:${status}` : ''
|
||||
@@ -404,16 +556,26 @@ function formatProgress(taskId: number) {
|
||||
if (info.file_index && info.file_total) parts.push(`文件:${info.file_index}/${info.file_total}`)
|
||||
if (info.file_name) parts.push(`当前文件:${info.file_name}`)
|
||||
if (info.current_line != null && info.total_lines != null) {
|
||||
parts.push(`进度:${info.current_line}/${info.total_lines}`)
|
||||
parts.push(`${isMultiFileTask(taskId) ? '当前文件进度' : '进度'}:${info.current_line}/${info.total_lines}`)
|
||||
}
|
||||
if (info.current_country) parts.push(`国家:${info.current_country}`)
|
||||
if (info.current_asin) parts.push(`ASIN:${info.current_asin}`)
|
||||
if (info.finished_files != null) parts.push(`已完成文件:${info.finished_files}`)
|
||||
if (info.finished_files != null) {
|
||||
const total = info.file_total && info.file_total > 0 ? `/${info.file_total}` : ''
|
||||
parts.push(`已完成文件:${info.finished_files}${total}`)
|
||||
}
|
||||
return parts.join('|')
|
||||
}
|
||||
|
||||
function formatProgressPercent(taskId: number) {
|
||||
function formatProgressPercent(item: DeleteBrandResultItem) {
|
||||
const taskId = item.taskId
|
||||
if (!taskId) return 0
|
||||
|
||||
const fp = getFileProgress(item)
|
||||
if (fp) {
|
||||
return Math.max(0, Math.min(100, fp.percent ?? 0))
|
||||
}
|
||||
|
||||
const info = taskDetails.value[taskId]?.line_progress?.info
|
||||
if (!info) return 0
|
||||
|
||||
@@ -446,34 +608,48 @@ function getPollingTaskIds() {
|
||||
function getQueueStatus(item: DeleteBrandResultItem) {
|
||||
const taskId = item.taskId
|
||||
if (!taskId) return ''
|
||||
const sessionTask = sessionTasks.value.find(t => t.taskId === taskId)
|
||||
const sessionTask = findSessionTask(taskId)
|
||||
if (!sessionTask) return ''
|
||||
|
||||
const sessionItem = findSessionItem(taskId, item)
|
||||
const status = getTaskStatus(taskId)
|
||||
if (status === 'SUCCESS' || status === 'FAILED' || item.taskStatus === 'SUCCESS' || item.taskStatus === 'FAILED') {
|
||||
const info = taskDetails.value[taskId]?.line_progress
|
||||
const multiFileTask = isMultiFileTask(taskId)
|
||||
|
||||
if (!multiFileTask && (status === 'SUCCESS' || status === 'FAILED' || item.taskStatus === 'SUCCESS' || item.taskStatus === 'FAILED')) {
|
||||
return '已被全局判定为终态'
|
||||
}
|
||||
|
||||
const sessionItem = sessionTask.items.find(i => i.resultId === item.resultId || i.sourceFilename === item.sourceFilename)
|
||||
// 先看前端会话态:一旦该文件已判定完成,优先展示完成,避免被 RUNNING 覆盖导致链式推进卡住。
|
||||
if (sessionItem?._completed) {
|
||||
return '本文件解析完毕'
|
||||
}
|
||||
|
||||
const info = taskDetails.value[taskId]?.line_progress
|
||||
|
||||
// 此时这个特定文件是否正好被后端报了进度名
|
||||
// 多文件场景下,如果后端进度明确指向当前文件,即使本地 _pushed 丢失也应判定为处理中(例如刷新页面后恢复场景)。
|
||||
if (status === 'RUNNING' && info?.has_progress && info?.info?.file_name === item.sourceFilename) {
|
||||
return '处理中'
|
||||
}
|
||||
|
||||
// 查查是否已经完结了这个特定文件
|
||||
const detail = taskDetails.value[taskId]
|
||||
if (detail?.items) {
|
||||
// 后端返回的 items 是 `resultJson` 解析来的,里面只含有已经做完的!
|
||||
const finished = detail.items.find((i: DeleteBrandResultItem) => i.sourceFilename === item.sourceFilename)
|
||||
if (finished) {
|
||||
return '本文件解析完毕'
|
||||
// 多文件:未推送且未完成,才视为未入队。
|
||||
if (multiFileTask && !sessionItem?._pushed && !sessionItem?._completed) {
|
||||
return '未入队'
|
||||
}
|
||||
|
||||
if (multiFileTask && status === 'SUCCESS') {
|
||||
return '本文件解析完毕'
|
||||
}
|
||||
|
||||
if (!multiFileTask) {
|
||||
const detail = taskDetails.value[taskId]
|
||||
if (detail?.items) {
|
||||
const finished = detail.items.find((i: DeleteBrandResultItem) => isSameDeleteBrandItem(i, item))
|
||||
if (finished && status === 'SUCCESS') {
|
||||
return '本文件解析完毕'
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (sessionItem && (sessionItem as any)._pushed) {
|
||||
if (sessionItem?._pushed) {
|
||||
return '已入队'
|
||||
}
|
||||
|
||||
@@ -497,12 +673,14 @@ function mergeTaskDetailItemsIntoSession(detail: DeleteBrandTaskDetailVo) {
|
||||
|
||||
let changed = false
|
||||
const mergedItems = sessionTask.items.map((item) => {
|
||||
const latest = detail.items?.find((candidate) => candidate.resultId === item.resultId || candidate.sourceFilename === item.sourceFilename)
|
||||
const latest = detail.items?.find((candidate) => isSameDeleteBrandItem(candidate, item))
|
||||
if (!latest) return item
|
||||
const merged = {
|
||||
...item,
|
||||
...latest,
|
||||
success: item.success ?? latest.success,
|
||||
_pushed: item._pushed,
|
||||
_completed: item._completed,
|
||||
}
|
||||
if (JSON.stringify(merged) !== JSON.stringify(item)) {
|
||||
changed = true
|
||||
@@ -516,8 +694,8 @@ function mergeTaskDetailItemsIntoSession(detail: DeleteBrandTaskDetailVo) {
|
||||
changed = true
|
||||
}
|
||||
|
||||
if (changed && typeof window !== 'undefined') {
|
||||
window.localStorage.setItem(getStorageKey(), JSON.stringify(sessionTasks.value))
|
||||
if (changed) {
|
||||
persistSessionTasks()
|
||||
}
|
||||
return changed
|
||||
}
|
||||
@@ -536,6 +714,9 @@ async function refreshTaskDetails(taskIds?: number[]) {
|
||||
if (mergeTaskDetailItemsIntoSession(detail)) {
|
||||
changed = true
|
||||
}
|
||||
if (updateCompletedItemsFromProgress(id)) {
|
||||
changed = true
|
||||
}
|
||||
}
|
||||
}
|
||||
if (changed) {
|
||||
@@ -812,6 +993,11 @@ function formatPreview(rows: DeleteBrandPreviewRow[]) {
|
||||
return hiddenCount > 0 ? `${preview.join('、')} 等 ${rows.length} 条` : preview.join('、')
|
||||
}
|
||||
|
||||
function debugAutoAdvance(step: string, payload?: Record<string, unknown>) {
|
||||
if (!import.meta.env.DEV) return
|
||||
console.log(`[DeleteBrandAutoAdvance] ${step}`, payload || {})
|
||||
}
|
||||
|
||||
function isPythonQueueBusy() {
|
||||
for (const task of sessionTasks.value) {
|
||||
for (const item of task.items) {
|
||||
@@ -833,12 +1019,19 @@ function findNextAutoRunnableItem() {
|
||||
return currentSectionItems.value.find((item) => {
|
||||
if (!item.openStoreUrl) return false
|
||||
if (getItemKey(item) === activeItemKey.value) return false
|
||||
const sessionItem = findSessionItem(item.taskId, item)
|
||||
if (sessionItem?._completed) return false
|
||||
return getQueueStatus(item) === '未入队'
|
||||
})
|
||||
}
|
||||
|
||||
async function runItem(item: DeleteBrandResultItem, options?: { auto?: boolean }) {
|
||||
if (options?.auto && isPythonQueueBusy()) {
|
||||
debugAutoAdvance('runItem blocked by busy queue', {
|
||||
mode: 'auto',
|
||||
taskId: item.taskId,
|
||||
sourceFilename: item.sourceFilename,
|
||||
})
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -870,9 +1063,10 @@ async function runItem(item: DeleteBrandResultItem, options?: { auto?: boolean }
|
||||
if (pushResult?.success) {
|
||||
qpr = `已推送文件 ${item.sourceFilename},当前队列长度:${pushResult.queue_size ?? '-'}`
|
||||
|
||||
const sessionItem = sessionTask.items.find(i => i.resultId === item.resultId || i.sourceFilename === item.sourceFilename)
|
||||
const sessionItem = findSessionItem(taskId, item)
|
||||
if (sessionItem) {
|
||||
(sessionItem as any)._pushed = true
|
||||
sessionItem._pushed = true
|
||||
sessionItem._completed = false
|
||||
}
|
||||
activeItemKey.value = getItemKey(item)
|
||||
saveSessionTask(taskId, sessionTask.items, qpr, queuePayloadText.value)
|
||||
@@ -885,22 +1079,63 @@ async function runItem(item: DeleteBrandResultItem, options?: { auto?: boolean }
|
||||
}
|
||||
|
||||
async function maybeAutoAdvance() {
|
||||
if (!chainStarted.value || autoAdvancing.value) return
|
||||
if (isPythonQueueBusy()) return
|
||||
if (!chainStarted.value || autoAdvancing.value) {
|
||||
debugAutoAdvance('skip: chain not started or auto advancing', {
|
||||
chainStarted: chainStarted.value,
|
||||
autoAdvancing: autoAdvancing.value,
|
||||
})
|
||||
return
|
||||
}
|
||||
if (isPythonQueueBusy()) {
|
||||
debugAutoAdvance('skip: queue busy')
|
||||
return
|
||||
}
|
||||
|
||||
const currentKey = activeItemKey.value
|
||||
if (currentKey) {
|
||||
const currentItem = currentSectionItems.value.find(item => getItemKey(item) === currentKey)
|
||||
if (currentItem) {
|
||||
const taskId = currentItem.taskId
|
||||
const multiFileTask = isMultiFileTask(taskId)
|
||||
if (multiFileTask) {
|
||||
const changed = updateCompletedItemsFromProgress(taskId)
|
||||
debugAutoAdvance('sync completed from progress', {
|
||||
taskId,
|
||||
changed,
|
||||
status: getQueueStatus(currentItem),
|
||||
})
|
||||
}
|
||||
const status = getQueueStatus(currentItem)
|
||||
if (status === '已入队' || status === '处理中') {
|
||||
debugAutoAdvance('skip: current item still running', {
|
||||
currentKey,
|
||||
status,
|
||||
taskId,
|
||||
sourceFilename: currentItem.sourceFilename,
|
||||
})
|
||||
return
|
||||
}
|
||||
if (multiFileTask) {
|
||||
const sessionItem = findSessionItem(taskId, currentItem)
|
||||
if (sessionItem && !sessionItem._completed) {
|
||||
debugAutoAdvance('skip: multi-file current not completed yet', {
|
||||
currentKey,
|
||||
taskId,
|
||||
sourceFilename: currentItem.sourceFilename,
|
||||
pushed: sessionItem._pushed,
|
||||
completed: sessionItem._completed,
|
||||
})
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const nextItem = findNextAutoRunnableItem()
|
||||
if (!nextItem) {
|
||||
debugAutoAdvance('stop chain: no next runnable item', {
|
||||
currentKey,
|
||||
})
|
||||
activeItemKey.value = ''
|
||||
chainStarted.value = false
|
||||
return
|
||||
@@ -908,8 +1143,17 @@ async function maybeAutoAdvance() {
|
||||
|
||||
autoAdvancing.value = true
|
||||
try {
|
||||
debugAutoAdvance('try start next item', {
|
||||
taskId: nextItem.taskId,
|
||||
sourceFilename: nextItem.sourceFilename,
|
||||
nextKey: getItemKey(nextItem),
|
||||
})
|
||||
const started = await runItem(nextItem, { auto: true })
|
||||
if (!started) {
|
||||
debugAutoAdvance('stop chain: runItem returned false', {
|
||||
taskId: nextItem.taskId,
|
||||
sourceFilename: nextItem.sourceFilename,
|
||||
})
|
||||
chainStarted.value = false
|
||||
}
|
||||
} finally {
|
||||
|
||||
@@ -200,10 +200,20 @@ export interface DeleteBrandTaskItem {
|
||||
downloadFilename?: string
|
||||
}
|
||||
|
||||
export interface DeleteBrandTaskFileProgress {
|
||||
fileKey?: string
|
||||
sourceFilename?: string
|
||||
processedRows?: number
|
||||
totalRows?: number
|
||||
percent?: number
|
||||
status?: 'PENDING' | 'RUNNING' | 'COMPLETED' | 'FAILED'
|
||||
}
|
||||
|
||||
export interface DeleteBrandTaskDetailVo {
|
||||
task: DeleteBrandTaskItem
|
||||
line_progress: DeleteBrandLineProgress
|
||||
items?: DeleteBrandResultItem[]
|
||||
fileProgress?: DeleteBrandTaskFileProgress[]
|
||||
}
|
||||
|
||||
export interface DeleteBrandTaskBatchVo {
|
||||
|
||||
Reference in New Issue
Block a user