From 8cd1a4e2b15f8b0ff864664b3e4bafb4ddebc96c Mon Sep 17 00:00:00 2001 From: super <2903208875@qq.com> Date: Wed, 6 May 2026 22:41:26 +0800 Subject: [PATCH] =?UTF-8?q?=E6=9B=B4=E6=96=B0=E5=A4=96=E8=A7=82=E6=A8=A1?= =?UTF-8?q?=E5=9D=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/amazon/similar_asin.py | 69 +++- .../config/TaskOperationLockConfig.java | 4 + .../service/DeleteBrandStaleTaskService.java | 38 ++ .../controller/PatrolDeleteController.java | 6 +- .../service/PatrolDeleteTaskService.java | 73 +++- .../service/PriceTrackTaskService.java | 77 +++- .../service/ProductRiskTaskService.java | 134 ++++--- .../service/QueryAsinTaskService.java | 55 ++- .../service/ShopMatchTaskService.java | 60 +++- .../controller/SimilarAsinController.java | 40 +++ .../SimilarAsinFilterConditionMapper.java | 9 + .../SimilarAsinFilterConditionAddRequest.java | 24 ++ .../SimilarAsinFilterConditionEntity.java | 19 + .../vo/SimilarAsinFilterConditionVo.java | 17 + .../service/SimilarAsinTaskService.java | 111 +++++- ...trol_delete_delete_and_history_indexes.sql | 35 ++ .../db/V47__similar_asin_filter_condition.sql | 15 + frontend-vue/similar-asin.html | 2 +- .../BrandApiSecretSettingsButton.vue | 337 ++++++++++++++++++ .../components/BrandAppearancePatentTab.vue | 16 +- .../brand/components/BrandSimilarAsinTab.vue | 203 ++++++----- .../pages/brand/components/BrandTopBar.vue | 12 +- frontend-vue/src/shared/api/java-modules.ts | 41 ++- .../src/shared/utils/api-secret-store.ts | 156 ++++++++ 24 files changed, 1378 insertions(+), 175 deletions(-) create mode 100644 backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/mapper/SimilarAsinFilterConditionMapper.java create mode 100644 backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinFilterConditionAddRequest.java create mode 100644 backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/entity/SimilarAsinFilterConditionEntity.java create mode 100644 backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/vo/SimilarAsinFilterConditionVo.java create mode 100644 backend-java/src/main/resources/db/V46__patrol_delete_delete_and_history_indexes.sql create mode 100644 backend-java/src/main/resources/db/V47__similar_asin_filter_condition.sql create mode 100644 frontend-vue/src/pages/brand/components/BrandApiSecretSettingsButton.vue create mode 100644 frontend-vue/src/shared/utils/api-secret-store.ts diff --git a/app/amazon/similar_asin.py b/app/amazon/similar_asin.py index 352d6c0..4c7afd8 100644 --- a/app/amazon/similar_asin.py +++ b/app/amazon/similar_asin.py @@ -276,8 +276,8 @@ class SimilarAsinTask(TaskBase): return list(grouped.values()) - @staticmethod - def normalize_groups(data): + @staticmethod + def normalize_groups(data): groups = data.get("groups") if isinstance(groups, list) and len(groups) > 0: return groups @@ -299,18 +299,71 @@ class SimilarAsinTask(TaskBase): "displayId": first.get("displayId") or item_id, "items": items, }) - return normalized + return normalized + + def fetch_parsed_payload(self, task_id, user_id=1): + url = f"{DELETE_BRAND_API_BASE}/api/similar-asin/tasks/{task_id}/parsed-payload" + response = requests.get(url, params={"user_id": user_id}, timeout=300, verify=False) + response.raise_for_status() + payload = response.json() if response.text else {} + if isinstance(payload, dict) and payload.get("success"): + return payload.get("data") or {} + raise RuntimeError(f"获取货源查询解析载荷失败: {payload}") - def process_task(self, task_data: dict): + @staticmethod + def _count_payload_rows(data): + rows = data.get("rows") or data.get("items") or [] + if isinstance(rows, list) and rows: + return len(rows) + + groups = data.get("groups") or [] + if not isinstance(groups, list): + return 0 + + total = 0 + for group in groups: + if isinstance(group, dict) and isinstance(group.get("items"), list): + total += len(group.get("items")) + return total + + def _merge_parsed_payload(self, data, parsed_payload): + payload_rows = parsed_payload.get("allItems") or parsed_payload.get("items") or [] + merged = { + **data, + "groups": parsed_payload.get("groups") or [], + "rows": payload_rows, + } + if parsed_payload.get("aiPrompt"): + merged["prompt"] = parsed_payload.get("aiPrompt") + if parsed_payload.get("apiKey") and not merged.get("api_key"): + merged["api_key"] = parsed_payload.get("apiKey") + return merged + + def process_task(self, task_data: dict): """处理审批任务主入口 Args: task_data: 任务数据 """ try: - data = task_data.get("data", {}) - task_id = data.get("taskId") - groups = self.normalize_groups(data) + data = task_data.get("data", {}) + task_id = data.get("taskId") + if task_id: + queued_rows = self._count_payload_rows(data) + try: + parsed_payload = self.fetch_parsed_payload(task_id, data.get("user_id") or data.get("userId") or 1) + data = self._merge_parsed_payload(data, parsed_payload) + self.log( + f"similar asin task {task_id} loaded full parsed payload from Java, queuedRows={queued_rows}, fullRows={self._count_payload_rows(data)}" + ) + except Exception as e: + if not queued_rows: + raise + self.log( + f"similar asin task {task_id} failed to load full parsed payload, fallback to queued rows={queued_rows}: {e}", + "WARNING" + ) + groups = self.normalize_groups(data) if not task_id: self.log("任务ID为空,跳过", "WARNING") @@ -447,7 +500,7 @@ class SimilarAsinTask(TaskBase): """回传处理结果到API """ - url = f"{DELETE_BRAND_API_BASE}/api/appearance-patent/tasks/{task_id}/result" + url = f"{DELETE_BRAND_API_BASE}/api/similar-asin/tasks/{task_id}/result" payload ={ "submissionId": f"{int(time.time())}", diff --git a/backend-java/src/main/java/com/nanri/aiimage/config/TaskOperationLockConfig.java b/backend-java/src/main/java/com/nanri/aiimage/config/TaskOperationLockConfig.java index b52ab42..c6c1601 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/config/TaskOperationLockConfig.java +++ b/backend-java/src/main/java/com/nanri/aiimage/config/TaskOperationLockConfig.java @@ -35,6 +35,7 @@ public class TaskOperationLockConfig implements WebMvcConfigurer { private static final Pattern TASK_RESULT_PATH = Pattern.compile(".*/api/([^/]+)/tasks/(\\d+)/result/?$"); private static final Set MUTATING_METHODS = Set.of("POST", "PUT", "PATCH", "DELETE"); private static final long RESULT_SUBMIT_WAIT_MILLIS = 10 * 60 * 1000L; + private static final long DELETE_WAIT_MILLIS = 60 * 1000L; private static final String LOCK_ATTRIBUTE = TaskOperationLockInterceptor.class.getName() + ".LOCK"; private final TaskDistributedLockService taskDistributedLockService; @@ -71,6 +72,9 @@ public class TaskOperationLockConfig implements WebMvcConfigurer { if ("POST".equals(method) && TASK_RESULT_PATH.matcher(uri == null ? "" : uri).matches()) { return RESULT_SUBMIT_WAIT_MILLIS; } + if ("DELETE".equals(method)) { + return DELETE_WAIT_MILLIS; + } return TaskDistributedLockService.DEFAULT_WAIT_MILLIS; } diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandStaleTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandStaleTaskService.java index d1f8e5a..440b009 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandStaleTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandStaleTaskService.java @@ -15,6 +15,7 @@ import com.nanri.aiimage.modules.shopmatch.service.ShopMatchTaskService; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity; import com.nanri.aiimage.modules.task.service.TaskDistributedLockService; +import com.nanri.aiimage.modules.task.service.TaskFileJobService; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; @@ -60,6 +61,7 @@ public class DeleteBrandStaleTaskService { private final DeleteBrandProgressProperties deleteBrandProgressProperties; private final DistributedJobLockService distributedJobLockService; private final TaskDistributedLockService taskDistributedLockService; + private final TaskFileJobService taskFileJobService; @Value("${aiimage.temp-dir.retention-hours:24}") private long tempDirRetentionHours; @@ -231,6 +233,14 @@ public class DeleteBrandStaleTaskService { if (!hasStartedProgress) { hasStartedProgress = productRiskTaskCacheService.hasAnyShopMergedPayload(task.getId()); } + if (hasPendingAssembleJobs(task.getId(), MODULE_TYPE_PRODUCT_RISK)) { + stats.skippedTaskCount++; + log.info("[stale-check] product-risk skip pending-assemble-jobs taskId={} unfinishedJobs={} activeJobs={}", + task.getId(), + taskFileJobService.countUnfinishedAssembleJobs(task.getId(), MODULE_TYPE_PRODUCT_RISK), + taskFileJobService.countActiveAssembleJobs(task.getId(), MODULE_TYPE_PRODUCT_RISK)); + continue; + } if (!hasStartedProgress && isWithinInitialGrace(task, initialThreshold)) { stats.skippedTaskCount++; log.info("[stale-check] product-risk skip initial-grace taskId={} createdAt={} initialThreshold={}", @@ -305,6 +315,14 @@ public class DeleteBrandStaleTaskService { if (!hasStartedProgress) { hasStartedProgress = priceTrackTaskCacheService.hasAnyShopMergedPayload(task.getId()); } + if (hasPendingAssembleJobs(task.getId(), MODULE_TYPE_PRICE_TRACK)) { + stats.skippedTaskCount++; + log.info("[stale-check] price-track skip pending-assemble-jobs taskId={} unfinishedJobs={} activeJobs={}", + task.getId(), + taskFileJobService.countUnfinishedAssembleJobs(task.getId(), MODULE_TYPE_PRICE_TRACK), + taskFileJobService.countActiveAssembleJobs(task.getId(), MODULE_TYPE_PRICE_TRACK)); + continue; + } if (!hasStartedProgress && isWithinInitialGrace(task, initialThreshold)) { stats.skippedTaskCount++; continue; @@ -380,6 +398,14 @@ public class DeleteBrandStaleTaskService { if (!hasStartedProgress) { hasStartedProgress = shopMatchTaskCacheService.hasAnyShopMergedPayload(task.getId()); } + if (hasPendingAssembleJobs(task.getId(), MODULE_TYPE_SHOP_MATCH)) { + stats.skippedTaskCount++; + log.info("[stale-check] shop-match skip pending-assemble-jobs taskId={} unfinishedJobs={} activeJobs={}", + task.getId(), + taskFileJobService.countUnfinishedAssembleJobs(task.getId(), MODULE_TYPE_SHOP_MATCH), + taskFileJobService.countActiveAssembleJobs(task.getId(), MODULE_TYPE_SHOP_MATCH)); + continue; + } if (!hasStartedProgress && isWithinInitialGrace(task, initialThreshold)) { stats.skippedTaskCount++; log.info("[stale-check] shop-match skip initial-grace taskId={} createdAt={} initialThreshold={}", @@ -453,6 +479,14 @@ public class DeleteBrandStaleTaskService { if (!hasStartedProgress) { hasStartedProgress = patrolDeleteTaskCacheService.hasAnyShopMergedPayload(task.getId()); } + if (hasPendingAssembleJobs(task.getId(), MODULE_TYPE_PATROL_DELETE)) { + stats.skippedTaskCount++; + log.info("[stale-check] patrol-delete skip pending-assemble-jobs taskId={} unfinishedJobs={} activeJobs={}", + task.getId(), + taskFileJobService.countUnfinishedAssembleJobs(task.getId(), MODULE_TYPE_PATROL_DELETE), + taskFileJobService.countActiveAssembleJobs(task.getId(), MODULE_TYPE_PATROL_DELETE)); + continue; + } if (!hasStartedProgress && isWithinInitialGrace(task, initialThreshold)) { stats.skippedTaskCount++; continue; @@ -612,6 +646,10 @@ public class DeleteBrandStaleTaskService { } } + private boolean hasPendingAssembleJobs(Long taskId, String moduleType) { + return taskFileJobService.countUnfinishedAssembleJobs(taskId, moduleType) > 0; + } + private static final class ProductRiskStaleCheckStats { private int scannedTaskCount; private int finalizedTaskCount; diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/controller/PatrolDeleteController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/controller/PatrolDeleteController.java index 98855b9..7776071 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/controller/PatrolDeleteController.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/controller/PatrolDeleteController.java @@ -155,8 +155,10 @@ public class PatrolDeleteController { @Operation(summary = "查询任务记录", description = "返回当前用户在巡店删除模块中的当前任务和历史任务记录。") public ApiResponse history( @Parameter(name = "user_id", description = "当前用户 ID", required = true, in = ParameterIn.QUERY, example = "1") - @RequestParam("user_id") Long userId) { - return ApiResponse.success(patrolDeleteTaskService.listHistory(userId)); + @RequestParam("user_id") Long userId, + @Parameter(description = "history limit, default 30, max 100", example = "30") + @RequestParam(value = "limit", required = false) Integer limit) { + return ApiResponse.success(patrolDeleteTaskService.listHistory(userId, limit)); } @PostMapping("/tasks/progress/batch") diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/service/PatrolDeleteTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/service/PatrolDeleteTaskService.java index 91e1b9b..6870888 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/service/PatrolDeleteTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/service/PatrolDeleteTaskService.java @@ -25,6 +25,7 @@ import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.model.entity.FileResultEntity; import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity; import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity; +import com.nanri.aiimage.modules.task.service.TaskDistributedLockService; import com.nanri.aiimage.modules.task.service.TaskFileJobService; import com.nanri.aiimage.modules.task.service.TaskProgressSnapshotService; import com.nanri.aiimage.modules.task.service.TaskResultItemService; @@ -68,6 +69,7 @@ public class PatrolDeleteTaskService { private final TaskResultItemService taskResultItemService; private final TaskProgressSnapshotService taskProgressSnapshotService; private final TaskScopePayloadStorageService taskScopePayloadStorageService; + private final TaskDistributedLockService taskDistributedLockService; private FileTaskEntity loadTaskForExecution(Long taskId) { Map cachedTasks = taskCacheService.getTaskCacheBatch(List.of(taskId)); @@ -137,8 +139,13 @@ public class PatrolDeleteTaskService { } public PatrolDeleteHistoryVo listHistory(Long userId) { + return listHistory(userId, null); + } + + public PatrolDeleteHistoryVo listHistory(Long userId, Integer limit) { long startedAt = System.nanoTime(); validateUserId(userId); + int normalizedLimit = normalizeHistoryLimit(limit); PatrolDeleteHistoryVo vo = new PatrolDeleteHistoryVo(); List entities = fileResultMapper.selectList(new LambdaQueryWrapper() .select(FileResultEntity::getId, @@ -155,7 +162,8 @@ public class PatrolDeleteTaskService { .eq(FileResultEntity::getModuleType, MODULE_TYPE) .eq(FileResultEntity::getUserId, userId) .orderByDesc(FileResultEntity::getCreatedAt) - .last("limit 100")); + .orderByDesc(FileResultEntity::getId) + .last("limit " + normalizedLimit)); long resultRowsLoadedAt = System.nanoTime(); if (entities.isEmpty()) { vo.setItems(List.of()); @@ -481,9 +489,20 @@ public class PatrolDeleteTaskService { throw new BusinessException("记录不存在"); } Long taskId = entity.getTaskId(); - fileResultMapper.deleteById(resultId); - cleanupResultAuxiliaryDataFast(taskId, resultId); - reconcileTaskAfterResultRemoval(taskId); + if (taskId == null || taskId <= 0) { + fileResultMapper.deleteById(resultId); + cleanupResultAuxiliaryDataFast(taskId, resultId); + return; + } + try (TaskDistributedLockService.LockHandle ignored = acquireTaskLockOrThrow(taskId)) { + FileResultEntity latestEntity = fileResultMapper.selectById(resultId); + if (latestEntity == null || !MODULE_TYPE.equals(latestEntity.getModuleType()) || !userId.equals(latestEntity.getUserId())) { + throw new BusinessException("record not found"); + } + fileResultMapper.deleteById(resultId); + cleanupResultAuxiliaryDataFast(taskId, resultId); + reconcileTaskAfterResultRemoval(taskId); + } } private void reconcileTaskAfterResultRemoval(Long taskId) { @@ -560,6 +579,13 @@ public class PatrolDeleteTaskService { return payloadByShop; } + private int normalizeHistoryLimit(Integer limit) { + if (limit == null || limit <= 0) { + return 30; + } + return Math.min(limit, 100); + } + private Map loadTaskMap(List entities) { List taskIds = entities.stream() .map(FileResultEntity::getTaskId) @@ -760,6 +786,9 @@ public class PatrolDeleteTaskService { long successCount = rows.stream().filter(row -> Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())).count(); long failedCount = rows.stream().filter(row -> Integer.valueOf(RESULT_FAILED).equals(row.getSuccess())).count(); long pendingCount = rows.stream().filter(row -> !isResultFinished(row)).count(); + if (pendingCount == 0 && successCount > 0 && isTaskWorkbookPending(task, rows)) { + pendingCount = 1; + } task.setSuccessFileCount((int) successCount); task.setFailedFileCount((int) failedCount); task.setUpdatedAt(LocalDateTime.now()); @@ -975,11 +1004,47 @@ public class PatrolDeleteTaskService { fileResultMapper.updateById(row); } } + updateTaskStatusFromRows(task, rows); + persistSnapshotJson(task, buildSnapshotFromDb(task, rows)); + fileTaskMapper.updateById(task); } finally { FileUtil.del(xlsx); } } + private boolean isTaskWorkbookPending(FileTaskEntity task, List rows) { + if (task == null || task.getId() == null || rows == null || rows.isEmpty()) { + return false; + } + FileResultEntity firstSuccess = rows.stream() + .filter(row -> Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) + .findFirst() + .orElse(null); + if (firstSuccess == null) { + return false; + } + if (!blank(firstSuccess.getResultFileUrl())) { + return false; + } + TaskFileJobEntity job = taskFileJobService.findAssembleJob(task.getId(), MODULE_TYPE, firstSuccess.getId()); + return job != null && !"SUCCESS".equals(job.getStatus()); + } + + private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) { + TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId); + if (lockHandle == null) { + throw new BusinessException(40901, "浠诲姟姝e湪澶勭悊涓紝璇风◢鍚庡啀璇?"); + } + return lockHandle; + } + + private TaskDistributedLockService.LockHandle acquireTaskLock(Long taskId) { + if (taskId == null || taskId <= 0) { + return null; + } + return taskDistributedLockService.acquire(MODULE_TYPE, taskId); + } + private void markResultSuccess(FileResultEntity row) { row.setSuccess(RESULT_SUCCESS); row.setErrorMessage(null); diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackTaskService.java index c12f046..03ef49e 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackTaskService.java @@ -27,6 +27,7 @@ import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.model.entity.FileResultEntity; import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity; import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity; +import com.nanri.aiimage.modules.task.service.TaskDistributedLockService; import com.nanri.aiimage.modules.task.service.TaskFileJobService; import com.nanri.aiimage.modules.task.service.TaskResultPayloadService; import com.nanri.aiimage.modules.ziniao.service.ZiniaoShopSwitchService; @@ -71,6 +72,7 @@ public class PriceTrackTaskService { private final TaskPressureProperties taskPressureProperties; private final TaskResultPayloadService taskResultPayloadService; private final TaskFileJobService taskFileJobService; + private final TaskDistributedLockService taskDistributedLockService; private FileTaskEntity loadTaskForExecution(Long taskId) { Map cachedTasks = priceTrackTaskCacheService.getTaskCacheBatch(List.of(taskId)); @@ -239,8 +241,16 @@ public class PriceTrackTaskService { throw new BusinessException("记录不存在"); } Long taskId = entity.getTaskId(); - fileResultMapper.deleteById(resultId); - if (taskId != null && taskId > 0) { + if (taskId == null || taskId <= 0) { + fileResultMapper.deleteById(resultId); + return; + } + try (TaskDistributedLockService.LockHandle ignored = acquireTaskLockOrThrow(taskId)) { + FileResultEntity latestEntity = fileResultMapper.selectById(resultId); + if (latestEntity == null || !MODULE_TYPE.equals(latestEntity.getModuleType()) || !userId.equals(latestEntity.getUserId())) { + throw new BusinessException("record not found"); + } + fileResultMapper.deleteById(resultId); reconcileTaskAfterResultRemoval(taskId); } } @@ -281,10 +291,18 @@ public class PriceTrackTaskService { FileTaskEntity task = loadTaskForExecution(fr.getTaskId()); if (task == null || !MODULE_TYPE.equals(task.getModuleType()) || !userId.equals(task.getUserId())) continue; if (!"RUNNING".equals(task.getStatus())) continue; - fileResultMapper.deleteById(fr.getId()); - reconcileTaskAfterResultRemoval(fr.getTaskId()); - vo.setRemoved(true); - return vo; + try (TaskDistributedLockService.LockHandle ignored = acquireTaskLock(fr.getTaskId())) { + if (ignored == null) continue; + FileTaskEntity lockedTask = loadTaskForExecution(fr.getTaskId()); + if (lockedTask == null || !MODULE_TYPE.equals(lockedTask.getModuleType()) || !userId.equals(lockedTask.getUserId())) continue; + if (!"RUNNING".equals(lockedTask.getStatus())) continue; + FileResultEntity lockedResult = fileResultMapper.selectById(fr.getId()); + if (lockedResult == null || !MODULE_TYPE.equals(lockedResult.getModuleType())) continue; + fileResultMapper.deleteById(fr.getId()); + reconcileTaskAfterResultRemoval(fr.getTaskId()); + vo.setRemoved(true); + return vo; + } } return vo; } @@ -687,6 +705,15 @@ public class PriceTrackTaskService { throw new BusinessException("结果文件载荷不存在"); } assembleShopResult(result, job.getScopeKey(), payload); + FileTaskEntity task = loadTaskForExecution(job.getTaskId()); + if (task != null && MODULE_TYPE.equals(task.getModuleType())) { + List latest = fileResultMapper.selectList( + new LambdaQueryWrapper() + .eq(FileResultEntity::getTaskId, job.getTaskId()) + .eq(FileResultEntity::getModuleType, MODULE_TYPE) + .orderByAsc(FileResultEntity::getId)); + updateTaskStatusFromLatestRows(task, latest); + } } private void enqueueResultFileAssembly(FileResultEntity result, @@ -1280,10 +1307,20 @@ public class PriceTrackTaskService { int fail = 0; List allErrors = new ArrayList<>(); boolean allDone = true; + Map jobMap = taskFileJobService.findAssembleJobsByResultIds( + MODULE_TYPE, + latest.stream() + .map(FileResultEntity::getId) + .filter(id -> id != null && id > 0) + .toList()); for (FileResultEntity fr : latest) { boolean success = fr.getSuccess() != null && fr.getSuccess() == 1; boolean failed = fr.getErrorMessage() != null && !fr.getErrorMessage().isBlank(); if (success) { + if (isResultAwaitingFileAssembly(fr, jobMap.get(fr.getId()))) { + allDone = false; + continue; + } ok++; } else if (failed) { fail++; @@ -1327,6 +1364,34 @@ public class PriceTrackTaskService { } } + private boolean isResultAwaitingFileAssembly(FileResultEntity result, TaskFileJobEntity job) { + if (result == null) { + return false; + } + if (result.getResultFileUrl() != null && !result.getResultFileUrl().isBlank()) { + return false; + } + if (result.getResultFilename() == null || result.getResultFilename().isBlank()) { + return false; + } + return job == null || !"SUCCESS".equals(job.getStatus()); + } + + private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) { + TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId); + if (lockHandle == null) { + throw new BusinessException(40901, "浠诲姟姝e湪澶勭悊涓紝璇风◢鍚庡啀璇?"); + } + return lockHandle; + } + + private TaskDistributedLockService.LockHandle acquireTaskLock(Long taskId) { + if (taskId == null || taskId <= 0) { + return null; + } + return taskDistributedLockService.acquire(MODULE_TYPE, taskId); + } + private void cleanupTaskCacheIfTerminal(Long taskId, String taskStatus) { if (!"SUCCESS".equals(taskStatus) && !"FAILED".equals(taskStatus) && !"DELETE_EMPTY".equals(taskStatus)) { return; diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/service/ProductRiskTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/service/ProductRiskTaskService.java index aea8c08..e4f15b5 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/service/ProductRiskTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/service/ProductRiskTaskService.java @@ -25,6 +25,7 @@ import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.model.entity.FileResultEntity; import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity; import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity; +import com.nanri.aiimage.modules.task.service.TaskDistributedLockService; import com.nanri.aiimage.modules.task.service.TaskFileJobService; import com.nanri.aiimage.modules.task.service.TaskProgressSnapshotService; import com.nanri.aiimage.modules.task.service.TaskResultItemService; @@ -68,6 +69,7 @@ public class ProductRiskTaskService { private final TaskFileJobService taskFileJobService; private final TaskResultItemService taskResultItemService; private final TaskProgressSnapshotService taskProgressSnapshotService; + private final TaskDistributedLockService taskDistributedLockService; private FileTaskEntity loadTaskForExecution(Long taskId) { Map cachedTasks = productRiskTaskCacheService.getTaskCacheBatch(List.of(taskId)); @@ -243,8 +245,16 @@ public class ProductRiskTaskService { throw new BusinessException("记录不存在"); } Long taskId = entity.getTaskId(); - fileResultMapper.deleteById(resultId); - if (taskId != null && taskId > 0) { + if (taskId == null || taskId <= 0) { + fileResultMapper.deleteById(resultId); + return; + } + try (TaskDistributedLockService.LockHandle ignored = acquireTaskLockOrThrow(taskId)) { + FileResultEntity latestEntity = fileResultMapper.selectById(resultId); + if (latestEntity == null || !MODULE_TYPE.equals(latestEntity.getModuleType()) || !userId.equals(latestEntity.getUserId())) { + throw new BusinessException("record not found"); + } + fileResultMapper.deleteById(resultId); reconcileTaskAfterResultRemoval(taskId); } } @@ -304,10 +314,26 @@ public class ProductRiskTaskService { continue; } Long tid = fr.getTaskId(); - fileResultMapper.deleteById(fr.getId()); - reconcileTaskAfterResultRemoval(tid); - vo.setRemoved(true); - return vo; + try (TaskDistributedLockService.LockHandle ignored = acquireTaskLock(tid)) { + if (ignored == null) { + continue; + } + FileTaskEntity lockedTask = loadTaskForExecution(tid); + if (lockedTask == null || !MODULE_TYPE.equals(lockedTask.getModuleType()) || !userId.equals(lockedTask.getUserId())) { + continue; + } + if (!"RUNNING".equals(lockedTask.getStatus())) { + continue; + } + FileResultEntity lockedResult = fileResultMapper.selectById(fr.getId()); + if (lockedResult == null || !MODULE_TYPE.equals(lockedResult.getModuleType())) { + continue; + } + fileResultMapper.deleteById(fr.getId()); + reconcileTaskAfterResultRemoval(tid); + vo.setRemoved(true); + return vo; + } } return vo; } @@ -631,54 +657,10 @@ public class ProductRiskTaskService { .eq(FileResultEntity::getModuleType, MODULE_TYPE) .orderByAsc(FileResultEntity::getId)); - int ok = 0; - int fail = 0; - List allErrors = new ArrayList<>(); - boolean allDone = true; - for (FileResultEntity fr : latest) { - boolean success = fr.getSuccess() != null && fr.getSuccess() == 1; - boolean failed = fr.getErrorMessage() != null && !fr.getErrorMessage().isBlank(); - if (success) { - ok++; - } else if (failed) { - fail++; - allErrors.add(fr.getSourceFilename() + ": " + fr.getErrorMessage()); - } else { - allDone = false; - } - } - - task.setSuccessFileCount(ok); - task.setFailedFileCount(fail); - task.setUpdatedAt(LocalDateTime.now()); - - if (!allDone) { - task.setStatus("RUNNING"); - task.setErrorMessage(null); - task.setFinishedAt(null); - } else if (ok > 0 && fail == 0) { - task.setStatus("SUCCESS"); - task.setErrorMessage(null); - task.setFinishedAt(LocalDateTime.now()); - } else if (ok > 0) { - task.setStatus("SUCCESS"); - task.setErrorMessage(String.join("; ", allErrors.isEmpty() ? batchErrors : allErrors)); - task.setFinishedAt(LocalDateTime.now()); - } else { - task.setStatus("FAILED"); - task.setErrorMessage(allErrors.isEmpty() ? "全部店铺处理失败" : String.join("; ", allErrors)); - task.setFinishedAt(LocalDateTime.now()); - } - try { - task.setResultJson(objectMapper.writeValueAsString(buildSnapshotFromDb(taskId, task.getStatus()))); - } catch (Exception ex) { - log.warn("[product-risk] compact result json failed: {}", ex.getMessage()); - } - fileTaskMapper.updateById(task); - log.warn("[product-risk] submitResult status evaluated taskId={} oldStatus={} newStatus={} ok={} fail={} allDone={} matchedShopCount={} skippedUnmatchedCount={} waitingCount={} assembledCount={} fallbackMatchedCount={} batchErrors={}", - taskId, oldStatus, task.getStatus(), ok, fail, allDone, matchedShopCount, skippedUnmatchedCount, + updateTaskStatusFromLatestRows(task, latest, batchErrors); + log.warn("[product-risk] submitResult status evaluated taskId={} oldStatus={} newStatus={} matchedShopCount={} skippedUnmatchedCount={} waitingCount={} assembledCount={} fallbackMatchedCount={} batchErrors={}", + taskId, oldStatus, task.getStatus(), matchedShopCount, skippedUnmatchedCount, waitingCount, assembledCount, fallbackMatchedCount, batchErrors); - cleanupTaskCacheIfTerminal(taskId, task.getStatus()); } public boolean tryFinalizeTask(Long taskId, boolean fromCompensation) { @@ -855,6 +837,14 @@ public class ProductRiskTaskService { } File workRoot = FileUtil.mkdir(FileUtil.file(System.getProperty("java.io.tmpdir"), "product-risk-result", String.valueOf(job.getTaskId()))); assembleShopResult(result, job.getScopeKey(), payload, workRoot); + FileTaskEntity task = loadTaskForExecution(job.getTaskId()); + if (task != null && MODULE_TYPE.equals(task.getModuleType())) { + List latest = fileResultMapper.selectList(new LambdaQueryWrapper() + .eq(FileResultEntity::getTaskId, job.getTaskId()) + .eq(FileResultEntity::getModuleType, MODULE_TYPE) + .orderByAsc(FileResultEntity::getId)); + updateTaskStatusFromLatestRows(task, latest, List.of()); + } } private void enqueueResultFileAssembly(FileResultEntity result, String shopKey, ProductRiskShopPayloadDto payload) { @@ -888,10 +878,20 @@ public class ProductRiskTaskService { int fail = 0; List allErrors = new ArrayList<>(); boolean allDone = true; + Map jobMap = taskFileJobService.findAssembleJobsByResultIds( + MODULE_TYPE, + latest.stream() + .map(FileResultEntity::getId) + .filter(id -> id != null && id > 0) + .toList()); for (FileResultEntity fr : latest) { boolean success = fr.getSuccess() != null && fr.getSuccess() == 1; boolean failed = fr.getErrorMessage() != null && !fr.getErrorMessage().isBlank(); if (success) { + if (isResultAwaitingFileAssembly(fr, jobMap.get(fr.getId()))) { + allDone = false; + continue; + } ok++; } else if (failed) { fail++; @@ -947,6 +947,34 @@ public class ProductRiskTaskService { productRiskTaskCacheService.deleteTaskCache(taskId); } + private boolean isResultAwaitingFileAssembly(FileResultEntity result, TaskFileJobEntity job) { + if (result == null) { + return false; + } + if (result.getResultFileUrl() != null && !result.getResultFileUrl().isBlank()) { + return false; + } + if (result.getResultFilename() == null || result.getResultFilename().isBlank()) { + return false; + } + return job == null || !"SUCCESS".equals(job.getStatus()); + } + + private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) { + TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId); + if (lockHandle == null) { + throw new BusinessException(40901, "浠诲姟姝e湪澶勭悊涓紝璇风◢鍚庡啀璇?"); + } + return lockHandle; + } + + private TaskDistributedLockService.LockHandle acquireTaskLock(Long taskId) { + if (taskId == null || taskId <= 0) { + return null; + } + return taskDistributedLockService.acquire(MODULE_TYPE, taskId); + } + private ProductRiskTaskDetailVo buildTaskDetail(FileTaskEntity task) { ProductRiskTaskDetailVo detail = new ProductRiskTaskDetailVo(); detail.setTask(toTaskItemVo(task)); diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/service/QueryAsinTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/service/QueryAsinTaskService.java index 6ef72b0..cac1068 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/service/QueryAsinTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/service/QueryAsinTaskService.java @@ -25,6 +25,7 @@ import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.model.entity.FileResultEntity; import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity; import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity; +import com.nanri.aiimage.modules.task.service.TaskDistributedLockService; import com.nanri.aiimage.modules.task.service.TaskFileJobService; import com.nanri.aiimage.modules.task.service.TaskProgressSnapshotService; import com.nanri.aiimage.modules.task.service.TaskResultItemService; @@ -66,6 +67,7 @@ public class QueryAsinTaskService { private final TaskFileJobService taskFileJobService; private final TaskResultItemService taskResultItemService; private final TaskProgressSnapshotService taskProgressSnapshotService; + private final TaskDistributedLockService taskDistributedLockService; private FileTaskEntity loadTaskForExecution(Long taskId) { Map cachedTasks = taskCacheService.getTaskCacheBatch(List.of(taskId)); @@ -468,8 +470,18 @@ public class QueryAsinTaskService { throw new BusinessException("记录不存在"); } Long taskId = entity.getTaskId(); - fileResultMapper.deleteById(resultId); - reconcileTaskAfterResultRemoval(taskId); + if (taskId == null || taskId <= 0) { + fileResultMapper.deleteById(resultId); + return; + } + try (TaskDistributedLockService.LockHandle ignored = acquireTaskLockOrThrow(taskId)) { + FileResultEntity latestEntity = fileResultMapper.selectById(resultId); + if (latestEntity == null || !MODULE_TYPE.equals(latestEntity.getModuleType()) || !userId.equals(latestEntity.getUserId())) { + throw new BusinessException("record not found"); + } + fileResultMapper.deleteById(resultId); + reconcileTaskAfterResultRemoval(taskId); + } } private void reconcileTaskAfterResultRemoval(Long taskId) { @@ -719,6 +731,9 @@ public class QueryAsinTaskService { long successCount = rows.stream().filter(row -> Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())).count(); long failedCount = rows.stream().filter(row -> Integer.valueOf(RESULT_FAILED).equals(row.getSuccess())).count(); long pendingCount = rows.stream().filter(row -> !isResultFinished(row)).count(); + if (pendingCount == 0 && successCount > 0 && isTaskWorkbookPending(task, rows)) { + pendingCount = 1; + } task.setSuccessFileCount((int) successCount); task.setFailedFileCount((int) failedCount); task.setUpdatedAt(LocalDateTime.now()); @@ -948,11 +963,47 @@ public class QueryAsinTaskService { fileResultMapper.updateById(row); } } + updateTaskStatusFromRows(task, rows); + persistSnapshotJson(task, buildSnapshotFromDb(task, rows)); + fileTaskMapper.updateById(task); } finally { FileUtil.del(xlsx); } } + private boolean isTaskWorkbookPending(FileTaskEntity task, List rows) { + if (task == null || task.getId() == null || rows == null || rows.isEmpty()) { + return false; + } + FileResultEntity firstSuccess = rows.stream() + .filter(row -> Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) + .findFirst() + .orElse(null); + if (firstSuccess == null) { + return false; + } + if (!blank(firstSuccess.getResultFileUrl())) { + return false; + } + TaskFileJobEntity job = taskFileJobService.findAssembleJob(task.getId(), MODULE_TYPE, firstSuccess.getId()); + return job != null && !"SUCCESS".equals(job.getStatus()); + } + + private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) { + TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId); + if (lockHandle == null) { + throw new BusinessException(40901, "浠诲姟姝e湪澶勭悊涓紝璇风◢鍚庡啀璇?"); + } + return lockHandle; + } + + private TaskDistributedLockService.LockHandle acquireTaskLock(Long taskId) { + if (taskId == null || taskId <= 0) { + return null; + } + return taskDistributedLockService.acquire(MODULE_TYPE, taskId); + } + private void markResultSuccess(FileResultEntity row) { row.setSuccess(RESULT_SUCCESS); row.setErrorMessage(null); diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/service/ShopMatchTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/service/ShopMatchTaskService.java index 87fdff0..4389880 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/service/ShopMatchTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/service/ShopMatchTaskService.java @@ -29,6 +29,7 @@ import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.model.entity.FileResultEntity; import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity; import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity; +import com.nanri.aiimage.modules.task.service.TaskDistributedLockService; import com.nanri.aiimage.modules.task.service.TaskFileJobService; import com.nanri.aiimage.modules.task.service.TaskResultPayloadService; import com.nanri.aiimage.modules.ziniao.service.ZiniaoShopSwitchService; @@ -67,6 +68,7 @@ public class ShopMatchTaskService { private final TaskPressureProperties taskPressureProperties; private final TaskResultPayloadService taskResultPayloadService; private final TaskFileJobService taskFileJobService; + private final TaskDistributedLockService taskDistributedLockService; private FileTaskEntity loadTaskForExecution(Long taskId) { Map cachedTasks = shopMatchTaskCacheService.getTaskCacheBatch(List.of(taskId)); @@ -250,8 +252,16 @@ public class ShopMatchTaskService { throw new BusinessException("记录不存在"); } Long taskId = entity.getTaskId(); - fileResultMapper.deleteById(resultId); - if (taskId != null && taskId > 0) { + if (taskId == null || taskId <= 0) { + fileResultMapper.deleteById(resultId); + return; + } + try (TaskDistributedLockService.LockHandle ignored = acquireTaskLockOrThrow(taskId)) { + FileResultEntity latestEntity = fileResultMapper.selectById(resultId); + if (latestEntity == null || !MODULE_TYPE.equals(latestEntity.getModuleType()) || !userId.equals(latestEntity.getUserId())) { + throw new BusinessException("record not found"); + } + fileResultMapper.deleteById(resultId); reconcileTaskAfterResultRemoval(taskId); } } @@ -727,6 +737,14 @@ public class ShopMatchTaskService { } File workRoot = FileUtil.mkdir(FileUtil.file(System.getProperty("java.io.tmpdir"), "shop-match-result", String.valueOf(job.getTaskId()))); assembleShopResult(result, job.getScopeKey(), payload, workRoot); + FileTaskEntity task = loadTaskForExecution(job.getTaskId()); + if (task != null && MODULE_TYPE.equals(task.getModuleType())) { + List latest = fileResultMapper.selectList(new LambdaQueryWrapper() + .eq(FileResultEntity::getTaskId, job.getTaskId()) + .eq(FileResultEntity::getModuleType, MODULE_TYPE) + .orderByAsc(FileResultEntity::getId)); + updateTaskStatusFromLatestRows(task, latest, List.of()); + } } private void enqueueResultFileAssembly(FileResultEntity result, String shopKey, ShopMatchShopPayloadDto payload) { @@ -771,10 +789,20 @@ public class ShopMatchTaskService { int fail = 0; List errors = new ArrayList<>(); boolean allDone = true; + Map jobMap = taskFileJobService.findAssembleJobsByResultIds( + MODULE_TYPE, + latest.stream() + .map(FileResultEntity::getId) + .filter(id -> id != null && id > 0) + .toList()); for (FileResultEntity row : latest) { boolean success = row.getSuccess() != null && row.getSuccess() == 1; boolean failed = row.getErrorMessage() != null && !row.getErrorMessage().isBlank(); if (success) { + if (isResultAwaitingFileAssembly(row, jobMap.get(row.getId()))) { + allDone = false; + continue; + } ok++; } else if (failed) { fail++; @@ -818,6 +846,34 @@ public class ShopMatchTaskService { } } + private boolean isResultAwaitingFileAssembly(FileResultEntity result, TaskFileJobEntity job) { + if (result == null) { + return false; + } + if (result.getResultFileUrl() != null && !result.getResultFileUrl().isBlank()) { + return false; + } + if (result.getResultFilename() == null || result.getResultFilename().isBlank()) { + return false; + } + return job == null || !"SUCCESS".equals(job.getStatus()); + } + + private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) { + TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId); + if (lockHandle == null) { + throw new BusinessException(40901, "浠诲姟姝e湪澶勭悊涓紝璇风◢鍚庡啀璇?"); + } + return lockHandle; + } + + private TaskDistributedLockService.LockHandle acquireTaskLock(Long taskId) { + if (taskId == null || taskId <= 0) { + return null; + } + return taskDistributedLockService.acquire(MODULE_TYPE, taskId); + } + private boolean shouldRemainRunningUntilNextStage(FileTaskEntity task) { ShopMatchCreateTaskRequest request = parseTaskRequestSilently(task); if (request == null || request.getScheduleTimes() == null || request.getScheduleTimes().isEmpty()) { diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/controller/SimilarAsinController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/controller/SimilarAsinController.java index 53de381..71c114e 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/controller/SimilarAsinController.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/controller/SimilarAsinController.java @@ -2,10 +2,13 @@ package com.nanri.aiimage.modules.similarasin.controller; import com.nanri.aiimage.common.api.ApiResponse; import com.nanri.aiimage.common.util.DownloadHeaderUtil; +import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinFilterConditionAddRequest; import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinParseRequest; +import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinParsedPayloadDto; import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinSubmitResultRequest; import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinTaskBatchRequest; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinDashboardVo; +import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinFilterConditionVo; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinHistoryVo; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinParseVo; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinTaskBatchVo; @@ -29,6 +32,7 @@ import org.springframework.web.server.ResponseStatusException; import java.io.InputStream; import java.net.URI; import java.nio.charset.StandardCharsets; +import java.util.List; @RestController @RequiredArgsConstructor @@ -38,12 +42,48 @@ public class SimilarAsinController { private final SimilarAsinTaskService service; + @GetMapping("/filter-conditions") + @Operation(summary = "查询货源查询筛选条件") + public ApiResponse> listFilterConditions( + @Parameter(description = "当前用户 ID", required = true, example = "1") + @RequestParam("user_id") Long userId) { + return ApiResponse.success(service.listFilterConditions(userId)); + } + + @PostMapping("/filter-conditions") + @Operation(summary = "保存货源查询筛选条件") + public ApiResponse addFilterCondition( + @Valid @RequestBody SimilarAsinFilterConditionAddRequest request) { + return ApiResponse.success(service.addFilterCondition(request)); + } + + @DeleteMapping("/filter-conditions/{id}") + @Operation(summary = "删除货源查询筛选条件") + public ApiResponse deleteFilterCondition( + @Parameter(description = "筛选条件主键", example = "10") + @PathVariable Long id, + @Parameter(description = "当前用户 ID", required = true, example = "1") + @RequestParam("user_id") Long userId) { + service.deleteFilterCondition(userId, id); + return ApiResponse.success(null); + } + @PostMapping("/parse") @Operation(summary = "解析 Excel 并创建任务", description = "解析上传后的 Excel 文件,提取 id、ASIN、国家、URL、标题等字段。返回给前端的数据只包含整数 id 和 n_1 行;n_2、n_3 等子行会保存在 OSS 解析载荷中,用于最终结果补齐。创建后的任务状态为 PENDING,不会自动推送 Python。") public ApiResponse parse(@Valid @RequestBody SimilarAsinParseRequest request) { return ApiResponse.success(service.parseAndCreateTask(request)); } + @GetMapping("/tasks/{taskId}/parsed-payload") + @Operation(summary = "获取货源查询完整解析载荷", description = "给 Python 队列消费端使用。parse 接口只返回轻量预览,大文件完整行数据通过该接口按 taskId 拉取。") + public ApiResponse parsedPayload( + @Parameter(description = "货源查询任务 ID", required = true, example = "7004") + @PathVariable Long taskId, + @Parameter(description = "当前用户 ID", required = true, example = "1") + @RequestParam("user_id") Long userId) { + return ApiResponse.success(service.parsedPayload(taskId, userId)); + } + @GetMapping("/dashboard") @Operation(summary = "查询相似ASIN检测总览", description = "查询当前用户的运行中、成功、失败和已结束任务数量。页面进入时请求一次即可,不需要持续轮询。") public ApiResponse dashboard( diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/mapper/SimilarAsinFilterConditionMapper.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/mapper/SimilarAsinFilterConditionMapper.java new file mode 100644 index 0000000..9b27a06 --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/mapper/SimilarAsinFilterConditionMapper.java @@ -0,0 +1,9 @@ +package com.nanri.aiimage.modules.similarasin.mapper; + +import com.baomidou.mybatisplus.core.mapper.BaseMapper; +import com.nanri.aiimage.modules.similarasin.model.entity.SimilarAsinFilterConditionEntity; +import org.apache.ibatis.annotations.Mapper; + +@Mapper +public interface SimilarAsinFilterConditionMapper extends BaseMapper { +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinFilterConditionAddRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinFilterConditionAddRequest.java new file mode 100644 index 0000000..e265712 --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinFilterConditionAddRequest.java @@ -0,0 +1,24 @@ +package com.nanri.aiimage.modules.similarasin.model.dto; + +import com.fasterxml.jackson.annotation.JsonAlias; +import com.fasterxml.jackson.annotation.JsonProperty; +import io.swagger.v3.oas.annotations.media.Schema; +import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.NotNull; +import lombok.Data; + +@Data +@Schema(description = "货源查询筛选条件新增请求") +public class SimilarAsinFilterConditionAddRequest { + + @NotNull + @JsonProperty("user_id") + @Schema(description = "当前用户 ID", example = "1", requiredMode = Schema.RequiredMode.REQUIRED) + private Long userId; + + @NotBlank + @JsonProperty("condition_text") + @JsonAlias({"conditionText", "filter_condition", "filterCondition"}) + @Schema(description = "筛选条件文本", example = "只筛选有货且不带品牌 logo 的商品", requiredMode = Schema.RequiredMode.REQUIRED) + private String conditionText; +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/entity/SimilarAsinFilterConditionEntity.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/entity/SimilarAsinFilterConditionEntity.java new file mode 100644 index 0000000..39f4c93 --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/entity/SimilarAsinFilterConditionEntity.java @@ -0,0 +1,19 @@ +package com.nanri.aiimage.modules.similarasin.model.entity; + +import com.baomidou.mybatisplus.annotation.IdType; +import com.baomidou.mybatisplus.annotation.TableId; +import com.baomidou.mybatisplus.annotation.TableName; +import lombok.Data; + +import java.time.LocalDateTime; + +@Data +@TableName("biz_similar_asin_filter_condition") +public class SimilarAsinFilterConditionEntity { + + @TableId(type = IdType.AUTO) + private Long id; + private Long userId; + private String conditionText; + private LocalDateTime createdAt; +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/vo/SimilarAsinFilterConditionVo.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/vo/SimilarAsinFilterConditionVo.java new file mode 100644 index 0000000..7c38f2f --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/vo/SimilarAsinFilterConditionVo.java @@ -0,0 +1,17 @@ +package com.nanri.aiimage.modules.similarasin.model.vo; + +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Data; + +import java.time.LocalDateTime; + +@Data +public class SimilarAsinFilterConditionVo { + + private Long id; + + @JsonProperty("conditionText") + private String conditionText; + + private LocalDateTime createdAt; +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java index fb268c1..9279684 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java @@ -11,14 +11,18 @@ import com.nanri.aiimage.common.service.DistributedJobLockService; import com.nanri.aiimage.config.InstanceMetadata; import com.nanri.aiimage.config.SimilarAsinProperties; import com.nanri.aiimage.config.StorageProperties; +import com.nanri.aiimage.modules.similarasin.mapper.SimilarAsinFilterConditionMapper; import com.nanri.aiimage.modules.similarasin.client.SimilarAsinCozeClient; +import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinFilterConditionAddRequest; import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinParseRequest; import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinParsedPayloadDto; import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinResultGroupDto; import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinResultRowDto; import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinSourceFileDto; import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinSubmitResultRequest; +import com.nanri.aiimage.modules.similarasin.model.entity.SimilarAsinFilterConditionEntity; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinDashboardVo; +import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinFilterConditionVo; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinHistoryItemVo; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinHistoryVo; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinParsedGroupVo; @@ -102,6 +106,7 @@ public class SimilarAsinTaskService { private static final long RESULT_ROWS_READ_RETRY_DELAY_MS = 500L; private static final Duration TASK_LOCK_TTL = Duration.ofMinutes(5); private static final long TASK_LOCK_WAIT_MILLIS = 10000L; + private static final int PARSE_RESPONSE_PREVIEW_LIMIT = 100; private static final List RESULT_HEADERS = List.of( "id", "asin", @@ -118,6 +123,7 @@ public class SimilarAsinTaskService { private final FileResultMapper fileResultMapper; private final TaskScopeStateMapper taskScopeStateMapper; private final TaskChunkMapper taskChunkMapper; + private final SimilarAsinFilterConditionMapper filterConditionMapper; private final ObjectMapper objectMapper; private final SimilarAsinCozeClient cozeClient; private final SimilarAsinTaskCacheService taskCacheService; @@ -133,6 +139,55 @@ public class SimilarAsinTaskService { @Qualifier("cozeTaskExecutor") private TaskExecutor cozeTaskExecutor; + public List listFilterConditions(Long userId) { + validateUserId(userId); + List rows = filterConditionMapper.selectList( + new LambdaQueryWrapper() + .eq(SimilarAsinFilterConditionEntity::getUserId, userId) + .orderByDesc(SimilarAsinFilterConditionEntity::getId)); + List list = new ArrayList<>(); + for (SimilarAsinFilterConditionEntity row : rows) { + list.add(toFilterConditionVo(row)); + } + return list; + } + + @Transactional + public SimilarAsinFilterConditionVo addFilterCondition(SimilarAsinFilterConditionAddRequest request) { + validateUserId(request.getUserId()); + String text = normalize(request.getConditionText()); + if (text.isBlank()) { + throw new BusinessException("筛选条件不能为空"); + } + SimilarAsinFilterConditionEntity existing = filterConditionMapper.selectOne( + new LambdaQueryWrapper() + .eq(SimilarAsinFilterConditionEntity::getUserId, request.getUserId()) + .eq(SimilarAsinFilterConditionEntity::getConditionText, text) + .last("limit 1")); + if (existing != null) { + return toFilterConditionVo(existing); + } + SimilarAsinFilterConditionEntity entity = new SimilarAsinFilterConditionEntity(); + entity.setUserId(request.getUserId()); + entity.setConditionText(text); + entity.setCreatedAt(LocalDateTime.now()); + filterConditionMapper.insert(entity); + return toFilterConditionVo(entity); + } + + @Transactional + public void deleteFilterCondition(Long userId, Long id) { + validateUserId(userId); + if (id == null || id <= 0) { + throw new BusinessException("id 不合法"); + } + SimilarAsinFilterConditionEntity row = filterConditionMapper.selectById(id); + if (row == null || !userId.equals(row.getUserId())) { + throw new BusinessException("记录不存在"); + } + filterConditionMapper.deleteById(id); + } + public SimilarAsinParseVo parseAndCreateTask(SimilarAsinParseRequest request) { long startedAt = System.nanoTime(); if (request.getUserId() == null || request.getUserId() <= 0) { @@ -240,8 +295,8 @@ public class SimilarAsinTaskService { vo.setDroppedRows(droppedRows); vo.setGroupCount(groups.size()); vo.setAiPrompt(normalize(request.getAiPrompt())); - vo.setItems(allRows); - vo.setGroups(groups); + vo.setItems(buildResponsePreviewRows(allRows)); + vo.setGroups(List.of()); long finishedAt = System.nanoTime(); log.info("[similar-asin] parse timing taskId={} files={} rows={} groups={} totalMs={} parseMs={} groupMs={} taskInsertMs={} payloadJsonMs={} payloadStoreMs={} persistMs={} responseMs={}", task.getId(), @@ -922,6 +977,36 @@ public class SimilarAsinTaskService { return groups; } + private List buildResponsePreviewRows(List rows) { + if (rows == null || rows.isEmpty()) { + return List.of(); + } + int limit = Math.min(PARSE_RESPONSE_PREVIEW_LIMIT, rows.size()); + List preview = new ArrayList<>(limit); + for (int i = 0; i < limit; i++) { + preview.add(copyPreviewRow(rows.get(i))); + } + return preview; + } + + private SimilarAsinParsedRowVo copyPreviewRow(SimilarAsinParsedRowVo row) { + SimilarAsinParsedRowVo vo = new SimilarAsinParsedRowVo(); + vo.setSourceFileKey(row.getSourceFileKey()); + vo.setSourceFilename(row.getSourceFilename()); + vo.setRowIndex(row.getRowIndex()); + vo.setSourceId(row.getSourceId()); + vo.setDisplayId(row.getDisplayId()); + vo.setRowToken(row.getRowToken()); + vo.setGroupKey(row.getGroupKey()); + vo.setAsin(row.getAsin()); + vo.setCountry(row.getCountry()); + vo.setPrice(row.getPrice()); + vo.setUrl(row.getUrl()); + vo.setTitle(row.getTitle()); + vo.setValues(new LinkedHashMap<>()); + return vo; + } + private Map> groupRowsByBaseId(List rows) { Map> result = new LinkedHashMap<>(); if (rows == null) { @@ -2228,6 +2313,28 @@ public class SimilarAsinTaskService { return preferred == null || preferred.isBlank() ? fallback : preferred.trim(); } + private void validateUserId(Long userId) { + if (userId == null || userId <= 0) { + throw new BusinessException("user_id 不合法"); + } + } + + private SimilarAsinFilterConditionVo toFilterConditionVo(SimilarAsinFilterConditionEntity row) { + SimilarAsinFilterConditionVo vo = new SimilarAsinFilterConditionVo(); + vo.setId(row.getId()); + vo.setConditionText(row.getConditionText()); + vo.setCreatedAt(row.getCreatedAt()); + return vo; + } + + public SimilarAsinParsedPayloadDto parsedPayload(Long taskId, Long userId) { + FileTaskEntity task = fileTaskMapper.selectById(taskId); + if (task == null || !MODULE_TYPE.equals(task.getModuleType()) || userId != null && !Objects.equals(userId, task.getUserId())) { + throw new BusinessException("任务不存在"); + } + return readParsedPayload(task); + } + private String normalize(String val) { return val == null ? "" : val.replace(String.valueOf((char) 0xFEFF), "").replace((char) 0x3000, ' ').trim().replaceAll("\\s+", " "); } diff --git a/backend-java/src/main/resources/db/V46__patrol_delete_delete_and_history_indexes.sql b/backend-java/src/main/resources/db/V46__patrol_delete_delete_and_history_indexes.sql new file mode 100644 index 0000000..8a63fae --- /dev/null +++ b/backend-java/src/main/resources/db/V46__patrol_delete_delete_and_history_indexes.sql @@ -0,0 +1,35 @@ +SET @schema_name = DATABASE(); + +SET @idx_file_result_patrol_history_desc_exists := ( + SELECT COUNT(1) + FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = @schema_name + AND TABLE_NAME = 'biz_file_result' + AND INDEX_NAME = 'idx_file_result_patrol_history_desc' +); + +SET @sql_add_file_result_patrol_history_desc_idx := IF( + @idx_file_result_patrol_history_desc_exists = 0, + 'ALTER TABLE biz_file_result ADD INDEX idx_file_result_patrol_history_desc (module_type, user_id, created_at DESC, id DESC)', + 'SELECT 1' +); +PREPARE stmt_add_file_result_patrol_history_desc_idx FROM @sql_add_file_result_patrol_history_desc_idx; +EXECUTE stmt_add_file_result_patrol_history_desc_idx; +DEALLOCATE PREPARE stmt_add_file_result_patrol_history_desc_idx; + +SET @idx_task_file_job_patrol_history_lookup_exists := ( + SELECT COUNT(1) + FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = @schema_name + AND TABLE_NAME = 'biz_task_file_job' + AND INDEX_NAME = 'idx_task_file_job_patrol_history_lookup' +); + +SET @sql_add_task_file_job_patrol_history_lookup_idx := IF( + @idx_task_file_job_patrol_history_lookup_exists = 0, + 'ALTER TABLE biz_task_file_job ADD INDEX idx_task_file_job_patrol_history_lookup (module_type, job_type, result_id, status)', + 'SELECT 1' +); +PREPARE stmt_add_task_file_job_patrol_history_lookup_idx FROM @sql_add_task_file_job_patrol_history_lookup_idx; +EXECUTE stmt_add_task_file_job_patrol_history_lookup_idx; +DEALLOCATE PREPARE stmt_add_task_file_job_patrol_history_lookup_idx; diff --git a/backend-java/src/main/resources/db/V47__similar_asin_filter_condition.sql b/backend-java/src/main/resources/db/V47__similar_asin_filter_condition.sql new file mode 100644 index 0000000..5535b1b --- /dev/null +++ b/backend-java/src/main/resources/db/V47__similar_asin_filter_condition.sql @@ -0,0 +1,15 @@ +CREATE TABLE IF NOT EXISTS biz_similar_asin_filter_condition ( + id BIGINT PRIMARY KEY AUTO_INCREMENT, + user_id BIGINT NOT NULL, + condition_text VARCHAR(500) NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + UNIQUE KEY uk_similar_asin_filter_condition_user_text (user_id, condition_text), + KEY idx_similar_asin_filter_condition_user_id (user_id) +); + +UPDATE `columns` +SET `name` = 'source query', + `menu_type` = 'app', + `route_path` = 'similar-asin', + `sort_order` = 112 +WHERE `column_key` = 'similar_asin'; diff --git a/frontend-vue/similar-asin.html b/frontend-vue/similar-asin.html index 79028f4..8665435 100644 --- a/frontend-vue/similar-asin.html +++ b/frontend-vue/similar-asin.html @@ -3,7 +3,7 @@ - 相似ASIN检测 + 货源查询
diff --git a/frontend-vue/src/pages/brand/components/BrandApiSecretSettingsButton.vue b/frontend-vue/src/pages/brand/components/BrandApiSecretSettingsButton.vue new file mode 100644 index 0000000..910c69a --- /dev/null +++ b/frontend-vue/src/pages/brand/components/BrandApiSecretSettingsButton.vue @@ -0,0 +1,337 @@ + + + + + diff --git a/frontend-vue/src/pages/brand/components/BrandAppearancePatentTab.vue b/frontend-vue/src/pages/brand/components/BrandAppearancePatentTab.vue index f09b551..0533283 100644 --- a/frontend-vue/src/pages/brand/components/BrandAppearancePatentTab.vue +++ b/frontend-vue/src/pages/brand/components/BrandAppearancePatentTab.vue @@ -24,11 +24,6 @@
{{ defaultAiPrompt }}
-
-
密钥
- -
-
-
AI 提示词
-