处理后台新需求
Build Backend JAR / build (push) Failing after 16m58s

This commit is contained in:
supernijia
2026-08-09 20:08:39 +08:00
parent f7e26ee266
commit 38216a1885
42 changed files with 3425 additions and 424 deletions
@@ -12,7 +12,7 @@ public class ModuleCleanupProperties {
private boolean enabled = true;
private String cron = "0 0 0 * * *";
private long retentionDays = 7;
// SHOP_DATA_CRAWL is governed by per-shop latest-three retention in its
// task service and must not be removed by the age-based sweep.
// SHOP_DATA_CRAWL keeps one per-shop daily workbook in its task service
// and must not be removed by the age-based sweep.
private List<String> moduleTypes = new ArrayList<>(List.of("DEDUPE", "SPLIT", "CONVERT", "DELETE_BRAND", "PRODUCT_RISK_RESOLVE", "PRICE_TRACK", "SHOP_MATCH", "PATROL_DELETE", "QUERY_ASIN", "WITHDRAW", "APPEARANCE_PATENT", "SIMILAR_ASIN", "COLLECT_DATA"));
}
@@ -805,6 +805,7 @@ public class CollectDataService {
InvalidAsinDataEntity entity = new InvalidAsinDataEntity();
entity.setDataValue(row.getAsin());
entity.setBrand(brand);
entity.setRecordSource("AUTO");
try {
invalidAsinDataMapper.insert(entity);
} catch (DuplicateKeyException ignored) {
@@ -1,6 +1,8 @@
package com.nanri.aiimage.modules.dedupe.controller;
import com.nanri.aiimage.common.api.ApiResponse;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.modules.admin.support.AdminAuthSupport;
import com.nanri.aiimage.common.util.DownloadHeaderUtil;
import com.nanri.aiimage.modules.dedupe.model.dto.DedupeTotalDataCreateRequest;
import com.nanri.aiimage.modules.dedupe.model.dto.DedupeTotalDataUpdateRequest;
@@ -10,6 +12,8 @@ import com.nanri.aiimage.modules.dedupe.model.vo.DedupeTotalDataImportVo;
import com.nanri.aiimage.modules.dedupe.model.vo.DedupeTotalDataItemVo;
import com.nanri.aiimage.modules.dedupe.model.vo.DedupeTotalDataPageVo;
import com.nanri.aiimage.modules.dedupe.service.DedupeTotalDataService;
import com.nanri.aiimage.modules.permission.model.entity.AdminUserEntity;
import com.nanri.aiimage.modules.permission.service.PermissionMenuService;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.Parameter;
import io.swagger.v3.oas.annotations.media.Content;
@@ -17,6 +21,8 @@ import io.swagger.v3.oas.annotations.media.Schema;
import io.swagger.v3.oas.annotations.responses.ApiResponses;
import io.swagger.v3.oas.annotations.tags.Tag;
import jakarta.validation.Valid;
import jakarta.servlet.http.HttpServletRequest;
import org.springframework.beans.factory.annotation.Value;
import lombok.RequiredArgsConstructor;
import org.springframework.format.annotation.DateTimeFormat;
import org.springframework.http.HttpHeaders;
@@ -36,6 +42,10 @@ import org.springframework.web.multipart.MultipartFile;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.security.MessageDigest;
@RestController
@RequiredArgsConstructor
@@ -43,9 +53,20 @@ import java.time.format.DateTimeFormatter;
@Tag(name = "数据去重总数据", description = "维护数据去重模块的总数据列表,支持增删改查。")
public class DedupeTotalDataController {
private static final String DEDUPE_TOTAL_DATA_COLUMN_KEY = "admin_dedupe_total_data";
private static final String DEDUPE_TOTAL_DATA_ROUTE_PATH = "dedupe-total-data";
private static final DateTimeFormatter EXPORT_FILENAME_FORMATTER = DateTimeFormatter.ofPattern("yyyyMMddHHmmss");
private final DedupeTotalDataService dedupeTotalDataService;
private final AdminAuthSupport adminAuthSupport;
private final PermissionMenuService permissionMenuService;
@Value("${aiimage.security.internal-token:}")
private String internalToken;
@Value("${aiimage.security.internal-token-file:}")
private String internalTokenFile;
@GetMapping
@Operation(summary = "分页查询总数据", description = "分页查询数据去重总数据,支持按值模糊搜索。")
@@ -61,9 +82,11 @@ public class DedupeTotalDataController {
@RequestParam(required = false) @DateTimeFormat(iso = DateTimeFormat.ISO.DATE) LocalDate startDate,
@Parameter(description = "结束日期(包含)")
@RequestParam(required = false) @DateTimeFormat(iso = DateTimeFormat.ISO.DATE) LocalDate endDate,
@Parameter(description = "当前操作人用户ID") @RequestParam Long operatorId) {
@Parameter(description = "分组 ID") @RequestParam(required = false) Long groupId,
HttpServletRequest request) {
RequestOperator operator = requireDedupeTotalDataAccess(request);
return ApiResponse.success(dedupeTotalDataService.page(
page, pageSize, keyword, username, startDate, endDate, operatorId));
page, pageSize, keyword, username, startDate, endDate, groupId, operator.id()));
}
@GetMapping("/export")
@@ -74,8 +97,10 @@ public class DedupeTotalDataController {
@RequestParam(required = false) @DateTimeFormat(iso = DateTimeFormat.ISO.DATE) LocalDate startDate,
@Parameter(description = "结束日期(包含)")
@RequestParam(required = false) @DateTimeFormat(iso = DateTimeFormat.ISO.DATE) LocalDate endDate,
@Parameter(description = "当前操作人用户ID") @RequestParam Long operatorId) {
byte[] bytes = dedupeTotalDataService.export(username, startDate, endDate, operatorId);
@Parameter(description = "分组 ID") @RequestParam(required = false) Long groupId,
HttpServletRequest request) {
RequestOperator operator = requireDedupeTotalDataAccess(request);
byte[] bytes = dedupeTotalDataService.export(username, startDate, endDate, groupId, operator.id());
String filename = "dedupe-total-data-" + LocalDateTime.now().format(EXPORT_FILENAME_FORMATTER) + ".xlsx";
return ResponseEntity.ok()
.header(HttpHeaders.CONTENT_DISPOSITION, DownloadHeaderUtil.contentDisposition(filename))
@@ -91,9 +116,10 @@ public class DedupeTotalDataController {
@io.swagger.v3.oas.annotations.responses.ApiResponse(responseCode = "400", description = "参数不合法或数据重复")
})
public ApiResponse<DedupeTotalDataItemVo> create(
@Parameter(description = "当前操作人用户ID") @RequestParam Long operatorId,
HttpServletRequest httpRequest,
@Valid @RequestBody DedupeTotalDataCreateRequest request) {
return ApiResponse.success("创建成功", dedupeTotalDataService.create(request, operatorId));
RequestOperator operator = requireDedupeTotalDataAccess(httpRequest);
return ApiResponse.success("创建成功", dedupeTotalDataService.create(request, operator.id()));
}
@PostMapping("/import")
@@ -104,16 +130,19 @@ public class DedupeTotalDataController {
})
public ApiResponse<DedupeTotalDataImportStartVo> importExcel(
@Parameter(description = "xlsx 文件", required = true) @RequestParam("file") MultipartFile file,
@Parameter(description = "当前操作人用户ID") @RequestParam Long operatorId) {
return ApiResponse.success("开始导入", dedupeTotalDataService.startImport(file, operatorId));
@Parameter(description = "分组 ID", required = true) @RequestParam Long groupId,
HttpServletRequest request) {
RequestOperator operator = requireDedupeTotalDataAccess(request);
return ApiResponse.success("开始导入", dedupeTotalDataService.startImport(file, groupId, operator.id()));
}
@GetMapping("/import/{importId}")
@Operation(summary = "查询导入进度", description = "根据导入任务 ID 查询当前进度。")
public ApiResponse<DedupeTotalDataImportProgressVo> importProgress(
@PathVariable String importId,
@Parameter(description = "当前操作人用户ID") @RequestParam Long operatorId) {
return ApiResponse.success(dedupeTotalDataService.getImportProgress(importId, operatorId));
HttpServletRequest request) {
RequestOperator operator = requireDedupeTotalDataAccess(request);
return ApiResponse.success(dedupeTotalDataService.getImportProgress(importId, operator.id()));
}
@PostMapping("/delete-import")
@@ -124,16 +153,19 @@ public class DedupeTotalDataController {
})
public ApiResponse<DedupeTotalDataImportStartVo> deleteImportExcel(
@Parameter(description = "xlsx 文件", required = true) @RequestParam("file") MultipartFile file,
@Parameter(description = "当前操作人用户ID") @RequestParam Long operatorId) {
return ApiResponse.success("开始删除", dedupeTotalDataService.startDeleteImport(file, operatorId));
@Parameter(description = "分组 ID", required = true) @RequestParam Long groupId,
HttpServletRequest request) {
RequestOperator operator = requireDedupeTotalDataAccess(request);
return ApiResponse.success("开始删除", dedupeTotalDataService.startDeleteImport(file, groupId, operator.id()));
}
@GetMapping("/delete-import/{importId}")
@Operation(summary = "查询删除导入进度", description = "根据删除导入任务 ID 查询当前进度。")
public ApiResponse<DedupeTotalDataImportProgressVo> deleteImportProgress(
@PathVariable String importId,
@Parameter(description = "当前操作人用户ID") @RequestParam Long operatorId) {
return ApiResponse.success(dedupeTotalDataService.getDeleteImportProgress(importId, operatorId));
HttpServletRequest request) {
RequestOperator operator = requireDedupeTotalDataAccess(request);
return ApiResponse.success(dedupeTotalDataService.getDeleteImportProgress(importId, operator.id()));
}
@PutMapping("/{id}")
@@ -145,9 +177,10 @@ public class DedupeTotalDataController {
})
public ApiResponse<DedupeTotalDataItemVo> update(
@Parameter(description = "主键ID", required = true) @PathVariable Long id,
@Parameter(description = "当前操作人用户ID") @RequestParam Long operatorId,
HttpServletRequest httpRequest,
@Valid @RequestBody DedupeTotalDataUpdateRequest request) {
return ApiResponse.success("更新成功", dedupeTotalDataService.update(id, request, operatorId));
RequestOperator operator = requireDedupeTotalDataAccess(httpRequest);
return ApiResponse.success("更新成功", dedupeTotalDataService.update(id, request, operator.id()));
}
@DeleteMapping("/{id}")
@@ -158,8 +191,105 @@ public class DedupeTotalDataController {
})
public ApiResponse<Void> delete(
@Parameter(description = "主键ID", required = true) @PathVariable Long id,
@Parameter(description = "当前操作人用户ID") @RequestParam Long operatorId) {
dedupeTotalDataService.delete(id, operatorId);
HttpServletRequest request) {
RequestOperator operator = requireDedupeTotalDataAccess(request);
dedupeTotalDataService.delete(id, operator.id());
return ApiResponse.success("删除成功", null);
}
private RequestOperator requireDedupeTotalDataAccess(HttpServletRequest request) {
AdminUserEntity operator = resolveOperator(request);
if (operator == null || operator.getId() == null || operator.getId() <= 0) {
throw new BusinessException(403, "无权访问数据去重总数据");
}
boolean superAdmin = "super_admin".equals(adminAuthSupport.currentRole(operator));
if (!superAdmin && !hasDedupeTotalDataPermission(operator)) {
throw new BusinessException(403, "无权访问数据去重总数据");
}
return new RequestOperator(operator.getId(), superAdmin);
}
private AdminUserEntity resolveOperator(HttpServletRequest request) {
try {
return adminAuthSupport.requireUser(request);
} catch (BusinessException authFailure) {
AdminUserEntity internalOperator = resolveInternalOperator(request);
if (internalOperator != null) {
return internalOperator;
}
throw authFailure;
}
}
private boolean hasDedupeTotalDataPermission(AdminUserEntity operator) {
return permissionMenuService.getUserColumnPermissions(operator.getId(), "admin")
.stream()
.anyMatch(item -> DEDUPE_TOTAL_DATA_COLUMN_KEY.equals(item.getColumnKey())
|| DEDUPE_TOTAL_DATA_ROUTE_PATH.equals(item.getRoutePath()));
}
private AdminUserEntity resolveInternalOperator(HttpServletRequest request) {
if (!isTrustedInternalRequest(request.getHeader("X-Internal-Token"))) {
return null;
}
String rawOperatorId = request.getParameter("operatorId");
if (rawOperatorId == null || rawOperatorId.isBlank()) {
rawOperatorId = request.getParameter("operator_id");
}
if (rawOperatorId == null || rawOperatorId.isBlank()) {
return null;
}
try {
return permissionMenuService.requireUserOperator(Long.parseLong(rawOperatorId.trim()));
} catch (NumberFormatException | BusinessException ex) {
return null;
}
}
private boolean isTrustedInternalRequest(String suppliedToken) {
String expectedToken = resolveExpectedInternalToken();
return !expectedToken.isBlank() && suppliedToken != null && !suppliedToken.isBlank()
&& MessageDigest.isEqual(
expectedToken.getBytes(StandardCharsets.UTF_8),
suppliedToken.trim().getBytes(StandardCharsets.UTF_8));
}
private String resolveExpectedInternalToken() {
if (internalToken != null && !internalToken.isBlank()) {
return internalToken.trim();
}
Path path = resolveInternalTokenFile();
if (path == null || !Files.isRegularFile(path)) {
return "";
}
try {
return Files.readString(path, StandardCharsets.UTF_8).trim();
} catch (Exception ignored) {
return "";
}
}
private Path resolveInternalTokenFile() {
String configuredPath = internalTokenFile == null ? "" : internalTokenFile.trim();
if (!configuredPath.isEmpty()) {
if (configuredPath.equals("~") || configuredPath.startsWith("~/") || configuredPath.startsWith("~\\")) {
String userHome = System.getProperty("user.home", "").trim();
if (userHome.isEmpty()) {
return null;
}
configuredPath = configuredPath.length() == 1
? userHome
: Path.of(userHome, configuredPath.substring(2)).toString();
}
Path configuredTokenPath = Path.of(configuredPath);
return configuredTokenPath.isAbsolute() ? configuredTokenPath.normalize() : null;
}
String userHome = System.getProperty("user.home", "").trim();
return userHome.isEmpty()
? null
: Path.of(userHome, ".aiimage", "internal-token").toAbsolutePath().normalize();
}
private record RequestOperator(Long id, boolean superAdmin) {
}
}
@@ -2,6 +2,7 @@ package com.nanri.aiimage.modules.dedupe.model.dto;
import io.swagger.v3.oas.annotations.media.Schema;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import jakarta.validation.constraints.Size;
import lombok.Data;
@@ -13,4 +14,8 @@ public class DedupeTotalDataCreateRequest {
@Size(max = 128, message = "数据长度不能超过128个字符")
@Schema(description = "总数据值", requiredMode = Schema.RequiredMode.REQUIRED)
private String dataValue;
@NotNull(message = "请选择分组")
@Schema(description = "分组 ID", requiredMode = Schema.RequiredMode.REQUIRED)
private Long groupId;
}
@@ -2,6 +2,7 @@ package com.nanri.aiimage.modules.dedupe.model.dto;
import io.swagger.v3.oas.annotations.media.Schema;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import jakarta.validation.constraints.Size;
import lombok.Data;
@@ -13,4 +14,8 @@ public class DedupeTotalDataUpdateRequest {
@Size(max = 128, message = "数据长度不能超过128个字符")
@Schema(description = "总数据值", requiredMode = Schema.RequiredMode.REQUIRED)
private String dataValue;
@NotNull(message = "请选择分组")
@Schema(description = "分组 ID", requiredMode = Schema.RequiredMode.REQUIRED)
private Long groupId;
}
@@ -14,6 +14,7 @@ public class DedupeTotalDataEntity {
@TableId(type = IdType.AUTO)
private Long id;
private String dataValue;
private Long groupId;
private Long uploaderUserId;
private String uploaderUsername;
private LocalDateTime createdAt;
@@ -15,6 +15,12 @@ public class DedupeTotalDataItemVo {
@Schema(description = "总数据值")
private String dataValue;
@Schema(description = "分组 ID")
private Long groupId;
@Schema(description = "分组名称")
private String groupName;
@Schema(description = "上传用户ID")
private Long uploaderUserId;
@@ -15,6 +15,7 @@ import com.nanri.aiimage.modules.dedupe.model.vo.DedupeTotalDataPageVo;
import com.nanri.aiimage.modules.permission.mapper.AdminUserMapper;
import com.nanri.aiimage.modules.permission.model.entity.AdminUserEntity;
import com.nanri.aiimage.modules.shopkey.mapper.ShopManageGroupMapper;
import com.nanri.aiimage.modules.shopkey.model.entity.ShopManageGroupEntity;
import lombok.RequiredArgsConstructor;
import org.apache.poi.ss.usermodel.Cell;
import org.apache.poi.ss.usermodel.DataFormatter;
@@ -42,6 +43,7 @@ import java.time.format.DateTimeFormatter;
import java.util.Collection;
import java.util.HashMap;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Locale;
@@ -65,6 +67,8 @@ public class DedupeTotalDataService {
private final Map<String, DedupeTotalDataImportProgressVo> deleteImportProgressMap = new ConcurrentHashMap<>();
private final Map<String, Long> importOwnerMap = new ConcurrentHashMap<>();
private final Map<String, Long> deleteImportOwnerMap = new ConcurrentHashMap<>();
private final Map<String, Long> importGroupMap = new ConcurrentHashMap<>();
private final Map<String, Long> deleteImportGroupMap = new ConcurrentHashMap<>();
private final Map<String, Long> importCompletedAtMap = new ConcurrentHashMap<>();
private final Map<String, Long> deleteImportCompletedAtMap = new ConcurrentHashMap<>();
@@ -76,6 +80,11 @@ public class DedupeTotalDataService {
public DedupeTotalDataPageVo page(long page, long pageSize, String keyword, String username,
LocalDate startDate, LocalDate endDate, Long operatorId) {
return page(page, pageSize, keyword, username, startDate, endDate, null, operatorId);
}
public DedupeTotalDataPageVo page(long page, long pageSize, String keyword, String username,
LocalDate startDate, LocalDate endDate, Long groupId, Long operatorId) {
if (startDate != null && endDate != null && startDate.isAfter(endDate)) {
throw new BusinessException("开始日期不能晚于结束日期");
}
@@ -86,17 +95,21 @@ public class DedupeTotalDataService {
AccessScope scope = resolveAccessScope(operatorId);
LambdaQueryWrapper<DedupeTotalDataEntity> query = new LambdaQueryWrapper<DedupeTotalDataEntity>()
.like(!safeKeyword.isEmpty(), DedupeTotalDataEntity::getDataValue, safeKeyword)
.in(!scope.allUsers(), DedupeTotalDataEntity::getUploaderUserId, scope.userIds())
.like(!safeUsername.isEmpty(), DedupeTotalDataEntity::getUploaderUsername, safeUsername)
.ge(startDate != null, DedupeTotalDataEntity::getCreatedAt,
startDate == null ? null : startDate.atStartOfDay())
.lt(endDate != null, DedupeTotalDataEntity::getCreatedAt,
endDate == null ? null : endDate.plusDays(1).atStartOfDay())
.orderByDesc(DedupeTotalDataEntity::getId);
applyGroupScope(query, scope, groupId);
Long total = dedupeTotalDataMapper.selectCount(query);
List<DedupeTotalDataItemVo> items = dedupeTotalDataMapper.selectList(query.last("LIMIT " + ((safePage - 1) * safePageSize) + ", " + safePageSize))
.stream()
.map(this::toItemVo)
List<DedupeTotalDataEntity> rows = dedupeTotalDataMapper.selectList(
query.last("LIMIT " + ((safePage - 1) * safePageSize) + ", " + safePageSize));
Map<Long, String> groupNames = loadGroupNames(rows);
List<DedupeTotalDataItemVo> items = rows.stream()
.map(row -> toItemVo(row, row.getGroupId() == null
? ""
: groupNames.getOrDefault(row.getGroupId(), "")))
.toList();
DedupeTotalDataPageVo vo = new DedupeTotalDataPageVo();
vo.setItems(items);
@@ -107,24 +120,29 @@ public class DedupeTotalDataService {
}
public byte[] export(String username, LocalDate startDate, LocalDate endDate, Long operatorId) {
return export(username, startDate, endDate, null, operatorId);
}
public byte[] export(String username, LocalDate startDate, LocalDate endDate,
Long groupId, Long operatorId) {
if (startDate != null && endDate != null && startDate.isAfter(endDate)) {
throw new BusinessException("开始日期不能晚于结束日期");
}
String safeUsername = username == null ? "" : username.trim();
AccessScope scope = resolveAccessScope(operatorId);
LambdaQueryWrapper<DedupeTotalDataEntity> query = new LambdaQueryWrapper<DedupeTotalDataEntity>()
.in(!scope.allUsers(), DedupeTotalDataEntity::getUploaderUserId, scope.userIds())
.like(!safeUsername.isEmpty(), DedupeTotalDataEntity::getUploaderUsername, safeUsername)
.ge(startDate != null, DedupeTotalDataEntity::getCreatedAt,
startDate == null ? null : startDate.atStartOfDay())
.lt(endDate != null, DedupeTotalDataEntity::getCreatedAt,
endDate == null ? null : endDate.plusDays(1).atStartOfDay())
.orderByDesc(DedupeTotalDataEntity::getId);
applyGroupScope(query, scope, groupId);
List<DedupeTotalDataEntity> rows = dedupeTotalDataMapper.selectList(query);
return buildExportWorkbook(rows);
return buildExportWorkbook(rows, loadGroupNames(rows));
}
private byte[] buildExportWorkbook(List<DedupeTotalDataEntity> rows) {
private byte[] buildExportWorkbook(List<DedupeTotalDataEntity> rows, Map<Long, String> groupNames) {
try (SXSSFWorkbook workbook = new SXSSFWorkbook(100);
ByteArrayOutputStream outputStream = new ByteArrayOutputStream()) {
Sheet sheet = workbook.createSheet("DedupeTotalData");
@@ -132,19 +150,24 @@ public class DedupeTotalDataService {
header.createCell(0).setCellValue("ID");
header.createCell(1).setCellValue("ASIN值");
header.createCell(2).setCellValue("用户名");
header.createCell(3).setCellValue("创建时间");
header.createCell(3).setCellValue("分组");
header.createCell(4).setCellValue("创建时间");
for (int index = 0; index < rows.size(); index++) {
DedupeTotalDataEntity entity = rows.get(index);
Row row = sheet.createRow(index + 1);
row.createCell(0).setCellValue(entity.getId() == null ? "" : String.valueOf(entity.getId()));
row.createCell(1).setCellValue(entity.getDataValue() == null ? "" : entity.getDataValue());
row.createCell(2).setCellValue(entity.getUploaderUsername() == null ? "" : entity.getUploaderUsername());
row.createCell(3).setCellValue(formatExportTime(entity.getCreatedAt()));
row.createCell(3).setCellValue(entity.getGroupId() == null
? ""
: groupNames.getOrDefault(entity.getGroupId(), ""));
row.createCell(4).setCellValue(formatExportTime(entity.getCreatedAt()));
}
sheet.setColumnWidth(0, 3600);
sheet.setColumnWidth(1, 5200);
sheet.setColumnWidth(2, 5200);
sheet.setColumnWidth(3, 5600);
sheet.setColumnWidth(3, 5200);
sheet.setColumnWidth(4, 5600);
workbook.write(outputStream);
workbook.dispose();
return outputStream.toByteArray();
@@ -160,14 +183,16 @@ public class DedupeTotalDataService {
@Transactional
public DedupeTotalDataItemVo create(DedupeTotalDataCreateRequest request, Long operatorId) {
AdminUserEntity uploader = getOperator(operatorId);
ShopManageGroupEntity group = resolveWritableGroup(request.getGroupId(), uploader);
String dataValue = normalizeComparableValue(request.getDataValue());
ensureUnique(dataValue, null);
DedupeTotalDataEntity entity = new DedupeTotalDataEntity();
entity.setDataValue(dataValue);
entity.setGroupId(group.getId());
entity.setUploaderUserId(uploader.getId());
entity.setUploaderUsername(uploader.getUsername());
dedupeTotalDataMapper.insert(entity);
return toItemVo(getById(entity.getId()));
return toItemVo(getById(entity.getId()), group.getGroupName());
}
public Set<String> listComparableValues() {
@@ -204,11 +229,16 @@ public class DedupeTotalDataService {
}
public DedupeTotalDataImportStartVo startImport(MultipartFile file, Long operatorId) {
return startImport(file, null, operatorId);
}
public DedupeTotalDataImportStartVo startImport(MultipartFile file, Long groupId, Long operatorId) {
cleanupExpiredProgress();
if (file == null || file.isEmpty()) {
throw new BusinessException("请上传 xlsx 文件");
}
AdminUserEntity uploader = getOperator(operatorId);
ShopManageGroupEntity group = resolveWritableGroup(groupId, uploader);
String importId = IdUtil.fastSimpleUUID();
DedupeTotalDataImportProgressVo progress = new DedupeTotalDataImportProgressVo();
progress.setStatus("pending");
@@ -219,15 +249,17 @@ public class DedupeTotalDataService {
progress.setSkippedCount(0);
importProgressMap.put(importId, progress);
importOwnerMap.put(importId, uploader.getId());
importGroupMap.put(importId, group.getId());
try {
File tempFile = saveMultipartToTempFile(file);
String filename = file.getOriginalFilename();
Thread.ofVirtual().start(() -> runImportTask(
importId, tempFile, filename, uploader.getId(), uploader.getUsername()));
importId, tempFile, filename, uploader.getId(), uploader.getUsername(), group.getId()));
} catch (Exception e) {
importProgressMap.remove(importId);
importOwnerMap.remove(importId);
importGroupMap.remove(importId);
throw new BusinessException("读取上传文件失败");
}
@@ -240,19 +272,25 @@ public class DedupeTotalDataService {
cleanupExpiredProgress();
DedupeTotalDataImportProgressVo progress = importProgressMap.get(importId);
Long ownerId = importOwnerMap.get(importId);
if (progress == null || ownerId == null) {
Long groupId = importGroupMap.get(importId);
if (progress == null || ownerId == null || groupId == null) {
throw new BusinessException("导入任务不存在");
}
ensureUploaderAccess(ownerId, resolveAccessScope(operatorId));
ensureImportTaskAccess(groupId, resolveAccessScope(operatorId));
return progress;
}
public DedupeTotalDataImportStartVo startDeleteImport(MultipartFile file, Long operatorId) {
return startDeleteImport(file, null, operatorId);
}
public DedupeTotalDataImportStartVo startDeleteImport(MultipartFile file, Long groupId, Long operatorId) {
cleanupExpiredProgress();
if (file == null || file.isEmpty()) {
throw new BusinessException("请上传 xlsx 文件");
}
AdminUserEntity operator = getOperator(operatorId);
ShopManageGroupEntity group = resolveWritableGroup(groupId, operator);
String importId = IdUtil.fastSimpleUUID();
DedupeTotalDataImportProgressVo progress = new DedupeTotalDataImportProgressVo();
progress.setStatus("pending");
@@ -263,14 +301,17 @@ public class DedupeTotalDataService {
progress.setSkippedCount(0);
deleteImportProgressMap.put(importId, progress);
deleteImportOwnerMap.put(importId, operator.getId());
deleteImportGroupMap.put(importId, group.getId());
try {
File tempFile = saveMultipartToTempFile(file);
String filename = file.getOriginalFilename();
Thread.ofVirtual().start(() -> runDeleteImportTask(importId, tempFile, filename, operator.getId()));
Thread.ofVirtual().start(() -> runDeleteImportTask(
importId, tempFile, filename, operator.getId(), group.getId()));
} catch (Exception e) {
deleteImportProgressMap.remove(importId);
deleteImportOwnerMap.remove(importId);
deleteImportGroupMap.remove(importId);
throw new BusinessException("读取上传文件失败");
}
@@ -283,14 +324,16 @@ public class DedupeTotalDataService {
cleanupExpiredProgress();
DedupeTotalDataImportProgressVo progress = deleteImportProgressMap.get(importId);
Long ownerId = deleteImportOwnerMap.get(importId);
if (progress == null || ownerId == null) {
Long groupId = deleteImportGroupMap.get(importId);
if (progress == null || ownerId == null || groupId == null) {
throw new BusinessException("删除任务不存在");
}
ensureUploaderAccess(ownerId, resolveAccessScope(operatorId));
ensureImportTaskAccess(groupId, resolveAccessScope(operatorId));
return progress;
}
private void runDeleteImportTask(String importId, File tempFile, String filename, Long operatorId) {
private void runDeleteImportTask(String importId, File tempFile, String filename,
Long operatorId, Long groupId) {
DedupeTotalDataImportProgressVo progress = deleteImportProgressMap.get(importId);
if (progress == null) {
return;
@@ -298,7 +341,8 @@ public class DedupeTotalDataService {
progress.setStatus("running");
try (InputStream inputStream = Files.newInputStream(tempFile.toPath())) {
AccessScope scope = resolveAccessScope(operatorId);
DedupeTotalDataImportVo result = deleteFromExcelInternal(inputStream, filename, progress, scope);
DedupeTotalDataImportVo result = deleteFromExcelInternal(
inputStream, filename, progress, scope, groupId);
progress.setTotalRows(result.getTotalRows());
progress.setAsinCount(result.getAsinCount());
progress.setInsertedCount(result.getInsertedCount());
@@ -315,7 +359,7 @@ public class DedupeTotalDataService {
}
private void runImportTask(String importId, File tempFile, String filename,
Long uploaderUserId, String uploaderUsername) {
Long uploaderUserId, String uploaderUsername, Long groupId) {
DedupeTotalDataImportProgressVo progress = importProgressMap.get(importId);
if (progress == null) {
return;
@@ -323,7 +367,7 @@ public class DedupeTotalDataService {
progress.setStatus("running");
try (InputStream inputStream = Files.newInputStream(tempFile.toPath())) {
DedupeTotalDataImportVo result = importFromExcelInternal(
inputStream, filename, progress, uploaderUserId, uploaderUsername);
inputStream, filename, progress, uploaderUserId, uploaderUsername, groupId);
progress.setTotalRows(result.getTotalRows());
progress.setAsinCount(result.getAsinCount());
progress.setInsertedCount(result.getInsertedCount());
@@ -364,13 +408,19 @@ public class DedupeTotalDataService {
}
public DedupeTotalDataImportVo importFromExcel(MultipartFile file, Long operatorId) {
return importFromExcel(file, null, operatorId);
}
public DedupeTotalDataImportVo importFromExcel(MultipartFile file, Long groupId, Long operatorId) {
if (file == null || file.isEmpty()) {
throw new BusinessException("请上传 xlsx 文件");
}
AdminUserEntity uploader = getOperator(operatorId);
ShopManageGroupEntity group = resolveWritableGroup(groupId, uploader);
try (InputStream inputStream = file.getInputStream()) {
return importFromExcelInternal(
inputStream, file.getOriginalFilename(), null, uploader.getId(), uploader.getUsername());
inputStream, file.getOriginalFilename(), null, uploader.getId(), uploader.getUsername(),
group.getId());
} catch (BusinessException e) {
throw e;
} catch (Exception e) {
@@ -379,12 +429,19 @@ public class DedupeTotalDataService {
}
public DedupeTotalDataImportVo deleteFromExcel(MultipartFile file, Long operatorId) {
return deleteFromExcel(file, null, operatorId);
}
public DedupeTotalDataImportVo deleteFromExcel(MultipartFile file, Long groupId, Long operatorId) {
if (file == null || file.isEmpty()) {
throw new BusinessException("请上传 xlsx 文件");
}
AdminUserEntity operator = getOperator(operatorId);
ShopManageGroupEntity group = resolveWritableGroup(groupId, operator);
AccessScope scope = resolveAccessScope(operatorId);
try (InputStream inputStream = file.getInputStream()) {
return deleteFromExcelInternal(inputStream, file.getOriginalFilename(), null, scope);
return deleteFromExcelInternal(
inputStream, file.getOriginalFilename(), null, scope, group.getId());
} catch (BusinessException e) {
throw e;
} catch (Exception e) {
@@ -394,7 +451,8 @@ public class DedupeTotalDataService {
private DedupeTotalDataImportVo importFromExcelInternal(InputStream inputStream, String filename,
DedupeTotalDataImportProgressVo progress,
Long uploaderUserId, String uploaderUsername) {
Long uploaderUserId, String uploaderUsername,
Long groupId) {
String lowerFilename = filename == null ? "" : filename.toLowerCase();
if (!(lowerFilename.endsWith(".xlsx") || lowerFilename.endsWith(".xls"))) {
throw new BusinessException("仅支持 .xlsx 或 .xls 文件");
@@ -490,6 +548,7 @@ public class DedupeTotalDataService {
}
DedupeTotalDataEntity entity = new DedupeTotalDataEntity();
entity.setDataValue(pendingDataValue);
entity.setGroupId(groupId);
entity.setUploaderUserId(uploaderUserId);
entity.setUploaderUsername(uploaderUsername);
dedupeTotalDataMapper.insert(entity);
@@ -526,7 +585,7 @@ public class DedupeTotalDataService {
private DedupeTotalDataImportVo deleteFromExcelInternal(InputStream inputStream, String filename,
DedupeTotalDataImportProgressVo progress,
AccessScope scope) {
AccessScope scope, Long groupId) {
String lowerFilename = filename == null ? "" : filename.toLowerCase();
if (!(lowerFilename.endsWith(".xlsx") || lowerFilename.endsWith(".xls"))) {
throw new BusinessException("仅支持 .xlsx 或 .xls 文件");
@@ -604,7 +663,8 @@ public class DedupeTotalDataService {
continue;
}
int deletedThisRow = newRequiresNewTemplate().execute(status -> deleteByDataValue(dataValue, scope));
int deletedThisRow = newRequiresNewTemplate().execute(
status -> deleteByDataValue(dataValue, scope, groupId));
if (deletedThisRow > 0) {
deletedCount += deletedThisRow;
} else {
@@ -634,11 +694,14 @@ public class DedupeTotalDataService {
@Transactional
public DedupeTotalDataItemVo update(Long id, DedupeTotalDataUpdateRequest request, Long operatorId) {
DedupeTotalDataEntity entity = getAccessibleById(id, operatorId);
AdminUserEntity operator = getOperator(operatorId);
String dataValue = normalizeComparableValue(request.getDataValue());
ensureUnique(dataValue, id);
entity.setDataValue(dataValue);
ShopManageGroupEntity group = resolveWritableGroup(request.getGroupId(), operator);
entity.setGroupId(group.getId());
dedupeTotalDataMapper.updateById(entity);
return toItemVo(getById(id));
return toItemVo(getById(id), group.getGroupName());
}
@Transactional
@@ -649,7 +712,7 @@ public class DedupeTotalDataService {
private DedupeTotalDataEntity getAccessibleById(Long id, Long operatorId) {
DedupeTotalDataEntity entity = getById(id);
ensureUploaderAccess(entity.getUploaderUserId(), resolveAccessScope(operatorId));
ensureRecordAccess(entity, resolveAccessScope(operatorId));
return entity;
}
@@ -690,26 +753,78 @@ public class DedupeTotalDataService {
return dedupeTotalDataMapper.selectOne(query) != null;
}
private int deleteByDataValue(String dataValue, AccessScope scope) {
return dedupeTotalDataMapper.delete(new LambdaQueryWrapper<DedupeTotalDataEntity>()
.eq(DedupeTotalDataEntity::getDataValue, dataValue)
.in(!scope.allUsers(), DedupeTotalDataEntity::getUploaderUserId, scope.userIds()));
private int deleteByDataValue(String dataValue, AccessScope scope, Long groupId) {
LambdaQueryWrapper<DedupeTotalDataEntity> query = new LambdaQueryWrapper<DedupeTotalDataEntity>()
.eq(DedupeTotalDataEntity::getDataValue, dataValue);
applyGroupScope(query, scope, groupId);
return dedupeTotalDataMapper.delete(query);
}
private AccessScope resolveAccessScope(Long operatorId) {
AdminUserEntity operator = getOperator(operatorId);
if (isSuperAdmin(operator)) {
return new AccessScope(true, Set.of());
return new AccessScope(true, Set.of(), Set.of());
}
List<Long> accessibleGroupIds = shopManageGroupMapper.selectAccessibleGroupIds(operator.getId());
LinkedHashSet<Long> visibleUserIds = new LinkedHashSet<>();
visibleUserIds.add(operator.getId());
if (accessibleGroupIds != null && !accessibleGroupIds.isEmpty()) {
List<Long> groupUserIds = shopManageGroupMapper.selectUserIdsByGroupIds(accessibleGroupIds);
if (groupUserIds != null) {
groupUserIds.stream()
.filter(id -> id != null && id > 0)
.forEach(visibleUserIds::add);
}
}
List<Long> memberUserIds = shopManageGroupMapper.selectManagedMemberUserIds(operator.getId());
if (memberUserIds != null) {
memberUserIds.stream()
.filter(id -> id != null && id > 0)
.forEach(visibleUserIds::add);
}
return new AccessScope(false, Set.copyOf(visibleUserIds));
return new AccessScope(
false,
Set.copyOf(accessibleGroupIds == null ? List.of() : accessibleGroupIds),
Set.copyOf(visibleUserIds));
}
private void applyGroupScope(LambdaQueryWrapper<DedupeTotalDataEntity> query,
AccessScope scope, Long groupId) {
if (groupId != null && groupId > 0) {
resolveAccessibleGroup(groupId, scope);
query.eq(DedupeTotalDataEntity::getGroupId, groupId);
return;
}
if (scope.allUsers()) {
return;
}
query.and(wrapper -> {
if (!scope.groupIds().isEmpty()) {
wrapper.in(DedupeTotalDataEntity::getGroupId, scope.groupIds())
.or();
}
wrapper.isNull(DedupeTotalDataEntity::getGroupId)
.in(DedupeTotalDataEntity::getUploaderUserId, scope.userIds());
});
}
private ShopManageGroupEntity resolveWritableGroup(Long groupId, AdminUserEntity operator) {
if (groupId == null || groupId <= 0) {
throw new BusinessException("请选择分组");
}
AccessScope scope = resolveAccessScope(operator.getId());
return resolveAccessibleGroup(groupId, scope);
}
private ShopManageGroupEntity resolveAccessibleGroup(Long groupId, AccessScope scope) {
ShopManageGroupEntity group = shopManageGroupMapper.selectById(groupId);
if (group == null) {
throw new BusinessException("分组不存在");
}
if (!scope.allUsers() && !scope.groupIds().contains(groupId)) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "无权操作该分组数据");
}
return group;
}
private AdminUserEntity getOperator(Long operatorId) {
@@ -729,27 +844,46 @@ public class DedupeTotalDataService {
|| (role.isEmpty() && Integer.valueOf(1).equals(user.getIsAdmin()) && user.getCreatedById() == null);
}
private void ensureUploaderAccess(Long uploaderUserId, AccessScope scope) {
if (!scope.allUsers() && (uploaderUserId == null || !scope.userIds().contains(uploaderUserId))) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "无权操作该总数据");
private void ensureRecordAccess(DedupeTotalDataEntity entity, AccessScope scope) {
if (scope.allUsers()) {
return;
}
if (entity.getGroupId() != null && scope.groupIds().contains(entity.getGroupId())) {
return;
}
if (entity.getGroupId() == null && entity.getUploaderUserId() != null
&& scope.userIds().contains(entity.getUploaderUserId())) {
return;
}
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "无权操作该总数据");
}
private void ensureImportTaskAccess(Long groupId, AccessScope scope) {
if (!scope.allUsers() && !scope.groupIds().contains(groupId)) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "无权查看该导入任务");
}
}
private void cleanupExpiredProgress() {
long cutoff = System.currentTimeMillis() - COMPLETED_PROGRESS_RETENTION_MILLIS;
cleanupExpiredProgressEntries(importCompletedAtMap, importProgressMap, importOwnerMap, cutoff);
cleanupExpiredProgressEntries(deleteImportCompletedAtMap, deleteImportProgressMap, deleteImportOwnerMap, cutoff);
cleanupExpiredProgressEntries(
importCompletedAtMap, importProgressMap, importOwnerMap, importGroupMap, cutoff);
cleanupExpiredProgressEntries(
deleteImportCompletedAtMap, deleteImportProgressMap, deleteImportOwnerMap,
deleteImportGroupMap, cutoff);
}
private void cleanupExpiredProgressEntries(
Map<String, Long> completedAtMap,
Map<String, DedupeTotalDataImportProgressVo> progressMap,
Map<String, Long> ownerMap,
Map<String, Long> groupMap,
long cutoff) {
completedAtMap.forEach((id, completedAt) -> {
if (completedAt != null && completedAt < cutoff && completedAtMap.remove(id, completedAt)) {
progressMap.remove(id);
ownerMap.remove(id);
groupMap.remove(id);
}
});
}
@@ -768,16 +902,36 @@ public class DedupeTotalDataService {
.replaceAll("\\s+", " ");
}
private DedupeTotalDataItemVo toItemVo(DedupeTotalDataEntity entity) {
private Map<Long, String> loadGroupNames(List<DedupeTotalDataEntity> rows) {
List<Long> groupIds = rows.stream()
.map(DedupeTotalDataEntity::getGroupId)
.filter(id -> id != null && id > 0)
.distinct()
.toList();
if (groupIds.isEmpty()) {
return Map.of();
}
return shopManageGroupMapper.selectBatchIds(groupIds).stream()
.filter(group -> group.getId() != null)
.collect(java.util.stream.Collectors.toMap(
ShopManageGroupEntity::getId,
group -> group.getGroupName() == null ? "" : group.getGroupName(),
(left, right) -> left,
LinkedHashMap::new));
}
private DedupeTotalDataItemVo toItemVo(DedupeTotalDataEntity entity, String groupName) {
DedupeTotalDataItemVo vo = new DedupeTotalDataItemVo();
vo.setId(entity.getId());
vo.setDataValue(entity.getDataValue());
vo.setGroupId(entity.getGroupId());
vo.setGroupName(groupName == null ? "" : groupName);
vo.setUploaderUserId(entity.getUploaderUserId());
vo.setUsername(entity.getUploaderUsername());
vo.setCreatedAt(entity.getCreatedAt());
return vo;
}
private record AccessScope(boolean allUsers, Set<Long> userIds) {
private record AccessScope(boolean allUsers, Set<Long> groupIds, Set<Long> userIds) {
}
}
@@ -147,6 +147,18 @@ public class OssStorageService {
}
}
/**
* Reads a managed result object from either its object key or stored public URL.
* This overload keeps callers from depending on the configured bucket name.
*/
public byte[] readObjectBytes(String value) {
StorageLocation location = resolveStorageLocation(value);
if (location == null) {
throw new IllegalArgumentException("object value must not be blank");
}
return readObjectBytes(location.bucket(), location.objectKey());
}
public boolean objectExists(String bucket, String objectKey) {
String normalizedBucket = requireStorageName(bucket, "bucket");
String normalizedObjectKey = requireStorageName(objectKey, "objectKey");
@@ -1,19 +1,25 @@
package com.nanri.aiimage.modules.invalidasin.controller;
import com.nanri.aiimage.common.api.ApiResponse;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.modules.admin.support.AdminAuthSupport;
import com.nanri.aiimage.modules.invalidasin.model.dto.InvalidAsinDataCreateRequest;
import com.nanri.aiimage.modules.invalidasin.model.dto.InvalidAsinDataUpdateRequest;
import com.nanri.aiimage.modules.invalidasin.model.vo.InvalidAsinDataItemVo;
import com.nanri.aiimage.modules.invalidasin.model.vo.InvalidAsinDataPageVo;
import com.nanri.aiimage.modules.invalidasin.service.InvalidAsinDataService;
import com.nanri.aiimage.modules.permission.model.entity.AdminUserEntity;
import com.nanri.aiimage.modules.permission.service.PermissionMenuService;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.Parameter;
import io.swagger.v3.oas.annotations.media.Content;
import io.swagger.v3.oas.annotations.media.Schema;
import io.swagger.v3.oas.annotations.responses.ApiResponses;
import io.swagger.v3.oas.annotations.tags.Tag;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.web.bind.annotation.DeleteMapping;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
@@ -24,13 +30,29 @@ import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.security.MessageDigest;
@RestController
@RequiredArgsConstructor
@RequestMapping("/api/admin/invalid-asin-data")
@Tag(name = "不符合ASIN数据", description = "维护不符合ASIN数据列表,支持增删改查。")
public class InvalidAsinDataController {
private static final String INVALID_ASIN_DATA_COLUMN_KEY = "admin_invalid_asin_data";
private static final String INVALID_ASIN_DATA_ROUTE_PATH = "invalid-asin-data";
@Value("${aiimage.security.internal-token:}")
private String internalToken;
@Value("${aiimage.security.internal-token-file:}")
private String internalTokenFile;
private final InvalidAsinDataService invalidAsinDataService;
private final AdminAuthSupport adminAuthSupport;
private final PermissionMenuService permissionMenuService;
@GetMapping
@Operation(summary = "分页查询不符合ASIN数据", description = "分页查询不符合ASIN数据,支持按 ASIN/品牌 模糊搜索。")
@@ -40,8 +62,12 @@ public class InvalidAsinDataController {
public ApiResponse<InvalidAsinDataPageVo> page(
@Parameter(description = "页码") @RequestParam(defaultValue = "1") Long page,
@Parameter(description = "每页数量") @RequestParam(defaultValue = "15") Long pageSize,
@Parameter(description = "模糊搜索关键字") @RequestParam(required = false) String keyword) {
return ApiResponse.success(invalidAsinDataService.page(page, pageSize, keyword));
@Parameter(description = "模糊搜索关键字") @RequestParam(required = false) String keyword,
@Parameter(description = "分组 ID") @RequestParam(required = false) Long groupId,
HttpServletRequest request) {
RequestOperator operator = requireInvalidAsinDataAccess(request);
return ApiResponse.success(invalidAsinDataService.page(
page, pageSize, keyword, groupId, operator.id(), operator.superAdmin()));
}
@PostMapping
@@ -50,8 +76,12 @@ public class InvalidAsinDataController {
@io.swagger.v3.oas.annotations.responses.ApiResponse(responseCode = "200", description = "创建成功", content = @Content(schema = @Schema(implementation = InvalidAsinDataItemVo.class))),
@io.swagger.v3.oas.annotations.responses.ApiResponse(responseCode = "400", description = "参数不合法或数据重复")
})
public ApiResponse<InvalidAsinDataItemVo> create(@Valid @RequestBody InvalidAsinDataCreateRequest request) {
return ApiResponse.success("创建成功", invalidAsinDataService.create(request));
public ApiResponse<InvalidAsinDataItemVo> create(
HttpServletRequest httpRequest,
@Valid @RequestBody InvalidAsinDataCreateRequest request) {
RequestOperator operator = requireInvalidAsinDataAccess(httpRequest);
return ApiResponse.success("创建成功", invalidAsinDataService.create(
request, operator.id(), operator.superAdmin()));
}
@PutMapping("/{id}")
@@ -63,8 +93,11 @@ public class InvalidAsinDataController {
})
public ApiResponse<InvalidAsinDataItemVo> update(
@Parameter(description = "主键ID", required = true) @PathVariable Long id,
HttpServletRequest httpRequest,
@Valid @RequestBody InvalidAsinDataUpdateRequest request) {
return ApiResponse.success("更新成功", invalidAsinDataService.update(id, request));
RequestOperator operator = requireInvalidAsinDataAccess(httpRequest);
return ApiResponse.success("更新成功", invalidAsinDataService.update(
id, request, operator.id(), operator.superAdmin()));
}
@DeleteMapping("/{id}")
@@ -73,8 +106,107 @@ public class InvalidAsinDataController {
@io.swagger.v3.oas.annotations.responses.ApiResponse(responseCode = "200", description = "删除成功"),
@io.swagger.v3.oas.annotations.responses.ApiResponse(responseCode = "404", description = "数据不存在")
})
public ApiResponse<Void> delete(@Parameter(description = "主键ID", required = true) @PathVariable Long id) {
invalidAsinDataService.delete(id);
public ApiResponse<Void> delete(
@Parameter(description = "主键ID", required = true) @PathVariable Long id,
HttpServletRequest request) {
RequestOperator operator = requireInvalidAsinDataAccess(request);
invalidAsinDataService.delete(id, operator.id(), operator.superAdmin());
return ApiResponse.success("删除成功", null);
}
private RequestOperator requireInvalidAsinDataAccess(HttpServletRequest request) {
AdminUserEntity operator = resolveOperator(request);
if (operator == null || operator.getId() == null || operator.getId() <= 0) {
throw new BusinessException(403, "无权访问品牌数据库");
}
boolean superAdmin = "super_admin".equals(adminAuthSupport.currentRole(operator));
if (!superAdmin && !hasInvalidAsinDataPermission(operator)) {
throw new BusinessException(403, "无权访问品牌数据库");
}
return new RequestOperator(operator.getId(), superAdmin);
}
private AdminUserEntity resolveOperator(HttpServletRequest request) {
try {
return adminAuthSupport.requireUser(request);
} catch (BusinessException authFailure) {
AdminUserEntity internalOperator = resolveInternalOperator(request);
if (internalOperator != null) {
return internalOperator;
}
throw authFailure;
}
}
private boolean hasInvalidAsinDataPermission(AdminUserEntity operator) {
return permissionMenuService.getUserColumnPermissions(operator.getId(), "admin")
.stream()
.anyMatch(item -> INVALID_ASIN_DATA_COLUMN_KEY.equals(item.getColumnKey())
|| INVALID_ASIN_DATA_ROUTE_PATH.equals(item.getRoutePath()));
}
private AdminUserEntity resolveInternalOperator(HttpServletRequest request) {
if (!isTrustedInternalRequest(request.getHeader("X-Internal-Token"))) {
return null;
}
String rawOperatorId = request.getParameter("operatorId");
if (rawOperatorId == null || rawOperatorId.isBlank()) {
rawOperatorId = request.getParameter("operator_id");
}
if (rawOperatorId == null || rawOperatorId.isBlank()) {
return null;
}
try {
return permissionMenuService.requireUserOperator(Long.parseLong(rawOperatorId.trim()));
} catch (NumberFormatException | BusinessException ex) {
return null;
}
}
private boolean isTrustedInternalRequest(String suppliedToken) {
String expectedToken = resolveExpectedInternalToken();
return !expectedToken.isBlank() && suppliedToken != null && !suppliedToken.isBlank()
&& MessageDigest.isEqual(
expectedToken.getBytes(StandardCharsets.UTF_8),
suppliedToken.trim().getBytes(StandardCharsets.UTF_8));
}
private String resolveExpectedInternalToken() {
if (internalToken != null && !internalToken.isBlank()) {
return internalToken.trim();
}
Path path = resolveInternalTokenFile();
if (path == null || !Files.isRegularFile(path)) {
return "";
}
try {
return Files.readString(path, StandardCharsets.UTF_8).trim();
} catch (Exception ignored) {
return "";
}
}
private Path resolveInternalTokenFile() {
String configuredPath = internalTokenFile == null ? "" : internalTokenFile.trim();
if (!configuredPath.isEmpty()) {
if (configuredPath.equals("~") || configuredPath.startsWith("~/") || configuredPath.startsWith("~\\")) {
String userHome = System.getProperty("user.home", "").trim();
if (userHome.isEmpty()) {
return null;
}
configuredPath = configuredPath.length() == 1
? userHome
: Path.of(userHome, configuredPath.substring(2)).toString();
}
Path configuredTokenPath = Path.of(configuredPath);
return configuredTokenPath.isAbsolute() ? configuredTokenPath.normalize() : null;
}
String userHome = System.getProperty("user.home", "").trim();
return userHome.isEmpty()
? null
: Path.of(userHome, ".aiimage", "internal-token").toAbsolutePath().normalize();
}
private record RequestOperator(Long id, boolean superAdmin) {
}
}
@@ -2,6 +2,7 @@ package com.nanri.aiimage.modules.invalidasin.model.dto;
import io.swagger.v3.oas.annotations.media.Schema;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import jakarta.validation.constraints.Size;
import lombok.Data;
@@ -17,4 +18,8 @@ public class InvalidAsinDataCreateRequest {
@Size(max = 128, message = "品牌长度不能超过128个字符")
@Schema(description = "品牌名称")
private String brand;
@NotNull(message = "请选择分组")
@Schema(description = "分组 ID", requiredMode = Schema.RequiredMode.REQUIRED)
private Long groupId;
}
@@ -17,4 +17,7 @@ public class InvalidAsinDataUpdateRequest {
@Size(max = 128, message = "品牌长度不能超过128个字符")
@Schema(description = "品牌名称")
private String brand;
@Schema(description = "分组 ID;自动导入数据无需填写")
private Long groupId;
}
@@ -15,6 +15,8 @@ public class InvalidAsinDataEntity {
private Long id;
private String dataValue;
private String brand;
private Long groupId;
private String recordSource;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}
@@ -18,6 +18,15 @@ public class InvalidAsinDataItemVo {
@Schema(description = "品牌名称")
private String brand;
@Schema(description = "分组 ID")
private Long groupId;
@Schema(description = "分组名称")
private String groupName;
@Schema(description = "记录来源:MANUAL/AUTO")
private String recordSource;
@Schema(description = "创建时间")
private LocalDateTime createdAt;
}
@@ -8,20 +8,30 @@ import com.nanri.aiimage.modules.invalidasin.model.dto.InvalidAsinDataUpdateRequ
import com.nanri.aiimage.modules.invalidasin.model.entity.InvalidAsinDataEntity;
import com.nanri.aiimage.modules.invalidasin.model.vo.InvalidAsinDataItemVo;
import com.nanri.aiimage.modules.invalidasin.model.vo.InvalidAsinDataPageVo;
import com.nanri.aiimage.modules.shopkey.model.entity.ShopManageGroupEntity;
import com.nanri.aiimage.modules.shopkey.service.ShopManageGroupService;
import lombok.RequiredArgsConstructor;
import org.springframework.http.HttpStatus;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.web.server.ResponseStatusException;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Set;
@Service
@RequiredArgsConstructor
public class InvalidAsinDataService {
private final InvalidAsinDataMapper invalidAsinDataMapper;
public static final String RECORD_SOURCE_MANUAL = "MANUAL";
public static final String RECORD_SOURCE_AUTO = "AUTO";
public InvalidAsinDataPageVo page(long page, long pageSize, String keyword) {
private final InvalidAsinDataMapper invalidAsinDataMapper;
private final ShopManageGroupService shopManageGroupService;
public InvalidAsinDataPageVo page(long page, long pageSize, String keyword, Long groupId, Long operatorId, boolean superAdmin) {
long safePage = Math.max(page, 1);
long safePageSize = Math.min(Math.max(pageSize, 1), 100);
String safeKeyword = keyword == null ? "" : keyword.trim();
@@ -31,11 +41,28 @@ public class InvalidAsinDataService {
.or()
.like(InvalidAsinDataEntity::getBrand, safeKeyword))
.orderByDesc(InvalidAsinDataEntity::getId);
if (!superAdmin) {
Long fixedGroupId = resolveFixedAccessibleGroupId(operatorId);
if (fixedGroupId == null) {
query.eq(InvalidAsinDataEntity::getId, -1L);
} else {
query.eq(InvalidAsinDataEntity::getRecordSource, RECORD_SOURCE_MANUAL)
.eq(InvalidAsinDataEntity::getGroupId, fixedGroupId);
}
} else if (groupId != null && groupId > 0) {
query.eq(InvalidAsinDataEntity::getGroupId, groupId);
}
Long total = invalidAsinDataMapper.selectCount(query);
List<InvalidAsinDataItemVo> items = invalidAsinDataMapper
.selectList(query.last("LIMIT " + ((safePage - 1) * safePageSize) + ", " + safePageSize))
.stream()
.map(this::toItemVo)
List<InvalidAsinDataEntity> rows = invalidAsinDataMapper
.selectList(query.last("LIMIT " + ((safePage - 1) * safePageSize) + ", " + safePageSize));
Map<Long, String> groupNameById = shopManageGroupService.buildGroupNameMap(rows.stream()
.map(InvalidAsinDataEntity::getGroupId)
.filter(rowGroupId -> rowGroupId != null && rowGroupId > 0)
.distinct()
.toList());
List<InvalidAsinDataItemVo> items = rows.stream()
.map(entity -> toItemVo(entity,
entity.getGroupId() == null ? "" : groupNameById.getOrDefault(entity.getGroupId(), "")))
.toList();
InvalidAsinDataPageVo vo = new InvalidAsinDataPageVo();
vo.setItems(items);
@@ -46,35 +73,63 @@ public class InvalidAsinDataService {
}
@Transactional
public InvalidAsinDataItemVo create(InvalidAsinDataCreateRequest request) {
public InvalidAsinDataItemVo create(InvalidAsinDataCreateRequest request, Long operatorId, boolean superAdmin) {
ShopManageGroupEntity group = resolveManualWriteGroup(request.getGroupId(), operatorId, superAdmin);
String dataValue = normalizeRequired(request.getDataValue(), "ASIN 不能为空");
String brand = normalizeRequired(request.getBrand(), "品牌不能为空").toLowerCase(Locale.ROOT);
ensureUnique(dataValue, null);
InvalidAsinDataEntity entity = new InvalidAsinDataEntity();
entity.setDataValue(dataValue);
entity.setBrand(brand);
entity.setGroupId(group.getId());
entity.setRecordSource(RECORD_SOURCE_MANUAL);
invalidAsinDataMapper.insert(entity);
return toItemVo(getById(entity.getId()));
return toItemVo(getById(entity.getId()), group.getGroupName());
}
@Transactional
public InvalidAsinDataItemVo update(Long id, InvalidAsinDataUpdateRequest request) {
public InvalidAsinDataItemVo update(Long id, InvalidAsinDataUpdateRequest request, Long operatorId, boolean superAdmin) {
InvalidAsinDataEntity entity = getById(id);
ShopManageGroupEntity fixedGroup = ensureRecordAccess(entity, operatorId, superAdmin);
ShopManageGroupEntity writeGroup = null;
if (isManualRecord(entity)) {
Long requestedGroupId = requireGroupId(request.getGroupId());
writeGroup = superAdmin
? shopManageGroupService.getAccessibleById(requestedGroupId, operatorId, true)
: resolveLockedGroup(requestedGroupId, fixedGroup);
}
String dataValue = normalizeRequired(request.getDataValue(), "ASIN 不能为空");
String brand = normalizeRequired(request.getBrand(), "品牌不能为空").toLowerCase(Locale.ROOT);
ensureUnique(dataValue, id);
entity.setDataValue(dataValue);
entity.setBrand(brand);
String groupName = "";
if (writeGroup != null) {
entity.setGroupId(writeGroup.getId());
groupName = writeGroup.getGroupName();
}
invalidAsinDataMapper.updateById(entity);
return toItemVo(getById(id));
return toItemVo(getById(id), groupName);
}
@Transactional
public void delete(Long id) {
public void delete(Long id, Long operatorId, boolean superAdmin) {
InvalidAsinDataEntity entity = getById(id);
ensureRecordAccess(entity, operatorId, superAdmin);
invalidAsinDataMapper.deleteById(entity.getId());
}
private ShopManageGroupEntity resolveManualWriteGroup(Long requestedGroupId, Long operatorId, boolean superAdmin) {
Long normalizedRequestedGroupId = requireGroupId(requestedGroupId);
if (superAdmin) {
return shopManageGroupService.getAccessibleById(
normalizedRequestedGroupId, operatorId, true);
}
Long fixedGroupId = requireFixedAccessibleGroupId(operatorId);
ensureRequestedLockedGroup(normalizedRequestedGroupId, fixedGroupId);
return shopManageGroupService.getAccessibleById(fixedGroupId, operatorId, false);
}
private InvalidAsinDataEntity getById(Long id) {
InvalidAsinDataEntity entity = invalidAsinDataMapper.selectById(id);
if (entity == null) {
@@ -115,11 +170,76 @@ public class InvalidAsinDataService {
.replaceAll("\\s+", " ");
}
private InvalidAsinDataItemVo toItemVo(InvalidAsinDataEntity entity) {
private ShopManageGroupEntity ensureRecordAccess(InvalidAsinDataEntity entity, Long operatorId, boolean superAdmin) {
if (superAdmin) {
return null;
}
if (!isManualRecord(entity) || entity.getGroupId() == null || entity.getGroupId() <= 0) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "无权操作该品牌数据");
}
Long fixedGroupId = resolveFixedAccessibleGroupId(operatorId);
if (fixedGroupId == null || !fixedGroupId.equals(entity.getGroupId())) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "无权操作非固定分组数据");
}
try {
return shopManageGroupService.getAccessibleById(fixedGroupId, operatorId, false);
} catch (BusinessException exception) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "无权操作该品牌数据");
}
}
private Long requireFixedAccessibleGroupId(Long operatorId) {
Long fixedGroupId = resolveFixedAccessibleGroupId(operatorId);
if (fixedGroupId == null) {
throw new BusinessException("当前用户没有可访问的分组");
}
return fixedGroupId;
}
private Long resolveFixedAccessibleGroupId(Long operatorId) {
Set<Long> accessibleGroupIds = shopManageGroupService.listAccessibleGroupIds(operatorId, false);
if (accessibleGroupIds == null) {
return null;
}
return accessibleGroupIds.stream()
.filter(groupId -> groupId != null && groupId > 0)
.min(Long::compareTo)
.orElse(null);
}
private ShopManageGroupEntity resolveLockedGroup(Long requestedGroupId, ShopManageGroupEntity fixedGroup) {
if (fixedGroup == null || fixedGroup.getId() == null) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "无权操作非固定分组数据");
}
ensureRequestedLockedGroup(requestedGroupId, fixedGroup.getId());
return fixedGroup;
}
private void ensureRequestedLockedGroup(Long requestedGroupId, Long fixedGroupId) {
if (fixedGroupId == null || !fixedGroupId.equals(requestedGroupId)) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "无权操作非固定分组数据");
}
}
private Long requireGroupId(Long groupId) {
if (groupId == null || groupId <= 0) {
throw new BusinessException("请选择分组");
}
return groupId;
}
private boolean isManualRecord(InvalidAsinDataEntity entity) {
return entity != null && RECORD_SOURCE_MANUAL.equalsIgnoreCase(normalizeText(entity.getRecordSource()));
}
private InvalidAsinDataItemVo toItemVo(InvalidAsinDataEntity entity, String groupName) {
InvalidAsinDataItemVo vo = new InvalidAsinDataItemVo();
vo.setId(entity.getId());
vo.setDataValue(entity.getDataValue());
vo.setBrand(entity.getBrand());
vo.setGroupId(entity.getGroupId());
vo.setGroupName(groupName == null ? "" : groupName);
vo.setRecordSource(isManualRecord(entity) ? RECORD_SOURCE_MANUAL : RECORD_SOURCE_AUTO);
vo.setCreatedAt(entity.getCreatedAt());
return vo;
}
@@ -0,0 +1,9 @@
package com.nanri.aiimage.modules.shopdatacrawl.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyFileEntity;
import org.apache.ibatis.annotations.Mapper;
@Mapper
public interface ShopDataCrawlDailyFileMapper extends BaseMapper<ShopDataCrawlDailyFileEntity> {
}
@@ -0,0 +1,9 @@
package com.nanri.aiimage.modules.shopdatacrawl.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyMemberEntity;
import org.apache.ibatis.annotations.Mapper;
@Mapper
public interface ShopDataCrawlDailyMemberMapper extends BaseMapper<ShopDataCrawlDailyMemberEntity> {
}
@@ -0,0 +1,32 @@
package com.nanri.aiimage.modules.shopdatacrawl.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.LocalDate;
import java.time.LocalDateTime;
@Data
@TableName("biz_shop_data_crawl_daily_file")
public class ShopDataCrawlDailyFileEntity {
@TableId(type = IdType.AUTO)
private Long id;
private Long userId;
private String shopKeyHash;
private String shopKey;
private LocalDate businessDate;
private Long latestTaskId;
private Long latestResultId;
private String resultFilename;
private String resultFileUrl;
private Long resultFileSize;
private String resultContentType;
private Integer rowCount;
private Long version;
private LocalDateTime lastSuccessAt;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}
@@ -0,0 +1,20 @@
package com.nanri.aiimage.modules.shopdatacrawl.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_shop_data_crawl_daily_member")
public class ShopDataCrawlDailyMemberEntity {
@TableId(type = IdType.AUTO)
private Long id;
private Long dailyFileId;
private Long taskId;
private Long resultId;
private LocalDateTime createdAt;
}
@@ -0,0 +1,218 @@
package com.nanri.aiimage.modules.shopdatacrawl.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.nanri.aiimage.modules.shopdatacrawl.mapper.ShopDataCrawlDailyFileMapper;
import com.nanri.aiimage.modules.shopdatacrawl.mapper.ShopDataCrawlDailyMemberMapper;
import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyFileEntity;
import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyMemberEntity;
import com.nanri.aiimage.modules.task.model.entity.FileResultEntity;
import com.nanri.aiimage.modules.task.service.TaskDistributedLockService;
import lombok.RequiredArgsConstructor;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.stereotype.Service;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.time.Duration;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.util.HexFormat;
import java.util.List;
import java.util.Set;
/**
* Persistence and locking primitives for the shop-data-crawl daily workbook.
* Workbook assembly remains in ShopDataCrawlTaskService so task state changes
* and the daily pointer are committed together.
*/
@Service
@RequiredArgsConstructor
public class ShopDataCrawlDailyFileService {
public static final ZoneId BUSINESS_ZONE = ZoneId.of("Asia/Shanghai");
private static final Duration DAILY_LOCK_TTL = Duration.ofMinutes(30);
private static final long DAILY_LOCK_WAIT_MILLIS = 15000L;
private static final String DAILY_LOCK_MODULE_PREFIX = "SHOP_DATA_CRAWL_DAILY_";
private final ShopDataCrawlDailyFileMapper dailyFileMapper;
private final ShopDataCrawlDailyMemberMapper dailyMemberMapper;
private final TaskDistributedLockService taskDistributedLockService;
public LocalDate currentBusinessDate() {
return LocalDate.now(BUSINESS_ZONE);
}
public LocalDateTime currentBusinessDateTime() {
return LocalDateTime.now(BUSINESS_ZONE);
}
public String shopKey(FileResultEntity row) {
if (row == null) {
return null;
}
String shopId = trimToNull(row.getSourceFileUrl());
if (shopId != null) {
return "shop-id:" + shopId;
}
String shopName = trimToNull(row.getSourceFilename());
return shopName == null ? null : "shop-name:" + shopName;
}
public String shopKeyHash(String shopKey) {
if (shopKey == null || shopKey.isBlank()) {
return null;
}
byte[] digest;
try {
digest = MessageDigest.getInstance("SHA-256")
.digest(shopKey.getBytes(StandardCharsets.UTF_8));
} catch (Exception ex) {
throw new IllegalStateException("SHA-256 is unavailable", ex);
}
return HexFormat.of().formatHex(digest);
}
public TaskDistributedLockService.LockHandle acquireLock(Long userId, String shopKey) {
if (userId == null || userId <= 0 || shopKey == null || shopKey.isBlank()) {
return null;
}
// The lock protects the shop's whole daily-file lifecycle. Including the
// date would allow yesterday and today to update the same shop together.
String identity = userId + "|" + shopKey;
String lockModule = DAILY_LOCK_MODULE_PREFIX + shopKeyHash(identity);
return taskDistributedLockService.acquire(lockModule, 1L, DAILY_LOCK_TTL, DAILY_LOCK_WAIT_MILLIS);
}
public ShopDataCrawlDailyFileEntity findForUpdate(Long userId, String shopKeyHash, LocalDate businessDate) {
if (userId == null || shopKeyHash == null || businessDate == null) {
return null;
}
return dailyFileMapper.selectOne(new LambdaQueryWrapper<ShopDataCrawlDailyFileEntity>()
.eq(ShopDataCrawlDailyFileEntity::getUserId, userId)
.eq(ShopDataCrawlDailyFileEntity::getShopKeyHash, shopKeyHash)
.eq(ShopDataCrawlDailyFileEntity::getBusinessDate, businessDate)
.last("FOR UPDATE"));
}
public List<ShopDataCrawlDailyFileEntity> findOlder(Long userId, String shopKeyHash, LocalDate businessDate) {
if (userId == null || shopKeyHash == null || businessDate == null) {
return List.of();
}
return dailyFileMapper.selectList(new LambdaQueryWrapper<ShopDataCrawlDailyFileEntity>()
.eq(ShopDataCrawlDailyFileEntity::getUserId, userId)
.eq(ShopDataCrawlDailyFileEntity::getShopKeyHash, shopKeyHash)
.lt(ShopDataCrawlDailyFileEntity::getBusinessDate, businessDate)
.orderByDesc(ShopDataCrawlDailyFileEntity::getBusinessDate)
.orderByDesc(ShopDataCrawlDailyFileEntity::getId)
.last("FOR UPDATE"));
}
public List<ShopDataCrawlDailyFileEntity> findByLatestResultId(Long resultId) {
if (resultId == null || resultId <= 0) {
return List.of();
}
return dailyFileMapper.selectList(new LambdaQueryWrapper<ShopDataCrawlDailyFileEntity>()
.eq(ShopDataCrawlDailyFileEntity::getLatestResultId, resultId)
.last("FOR UPDATE"));
}
public ShopDataCrawlDailyFileEntity findById(Long dailyFileId) {
return dailyFileId == null || dailyFileId <= 0 ? null : dailyFileMapper.selectById(dailyFileId);
}
public List<ShopDataCrawlDailyMemberEntity> findMembersByResultId(Long resultId) {
if (resultId == null || resultId <= 0) {
return List.of();
}
return dailyMemberMapper.selectList(new LambdaQueryWrapper<ShopDataCrawlDailyMemberEntity>()
.eq(ShopDataCrawlDailyMemberEntity::getResultId, resultId));
}
public boolean containsResult(Long dailyFileId, Long resultId) {
if (dailyFileId == null || dailyFileId <= 0 || resultId == null || resultId <= 0) {
return false;
}
return dailyMemberMapper.selectCount(new LambdaQueryWrapper<ShopDataCrawlDailyMemberEntity>()
.eq(ShopDataCrawlDailyMemberEntity::getDailyFileId, dailyFileId)
.eq(ShopDataCrawlDailyMemberEntity::getResultId, resultId)) > 0;
}
public boolean addMember(Long dailyFileId, Long taskId, Long resultId) {
if (dailyFileId == null || dailyFileId <= 0 || taskId == null || taskId <= 0
|| resultId == null || resultId <= 0) {
return false;
}
ShopDataCrawlDailyMemberEntity member = new ShopDataCrawlDailyMemberEntity();
member.setDailyFileId(dailyFileId);
member.setTaskId(taskId);
member.setResultId(resultId);
member.setCreatedAt(currentBusinessDateTime());
try {
dailyMemberMapper.insert(member);
return true;
} catch (DuplicateKeyException ignored) {
return false;
}
}
public List<ShopDataCrawlDailyMemberEntity> listMembers(Long dailyFileId) {
if (dailyFileId == null || dailyFileId <= 0) {
return List.of();
}
return dailyMemberMapper.selectList(new LambdaQueryWrapper<ShopDataCrawlDailyMemberEntity>()
.eq(ShopDataCrawlDailyMemberEntity::getDailyFileId, dailyFileId)
.orderByDesc(ShopDataCrawlDailyMemberEntity::getCreatedAt)
.orderByDesc(ShopDataCrawlDailyMemberEntity::getId)
.last("FOR UPDATE"));
}
public void deleteMembers(Long dailyFileId) {
if (dailyFileId == null || dailyFileId <= 0) {
return;
}
dailyMemberMapper.delete(new LambdaQueryWrapper<ShopDataCrawlDailyMemberEntity>()
.eq(ShopDataCrawlDailyMemberEntity::getDailyFileId, dailyFileId));
}
public void deleteMembersForResults(Set<Long> resultIds) {
if (resultIds == null || resultIds.isEmpty()) {
return;
}
dailyMemberMapper.delete(new LambdaQueryWrapper<ShopDataCrawlDailyMemberEntity>()
.in(ShopDataCrawlDailyMemberEntity::getResultId, resultIds));
}
public long countObjectReferences(String objectKey) {
if (objectKey == null || objectKey.isBlank()) {
return 0L;
}
Long count = dailyFileMapper.selectCount(new LambdaQueryWrapper<ShopDataCrawlDailyFileEntity>()
.eq(ShopDataCrawlDailyFileEntity::getResultFileUrl, objectKey));
return count == null ? 0L : count;
}
public void deleteDailyFile(Long dailyFileId) {
if (dailyFileId == null || dailyFileId <= 0) {
return;
}
deleteMembers(dailyFileId);
dailyFileMapper.deleteById(dailyFileId);
}
public void insert(ShopDataCrawlDailyFileEntity entity) {
dailyFileMapper.insert(entity);
}
public void update(ShopDataCrawlDailyFileEntity entity) {
dailyFileMapper.updateById(entity);
}
private String trimToNull(String value) {
if (value == null) {
return null;
}
String normalized = value.trim();
return normalized.isEmpty() ? null : normalized;
}
}
@@ -19,6 +19,7 @@ import org.springframework.core.io.ClassPathResource;
import org.springframework.stereotype.Service;
import java.io.File;
import java.io.FileInputStream;
import java.io.FileOutputStream;
import java.io.InputStream;
import java.util.ArrayList;
@@ -64,6 +65,36 @@ public class ShopDataCrawlExcelAssemblyService {
}
}
/**
* Appends the supplied task rows to an already assembled daily workbook.
* Existing rows, styles and drawings are deliberately left untouched.
*/
public void appendWorkbook(File baseXlsx, File outputXlsx, List<ShopDataCrawlResultItemVo> items) {
if (baseXlsx == null || !baseXlsx.isFile()) {
throw new BusinessException("当天累计文件不存在");
}
try (InputStream input = new FileInputStream(baseXlsx);
XSSFWorkbook workbook = new XSSFWorkbook(input);
FileOutputStream output = new FileOutputStream(outputXlsx)) {
validateTemplate(workbook);
Map<String, List<ShopDataCrawlRowDto>> rowsByCountry = rowsByCountry(items);
Map<String, SimilarAsinImageEmbedder.ResizedImage> imageCache = new ConcurrentHashMap<>();
imageEmbedder.prefetch(imageUrls(rowsByCountry), imageCache);
Map<String, Integer> pictureIndexes = new LinkedHashMap<>();
for (int i = 0; i < COUNTRIES.size(); i++) {
List<ShopDataCrawlRowDto> rows = rowsByCountry.get(COUNTRIES.get(i));
if (rows != null && !rows.isEmpty()) {
appendSheet(workbook, workbook.getSheetAt(i), rows, imageCache, pictureIndexes);
}
}
workbook.write(output);
} catch (BusinessException ex) {
throw ex;
} catch (Exception ex) {
throw new BusinessException("店铺数据抓取累计 Excel 追加失败: " + ex.getMessage());
}
}
public int countRows(List<ShopDataCrawlResultItemVo> items) {
return rowsByCountry(items).values().stream().mapToInt(List::size).sum();
}
@@ -141,6 +172,50 @@ public class ShopDataCrawlExcelAssemblyService {
}
}
private void appendSheet(XSSFWorkbook workbook,
Sheet sheet,
List<ShopDataCrawlRowDto> rows,
Map<String, SimilarAsinImageEmbedder.ResizedImage> imageCache,
Map<String, Integer> pictureIndexes) {
Row header = sheet.getRow(0);
Row styleRow = sheet.getRow(1);
boolean currentTemplate = header != null && "商品图片".equals(cellText(header, IMAGE_COLUMN));
boolean templateHasBrand = header != null && "品牌".equals(cellText(header, BRAND_COLUMN));
CellStyle[] styles = new CellStyle[HEADERS.size()];
for (int column = 0; column < styles.length; column++) {
int sourceColumn = templateColumnForOutput(column, currentTemplate, templateHasBrand);
Cell cell = styleRow == null ? null : styleRow.getCell(sourceColumn);
styles[column] = cell == null ? null : cell.getCellStyle();
}
writeHeaders(sheet, currentTemplate, templateHasBrand);
sheet.setColumnWidth(IMAGE_COLUMN, IMAGE_COLUMN_WIDTH);
int rowIndex = Math.max(1, sheet.getLastRowNum() + 1);
for (ShopDataCrawlRowDto value : rows) {
Row row = sheet.createRow(rowIndex++);
writeRow(workbook, sheet, row, value, styles, imageCache, pictureIndexes);
}
}
private void writeRow(XSSFWorkbook workbook,
Sheet sheet,
Row row,
ShopDataCrawlRowDto value,
CellStyle[] styles,
Map<String, SimilarAsinImageEmbedder.ResizedImage> imageCache,
Map<String, Integer> pictureIndexes) {
String[] values = {value.getDate(), value.getAsin(), "", value.getInventorySales(), value.getSalesRank(),
value.getPageViews(), value.getUnitsSold(), value.getPrice(), value.getRecommendedOffer(), value.getBrand()};
for (int column = 0; column < values.length; column++) {
Cell cell = row.createCell(column);
if (styles[column] != null) cell.setCellStyle(styles[column]);
cell.setCellValue(values[column] == null ? "" : values[column]);
}
if (!blank(value.getCommodityImage())) {
row.setHeightInPoints(IMAGE_ROW_HEIGHT_POINTS);
embedImage(workbook, sheet, row, value.getCommodityImage(), imageCache, pictureIndexes);
}
}
private void writeHeaders(Sheet sheet, boolean currentTemplate, boolean templateHasBrand) {
Row header = sheet.getRow(0);
if (header == null) header = sheet.createRow(0);
@@ -17,6 +17,8 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlCreateTask
import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlShopPayloadDto;
import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlSubmitResultRequest;
import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlTaskItemDto;
import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyFileEntity;
import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyMemberEntity;
import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlTaskBatchVo;
import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlCreateTaskVo;
import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlHistoryVo;
@@ -45,16 +47,24 @@ import org.springframework.transaction.annotation.Transactional;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.transaction.support.TransactionSynchronization;
import org.springframework.transaction.support.TransactionSynchronizationManager;
import java.io.File;
import java.nio.file.Files;
import java.time.LocalDateTime;
import java.time.Duration;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeMap;
import java.util.UUID;
@Service
@RequiredArgsConstructor
@@ -66,7 +76,6 @@ public class ShopDataCrawlTaskService {
private static final int RESULT_PENDING = -1;
private static final int RESULT_FAILED = 0;
private static final int RESULT_SUCCESS = 1;
private static final int SHOP_HISTORY_RETENTION_LIMIT = 3;
private static final String INTERRUPTED_MESSAGE = "Python 在该店铺结果提交完成前中断";
private static final String PARTIAL_RESULT_MESSAGE = "Python 中断,已保留已回传的部分数据";
private static final String RESULT_CHUNK_SCOPE_PREFIX = "result-chunks:";
@@ -88,6 +97,7 @@ public class ShopDataCrawlTaskService {
private final TaskScopeStateMapper taskScopeStateMapper;
private final TransientPayloadStorageService transientPayloadStorageService;
private final InstanceMetadata instanceMetadata;
private final ShopDataCrawlDailyFileService dailyFileService;
@Value("${aiimage.shop-data-crawl.stale-timeout-minutes:30}")
private long staleTimeoutMinutes;
@@ -601,18 +611,29 @@ public class ShopDataCrawlTaskService {
}
ensureTaskOwnedByCurrentInstance(task, "delete shop data crawl task");
try (TaskDistributedLockService.LockHandle ignored = acquireTaskLockOrThrow(taskId)) {
List<String> resultFileUrls = listTaskRows(taskId).stream()
.map(FileResultEntity::getResultFileUrl).filter(url -> !blank(url)).distinct().toList();
fileResultMapper.delete(new LambdaQueryWrapper<FileResultEntity>()
.eq(FileResultEntity::getTaskId, taskId)
.eq(FileResultEntity::getModuleType, MODULE_TYPE));
taskFileJobService.deleteTaskJobs(taskId, MODULE_TYPE);
taskResultItemService.deleteTaskItems(taskId, MODULE_TYPE);
taskProgressSnapshotService.delete(taskId, MODULE_TYPE);
fileTaskMapper.deleteById(taskId);
taskCacheService.deleteTaskCache(taskId);
deleteTransientResultChunks(taskId);
resultFileUrls.forEach(this::deleteResultObjectIfUnreferenced);
List<FileResultEntity> taskRows = listTaskRows(taskId);
try (DailyLockSet dailyLocks = acquireDailyLocks(task.getUserId(), taskRows)) {
Set<Long> removedResultIds = taskRows.stream()
.map(FileResultEntity::getId)
.filter(id -> id != null && id > 0)
.collect(java.util.stream.Collectors.toSet());
DailyDeletionResult dailyResult = prepareDailyForDeletion(removedResultIds);
registerUploadedObjectRollback(dailyResult.uploadedObjectKeys());
List<String> resultFileUrls = new ArrayList<>(dailyResult.obsoleteObjectKeys());
resultFileUrls.addAll(taskRows.stream()
.map(FileResultEntity::getResultFileUrl).filter(url -> !blank(url)).distinct().toList());
resultFileUrls = resultFileUrls.stream().filter(url -> !blank(url)).distinct().toList();
fileResultMapper.delete(new LambdaQueryWrapper<FileResultEntity>()
.eq(FileResultEntity::getTaskId, taskId)
.eq(FileResultEntity::getModuleType, MODULE_TYPE));
taskFileJobService.deleteTaskJobs(taskId, MODULE_TYPE);
taskResultItemService.deleteTaskItems(taskId, MODULE_TYPE);
taskProgressSnapshotService.delete(taskId, MODULE_TYPE);
fileTaskMapper.deleteById(taskId);
taskCacheService.deleteTaskCache(taskId);
deleteTransientResultChunks(taskId);
resultFileUrls.forEach(this::deleteResultObjectIfUnreferenced);
}
}
}
@@ -632,7 +653,9 @@ public class ShopDataCrawlTaskService {
if (latestEntity == null || !MODULE_TYPE.equals(latestEntity.getModuleType()) || !userId.equals(latestEntity.getUserId())) {
throw new BusinessException("记录不存在");
}
deleteResultHistoryRow(latestEntity);
try (DailyLockSet dailyLocks = acquireDailyLocks(userId, List.of(latestEntity))) {
deleteResultHistoryRow(latestEntity);
}
}
}
@@ -661,6 +684,8 @@ public class ShopDataCrawlTaskService {
Long taskId = entity.getTaskId();
Long resultId = entity.getId();
String resultFileUrl = entity.getResultFileUrl();
DailyDeletionResult dailyResult = prepareDailyForDeletion(Set.of(resultId));
registerUploadedObjectRollback(dailyResult.uploadedObjectKeys());
taskFileJobService.deleteResultJobs(taskId, MODULE_TYPE, resultId);
taskResultItemService.deleteResultItem(taskId, MODULE_TYPE, resultId);
fileResultMapper.deleteById(resultId);
@@ -673,7 +698,161 @@ public class ShopDataCrawlTaskService {
taskId, resultId, safeMessage(ex));
}
// Run this after the row delete so shared object references are counted correctly.
deleteResultObjectIfUnreferenced(resultFileUrl);
Set<String> obsoleteObjectKeys = new HashSet<>(dailyResult.obsoleteObjectKeys());
collectObjectKey(obsoleteObjectKeys, resultFileUrl);
obsoleteObjectKeys.forEach(this::deleteResultObjectIfUnreferenced);
}
private DailyDeletionResult prepareDailyForDeletion(Set<Long> removedResultIds) {
if (removedResultIds == null || removedResultIds.isEmpty()) {
return new DailyDeletionResult(List.of(), List.of());
}
Set<Long> handledDailyIds = new HashSet<>();
Set<String> obsoleteObjectKeys = new HashSet<>();
List<String> uploadedObjectKeys = new ArrayList<>();
try {
for (Long resultId : removedResultIds) {
List<ShopDataCrawlDailyFileEntity> latestFiles = dailyFileService.findByLatestResultId(resultId);
List<ShopDataCrawlDailyFileEntity> affectedFiles = new ArrayList<>(
latestFiles == null ? List.of() : latestFiles);
List<ShopDataCrawlDailyMemberEntity> memberRows = dailyFileService.findMembersByResultId(resultId);
for (ShopDataCrawlDailyMemberEntity member : memberRows == null ? List.<ShopDataCrawlDailyMemberEntity>of() : memberRows) {
if (member == null) {
continue;
}
ShopDataCrawlDailyFileEntity memberFile = dailyFileService.findById(member.getDailyFileId());
if (memberFile != null) {
affectedFiles.add(memberFile);
}
}
for (ShopDataCrawlDailyFileEntity dailyFile : affectedFiles) {
if (dailyFile == null || dailyFile.getId() == null || !handledDailyIds.add(dailyFile.getId())) {
continue;
}
ShopDataCrawlDailyFileEntity lockedDailyFile = dailyFileService.findForUpdate(
dailyFile.getUserId(), dailyFile.getShopKeyHash(), dailyFile.getBusinessDate());
if (lockedDailyFile != null) {
dailyFile = lockedDailyFile;
}
List<DailyMemberData> survivors = loadDailyMemberData(dailyFile, removedResultIds);
if (survivors.isEmpty()) {
collectObjectKey(obsoleteObjectKeys, dailyFile.getResultFileUrl());
dailyFileService.deleteDailyFile(dailyFile.getId());
continue;
}
rebuildDailyFileAfterDeletion(dailyFile, survivors, uploadedObjectKeys, obsoleteObjectKeys);
}
}
dailyFileService.deleteMembersForResults(removedResultIds);
return new DailyDeletionResult(new ArrayList<>(obsoleteObjectKeys), uploadedObjectKeys);
} catch (RuntimeException ex) {
registerRollbackObjectCleanup(uploadedObjectKeys);
throw ex;
}
}
private List<DailyMemberData> loadDailyMemberData(ShopDataCrawlDailyFileEntity dailyFile,
Set<Long> removedResultIds) {
List<ShopDataCrawlDailyMemberEntity> memberRows = dailyFileService.listMembers(dailyFile.getId());
List<ShopDataCrawlDailyMemberEntity> members = new ArrayList<>(
memberRows == null ? List.of() : memberRows);
members.removeIf(Objects::isNull);
members.sort(Comparator
.comparing(ShopDataCrawlDailyMemberEntity::getCreatedAt,
Comparator.nullsLast(Comparator.naturalOrder()))
.thenComparing(ShopDataCrawlDailyMemberEntity::getId,
Comparator.nullsLast(Comparator.naturalOrder())));
List<DailyMemberData> survivors = new ArrayList<>();
for (ShopDataCrawlDailyMemberEntity member : members) {
if (removedResultIds.contains(member.getResultId())) {
continue;
}
FileResultEntity result = fileResultMapper.selectById(member.getResultId());
if (result == null || !Integer.valueOf(RESULT_SUCCESS).equals(result.getSuccess())) {
continue;
}
ShopDataCrawlResultItemVo snapshot = loadSnapshotForDailyMember(result);
if (snapshot == null || !Boolean.TRUE.equals(snapshot.getSuccess())) {
throw new BusinessException("无法读取累计文件中的剩余结果");
}
survivors.add(new DailyMemberData(member, result, snapshot));
}
return survivors;
}
private ShopDataCrawlResultItemVo loadSnapshotForDailyMember(FileResultEntity result) {
ShopDataCrawlResultItemVo snapshot = taskResultItemService.getResultSnapshot(
result.getTaskId(), MODULE_TYPE, result.getId(), ShopDataCrawlResultItemVo.class);
if (snapshot != null) {
return snapshot;
}
FileTaskEntity task = fileTaskMapper.selectById(result.getTaskId());
if (task == null) {
return null;
}
return indexSnapshotByResultId(buildSnapshotFromDb(task, List.of(result))).get(result.getId());
}
private void rebuildDailyFileAfterDeletion(ShopDataCrawlDailyFileEntity dailyFile,
List<DailyMemberData> survivors,
List<String> uploadedObjectKeys,
Set<String> obsoleteObjectKeys) {
DailyMemberData latest = survivors.get(survivors.size() - 1);
FileTaskEntity latestTask = fileTaskMapper.selectById(latest.result().getTaskId());
String filename = blank(dailyFile.getResultFilename())
? buildTaskWorkbookFilename(latestTask)
: dailyFile.getResultFilename();
File workRoot = FileUtil.mkdir(FileUtil.file(
System.getProperty("java.io.tmpdir"),
"shop-data-crawl-result",
"daily-delete-" + UUID.randomUUID()));
File outputXlsx = FileUtil.file(workRoot, filename);
String oldObjectKey = dailyFile.getResultFileUrl();
List<ShopDataCrawlResultItemVo> snapshots = survivors.stream()
.map(DailyMemberData::snapshot)
.toList();
try {
excelAssemblyService.writeWorkbook(outputXlsx, snapshots);
String newObjectKey = ossStorageService.uploadResultFile(outputXlsx, MODULE_TYPE);
if (blank(newObjectKey)) {
throw new BusinessException("累计文件上传结果为空");
}
uploadedObjectKeys.add(newObjectKey);
int rowCount = excelAssemblyService.countRows(snapshots);
for (DailyMemberData survivor : survivors) {
FileResultEntity result = survivor.result();
if (!Objects.equals(result.getId(), latest.result().getId()) && !blank(result.getResultFileUrl())) {
result.setResultFileUrl(null);
result.setResultFileSize(null);
result.setResultContentType(null);
fileResultMapper.updateById(result);
}
}
latest.result().setResultFilename(filename);
latest.result().setResultFileUrl(newObjectKey);
latest.result().setResultFileSize(outputXlsx.length());
latest.result().setResultContentType(CONTENT_TYPE_XLSX);
latest.result().setRowCount(rowCount);
fileResultMapper.updateById(latest.result());
LocalDateTime now = dailyFileService.currentBusinessDateTime();
dailyFile.setLatestTaskId(latest.result().getTaskId());
dailyFile.setLatestResultId(latest.result().getId());
dailyFile.setResultFilename(filename);
dailyFile.setResultFileUrl(newObjectKey);
dailyFile.setResultFileSize(latest.result().getResultFileSize());
dailyFile.setResultContentType(CONTENT_TYPE_XLSX);
dailyFile.setRowCount(rowCount);
dailyFile.setVersion(Math.max(0L, Objects.requireNonNullElse(dailyFile.getVersion(), 0L)) + 1L);
dailyFile.setLastSuccessAt(now);
dailyFile.setUpdatedAt(now);
dailyFileService.update(dailyFile);
collectObjectKey(obsoleteObjectKeys, oldObjectKey);
obsoleteObjectKeys.remove(newObjectKey);
} finally {
FileUtil.del(outputXlsx);
FileUtil.del(workRoot);
}
}
private void reconcileTaskAfterResultRemoval(Long taskId) {
@@ -700,131 +879,6 @@ public class ShopDataCrawlTaskService {
taskCacheService.saveTaskCache(task);
}
private void pruneCompletedHistoryQuietly(FileTaskEntity currentTask, List<FileResultEntity> currentRows) {
if (currentRows == null || currentRows.isEmpty()) {
return;
}
for (FileResultEntity row : currentRows) {
if (!isRetentionCandidate(row)) {
continue;
}
Long userId = row.getUserId() != null
? row.getUserId()
: currentTask == null ? null : currentTask.getUserId();
String shopKey = retentionShopKey(row);
if (userId == null || shopKey == null) {
log.warn("[shop-data-crawl] skip history retention because ownership key is incomplete taskId={} resultId={}",
currentTask == null ? null : currentTask.getId(), row.getId());
continue;
}
try {
pruneCompletedHistoryForShop(userId, shopKey);
} catch (Exception ex) {
// Retention is best effort. A cleanup failure must not fail the newly assembled workbook job.
log.warn("[shop-data-crawl] history retention failed taskId={} resultId={} msg={}",
currentTask == null ? null : currentTask.getId(), row.getId(), safeMessage(ex));
}
}
}
void pruneCompletedHistoryForShop(Long userId, String shopKey) {
if (userId == null || shopKey == null) {
return;
}
List<FileResultEntity> candidates = fileResultMapper.selectList(new LambdaQueryWrapper<FileResultEntity>()
.select(FileResultEntity::getId,
FileResultEntity::getTaskId,
FileResultEntity::getModuleType,
FileResultEntity::getSourceFilename,
FileResultEntity::getSourceFileUrl,
FileResultEntity::getResultFileUrl,
FileResultEntity::getSuccess,
FileResultEntity::getUserId,
FileResultEntity::getCreatedAt)
.eq(FileResultEntity::getModuleType, MODULE_TYPE)
.eq(FileResultEntity::getUserId, userId)
.eq(FileResultEntity::getSuccess, RESULT_SUCCESS)
.isNotNull(FileResultEntity::getResultFileUrl)
.ne(FileResultEntity::getResultFileUrl, "")
.orderByDesc(FileResultEntity::getCreatedAt)
.orderByDesc(FileResultEntity::getId));
if (candidates == null || candidates.isEmpty()) {
return;
}
List<FileResultEntity> shopResults = new ArrayList<>();
for (FileResultEntity candidate : candidates) {
if (!isRetentionCandidate(candidate)
|| !Objects.equals(userId, candidate.getUserId())) {
continue;
}
if (!Objects.equals(shopKey, retentionShopKey(candidate))) {
continue;
}
shopResults.add(candidate);
}
Comparator<FileResultEntity> newestFirst = Comparator
.comparing(FileResultEntity::getCreatedAt, Comparator.nullsLast(Comparator.reverseOrder()))
.thenComparing(FileResultEntity::getId, Comparator.nullsLast(Comparator.reverseOrder()));
shopResults.sort(newestFirst);
for (int index = SHOP_HISTORY_RETENTION_LIMIT; index < shopResults.size(); index++) {
deleteRetentionResultQuietly(shopResults.get(index), userId, shopKey);
}
}
private void deleteRetentionResultQuietly(FileResultEntity candidate, Long userId, String shopKey) {
if (candidate == null || candidate.getId() == null || candidate.getId() <= 0
|| candidate.getTaskId() == null || candidate.getTaskId() <= 0) {
return;
}
try {
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(candidate.getTaskId());
if (lockHandle == null) {
log.info("[shop-data-crawl] skip retained-history deletion because task lock is busy taskId={} resultId={}",
candidate.getTaskId(), candidate.getId());
return;
}
try (lockHandle) {
FileResultEntity latest = fileResultMapper.selectById(candidate.getId());
if (!isRetentionCandidate(latest)
|| !Objects.equals(userId, latest.getUserId())
|| !Objects.equals(shopKey, retentionShopKey(latest))) {
return;
}
FileTaskEntity task = fileTaskMapper.selectById(candidate.getTaskId());
if (task == null || !MODULE_TYPE.equals(task.getModuleType())
|| !isTerminalTaskStatus(task.getStatus())) {
log.info("[shop-data-crawl] skip retained-history deletion because task is not terminal taskId={} resultId={}",
candidate.getTaskId(), candidate.getId());
return;
}
deleteResultHistoryRow(latest);
}
} catch (Exception ex) {
log.warn("[shop-data-crawl] retained-history deletion failed taskId={} resultId={} msg={}",
candidate.getTaskId(), candidate.getId(), safeMessage(ex));
}
}
private boolean isRetentionCandidate(FileResultEntity row) {
return row != null
&& Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())
&& !blank(row.getResultFileUrl());
}
private String retentionShopKey(FileResultEntity row) {
if (row == null) {
return null;
}
String shopId = trimToNull(row.getSourceFileUrl());
if (shopId != null) {
return "shop-id:" + shopId;
}
String shopName = trimToNull(row.getSourceFilename());
return shopName == null ? null : "shop-name:" + shopName;
}
private String trimToNull(String value) {
if (value == null) {
return null;
@@ -1519,6 +1573,7 @@ public class ShopDataCrawlTaskService {
}
}
@Transactional
public void processResultFileJob(TaskFileJobEntity job) {
if (job == null || job.getTaskId() == null) {
throw new BusinessException("结果文件任务参数不完整");
@@ -1539,31 +1594,434 @@ public class ShopDataCrawlTaskService {
if (successItems.isEmpty()) {
throw new BusinessException("没有可生成的店铺数据抓取结果");
}
File workRoot = FileUtil.mkdir(FileUtil.file(System.getProperty("java.io.tmpdir"), "shop-data-crawl-result", String.valueOf(task.getId())));
String filename = buildTaskWorkbookFilename(task);
File xlsx = FileUtil.file(workRoot, filename);
Map<Long, ShopDataCrawlResultItemVo> snapshotsByResultId = indexSnapshotByResultId(successItems);
LocalDate businessDate = dailyFileService.currentBusinessDate();
List<String> uploadedObjectKeys = new ArrayList<>();
List<String> obsoleteObjectKeys = new ArrayList<>();
try {
excelAssemblyService.writeWorkbook(xlsx, successItems);
String objectKey = ossStorageService.uploadResultFile(xlsx, MODULE_TYPE);
long fileSize = xlsx.length();
int rowCount = excelAssemblyService.countRows(successItems);
for (FileResultEntity row : rows) {
if (Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
row.setResultFilename(filename);
row.setResultFileUrl(objectKey);
row.setResultFileSize(fileSize);
row.setResultContentType(CONTENT_TYPE_XLSX);
row.setRowCount(rowCount);
fileResultMapper.updateById(row);
if (!Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
continue;
}
ShopDataCrawlResultItemVo snapshot = snapshotsByResultId.get(row.getId());
if (snapshot == null) {
snapshot = successItems.size() == 1 ? successItems.get(0) : null;
}
if (snapshot == null) {
throw new BusinessException("店铺结果快照不存在,无法生成累计文件");
}
DailyAggregationResult result = aggregateDailyResult(
task, row, snapshot, businessDate, uploadedObjectKeys);
obsoleteObjectKeys.addAll(result.obsoleteObjectKeys());
}
updateTaskStatusFromRows(task, rows);
persistSnapshotJson(task, buildSnapshotFromDb(task, rows));
fileTaskMapper.updateById(task);
} finally {
FileUtil.del(xlsx);
} catch (RuntimeException ex) {
registerRollbackObjectCleanup(uploadedObjectKeys);
throw ex;
}
pruneCompletedHistoryQuietly(task, rows);
registerDailyObjectLifecycle(uploadedObjectKeys, obsoleteObjectKeys);
}
private DailyAggregationResult aggregateDailyResult(FileTaskEntity task,
FileResultEntity row,
ShopDataCrawlResultItemVo snapshot,
LocalDate businessDate,
List<String> uploadedObjectKeys) {
Long userId = row.getUserId() != null ? row.getUserId() : task.getUserId();
String shopKey = dailyFileService.shopKey(row);
String shopKeyHash = dailyFileService.shopKeyHash(shopKey);
if (userId == null || shopKeyHash == null) {
throw new BusinessException("店铺累计文件归属信息不完整");
}
TaskDistributedLockService.LockHandle lock = dailyFileService.acquireLock(userId, shopKey);
if (lock == null) {
throw new BusinessException("店铺当天累计文件正在处理中,请稍后重试");
}
try {
ShopDataCrawlDailyFileEntity dailyFile = dailyFileService.findForUpdate(userId, shopKeyHash, businessDate);
ShopDataCrawlDailyFileEntity existingMembership = findExistingDailyMembership(row.getId());
if (existingMembership != null) {
boolean currentFileOwnsMembership = dailyFile != null
&& Objects.equals(dailyFile.getId(), existingMembership.getId());
if (Objects.equals(existingMembership.getLatestResultId(), row.getId())
&& (dailyFile == null || currentFileOwnsMembership)) {
attachCanonicalResult(row, existingMembership);
} else if (!blank(row.getResultFileUrl())) {
row.setResultFileUrl(null);
row.setResultFileSize(null);
row.setResultContentType(null);
fileResultMapper.updateById(row);
}
return new DailyAggregationResult(List.of());
}
if (dailyFile != null && dailyFileService.containsResult(dailyFile.getId(), row.getId())) {
if (Objects.equals(dailyFile.getLatestResultId(), row.getId())) {
attachCanonicalResult(row, dailyFile);
} else if (!blank(row.getResultFileUrl())) {
row.setResultFileUrl(null);
row.setResultFileSize(null);
row.setResultContentType(null);
fileResultMapper.updateById(row);
}
return new DailyAggregationResult(List.of());
}
List<ShopDataCrawlDailyFileEntity> olderFiles = dailyFileService.findOlder(userId, shopKeyHash, businessDate);
Set<String> obsoleteObjectKeys = new HashSet<>();
collectObjectKey(obsoleteObjectKeys, dailyFile == null ? null : dailyFile.getResultFileUrl());
for (ShopDataCrawlDailyFileEntity older : olderFiles) {
collectObjectKey(obsoleteObjectKeys, older.getResultFileUrl());
}
int addedRowCount = excelAssemblyService.countRows(List.of(snapshot));
String filename = dailyFile != null && !blank(dailyFile.getResultFilename())
? dailyFile.getResultFilename()
: buildTaskWorkbookFilename(task);
File workRoot = FileUtil.mkdir(FileUtil.file(
System.getProperty("java.io.tmpdir"),
"shop-data-crawl-result",
String.valueOf(task.getId()),
"daily-" + UUID.randomUUID()));
File baseXlsx = FileUtil.file(workRoot, "base.xlsx");
File outputXlsx = FileUtil.file(workRoot, filename);
String objectKey = null;
try {
if (dailyFile != null && !blank(dailyFile.getResultFileUrl())) {
try {
Files.write(baseXlsx.toPath(), ossStorageService.readObjectBytes(dailyFile.getResultFileUrl()));
} catch (Exception ex) {
throw new BusinessException("读取当天累计文件失败: " + safeMessage(ex));
}
if (addedRowCount > 0) {
excelAssemblyService.appendWorkbook(baseXlsx, outputXlsx, List.of(snapshot));
} else {
try {
Files.copy(baseXlsx.toPath(), outputXlsx.toPath());
} catch (Exception ex) {
throw new BusinessException("复制当天累计文件失败: " + safeMessage(ex));
}
}
} else {
excelAssemblyService.writeWorkbook(outputXlsx, List.of(snapshot));
}
if (dailyFile != null && addedRowCount == 0) {
objectKey = dailyFile.getResultFileUrl();
} else {
objectKey = ossStorageService.uploadResultFile(outputXlsx, MODULE_TYPE);
uploadedObjectKeys.add(objectKey);
}
List<FileResultEntity> shopRows = findShopResultRows(userId, row);
clearShopResultPointers(shopRows, row.getId());
row.setResultFilename(filename);
row.setResultFileUrl(objectKey);
row.setResultFileSize(outputXlsx.length());
row.setResultContentType(CONTENT_TYPE_XLSX);
row.setRowCount(dailyFile == null
? addedRowCount
: Math.max(0, Objects.requireNonNullElse(dailyFile.getRowCount(), 0)) + addedRowCount);
fileResultMapper.updateById(row);
LocalDateTime now = dailyFileService.currentBusinessDateTime();
if (dailyFile == null) {
dailyFile = new ShopDataCrawlDailyFileEntity();
dailyFile.setUserId(userId);
dailyFile.setShopKeyHash(shopKeyHash);
dailyFile.setShopKey(shopKey);
dailyFile.setBusinessDate(businessDate);
dailyFile.setVersion(1L);
dailyFile.setCreatedAt(now);
} else {
dailyFile.setVersion(Math.max(0L, Objects.requireNonNullElse(dailyFile.getVersion(), 0L)) + 1L);
}
dailyFile.setLatestTaskId(row.getTaskId());
dailyFile.setLatestResultId(row.getId());
dailyFile.setResultFilename(filename);
dailyFile.setResultFileUrl(objectKey);
dailyFile.setResultFileSize(row.getResultFileSize());
dailyFile.setResultContentType(CONTENT_TYPE_XLSX);
dailyFile.setRowCount(row.getRowCount());
dailyFile.setLastSuccessAt(now);
dailyFile.setUpdatedAt(now);
if (dailyFile.getId() == null) {
dailyFileService.insert(dailyFile);
} else {
dailyFileService.update(dailyFile);
}
if (!dailyFileService.addMember(dailyFile.getId(), row.getTaskId(), row.getId())) {
throw new BusinessException("结果已归档,请重试文件任务");
}
for (ShopDataCrawlDailyFileEntity older : olderFiles) {
dailyFileService.deleteDailyFile(older.getId());
}
obsoleteObjectKeys.remove(objectKey);
return new DailyAggregationResult(new ArrayList<>(obsoleteObjectKeys));
} finally {
FileUtil.del(baseXlsx);
FileUtil.del(outputXlsx);
FileUtil.del(workRoot);
}
} finally {
releaseDailyLockAfterTransaction(lock);
}
}
private ShopDataCrawlDailyFileEntity findExistingDailyMembership(Long resultId) {
List<ShopDataCrawlDailyMemberEntity> members = dailyFileService.findMembersByResultId(resultId);
for (ShopDataCrawlDailyMemberEntity member : members == null ? List.<ShopDataCrawlDailyMemberEntity>of() : members) {
if (member == null) {
continue;
}
ShopDataCrawlDailyFileEntity dailyFile = dailyFileService.findById(member.getDailyFileId());
if (dailyFile != null) {
return dailyFile;
}
}
return null;
}
private void attachCanonicalResult(FileResultEntity row, ShopDataCrawlDailyFileEntity dailyFile) {
if (row == null || dailyFile == null) {
return;
}
row.setResultFilename(dailyFile.getResultFilename());
row.setResultFileUrl(dailyFile.getResultFileUrl());
row.setResultFileSize(dailyFile.getResultFileSize());
row.setResultContentType(dailyFile.getResultContentType());
row.setRowCount(dailyFile.getRowCount());
fileResultMapper.updateById(row);
}
private List<FileResultEntity> findShopResultRows(Long userId, FileResultEntity sourceRow) {
if (userId == null || sourceRow == null) {
return List.of();
}
String shopId = trimToNull(sourceRow.getSourceFileUrl());
String shopName = trimToNull(sourceRow.getSourceFilename());
LambdaQueryWrapper<FileResultEntity> wrapper = new LambdaQueryWrapper<FileResultEntity>()
.eq(FileResultEntity::getModuleType, MODULE_TYPE)
.and(owner -> owner.eq(FileResultEntity::getUserId, userId)
.or().isNull(FileResultEntity::getUserId));
if (shopId != null) {
wrapper.eq(FileResultEntity::getSourceFileUrl, shopId);
} else if (shopName != null) {
wrapper.eq(FileResultEntity::getSourceFilename, shopName)
.apply("TRIM(COALESCE(source_file_url, '')) = ''");
} else {
return List.of();
}
List<FileResultEntity> candidates = fileResultMapper.selectList(wrapper);
if (candidates == null || candidates.isEmpty()) {
return List.of();
}
Map<Long, FileTaskEntity> legacyTaskOwners = loadTaskMapByIds(candidates.stream()
.filter(candidate -> candidate != null && candidate.getUserId() == null)
.map(FileResultEntity::getTaskId)
.filter(Objects::nonNull)
.distinct()
.toList());
return candidates.stream()
.filter(Objects::nonNull)
.filter(candidate -> Objects.equals(userId, candidate.getUserId())
|| Objects.equals(userId,
legacyTaskOwners.get(candidate.getTaskId()) == null
? null
: legacyTaskOwners.get(candidate.getTaskId()).getUserId()))
.toList();
}
private void clearShopResultPointers(List<FileResultEntity> rows, Long keepResultId) {
if (rows == null) {
return;
}
for (FileResultEntity candidate : rows) {
if (candidate == null || Objects.equals(candidate.getId(), keepResultId)
|| blank(candidate.getResultFileUrl())) {
continue;
}
candidate.setResultFileUrl(null);
candidate.setResultFileSize(null);
candidate.setResultContentType(null);
fileResultMapper.updateById(candidate);
}
}
private void collectObjectKey(Set<String> target, String value) {
if (target == null || blank(value)) {
return;
}
target.add(value);
}
private void registerDailyObjectLifecycle(List<String> uploadedObjectKeys, List<String> obsoleteObjectKeys) {
Set<String> uploaded = uploadedObjectKeys == null ? Set.of() : new HashSet<>(uploadedObjectKeys);
Set<String> obsolete = obsoleteObjectKeys == null ? new HashSet<>() : new HashSet<>(obsoleteObjectKeys);
obsolete.removeAll(uploaded);
if (!TransactionSynchronizationManager.isSynchronizationActive()) {
obsolete.forEach(this::deleteObjectQuietly);
return;
}
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCommit() {
obsolete.forEach(ShopDataCrawlTaskService.this::deleteObjectQuietly);
}
@Override
public void afterCompletion(int status) {
if (status != STATUS_COMMITTED) {
uploaded.forEach(ShopDataCrawlTaskService.this::deleteObjectQuietly);
}
}
});
}
private void registerRollbackObjectCleanup(List<String> uploadedObjectKeys) {
if (uploadedObjectKeys == null || uploadedObjectKeys.isEmpty()) {
return;
}
Set<String> uploaded = new HashSet<>(uploadedObjectKeys);
if (!TransactionSynchronizationManager.isSynchronizationActive()) {
uploaded.forEach(this::deleteObjectQuietly);
return;
}
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCompletion(int status) {
if (status != STATUS_COMMITTED) {
uploaded.forEach(ShopDataCrawlTaskService.this::deleteObjectQuietly);
}
}
});
}
private void registerUploadedObjectRollback(List<String> uploadedObjectKeys) {
if (uploadedObjectKeys == null || uploadedObjectKeys.isEmpty()
|| !TransactionSynchronizationManager.isSynchronizationActive()) {
return;
}
Set<String> uploaded = new HashSet<>(uploadedObjectKeys);
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCompletion(int status) {
if (status != STATUS_COMMITTED) {
uploaded.forEach(ShopDataCrawlTaskService.this::deleteObjectQuietly);
}
}
});
}
private void deleteObjectQuietly(String objectKey) {
if (blank(objectKey)) {
return;
}
try {
deleteResultObjectNowIfUnreferenced(objectKey);
} catch (Exception ex) {
log.warn("[shop-data-crawl] daily object cleanup failed object={} msg={}", objectKey, safeMessage(ex));
}
}
private DailyLockSet acquireDailyLocks(Long fallbackUserId, List<FileResultEntity> rows) {
Map<String, DailyLockRequest> requests = new TreeMap<>();
if (rows != null) {
for (FileResultEntity row : rows) {
if (row == null) {
continue;
}
Long userId = row.getUserId() != null ? row.getUserId() : fallbackUserId;
String shopKey = dailyFileService.shopKey(row);
if (userId == null || userId <= 0 || blank(shopKey)) {
continue;
}
requests.putIfAbsent(userId + "|" + shopKey, new DailyLockRequest(userId, shopKey));
}
}
if (requests.isEmpty()) {
return new DailyLockSet(List.of());
}
List<TaskDistributedLockService.LockHandle> handles = new ArrayList<>();
try {
for (DailyLockRequest request : requests.values()) {
TaskDistributedLockService.LockHandle handle = dailyFileService.acquireLock(
request.userId(), request.shopKey());
if (handle == null) {
throw new BusinessException("店铺累计文件正在处理中,请稍后重试");
}
handles.add(handle);
}
return new DailyLockSet(handles);
} catch (RuntimeException ex) {
handles.forEach(TaskDistributedLockService.LockHandle::close);
throw ex;
}
}
private void releaseDailyLockAfterTransaction(TaskDistributedLockService.LockHandle lock) {
if (lock == null) {
return;
}
if (!TransactionSynchronizationManager.isSynchronizationActive()) {
lock.close();
return;
}
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCompletion(int status) {
lock.close();
}
});
}
private record DailyLockRequest(Long userId, String shopKey) {
}
private record DailyMemberData(ShopDataCrawlDailyMemberEntity member,
FileResultEntity result,
ShopDataCrawlResultItemVo snapshot) {
}
private record DailyDeletionResult(List<String> obsoleteObjectKeys,
List<String> uploadedObjectKeys) {
}
private static final class DailyLockSet implements AutoCloseable {
private final List<TaskDistributedLockService.LockHandle> handles;
private boolean closed;
private DailyLockSet(List<TaskDistributedLockService.LockHandle> handles) {
this.handles = handles == null ? List.of() : handles;
}
@Override
public void close() {
if (closed) {
return;
}
closed = true;
if (TransactionSynchronizationManager.isSynchronizationActive()) {
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCompletion(int status) {
closeNow();
}
});
} else {
closeNow();
}
}
private void closeNow() {
for (int index = handles.size() - 1; index >= 0; index--) {
handles.get(index).close();
}
}
}
private record DailyAggregationResult(List<String> obsoleteObjectKeys) {
}
public void cleanupResultFileJob(TaskFileJobEntity job) {
@@ -1813,9 +2271,29 @@ public class ShopDataCrawlTaskService {
void deleteResultObjectIfUnreferenced(String resultFileUrl) {
if (blank(resultFileUrl)) return;
if (!isResultObjectUnreferenced(resultFileUrl)) return;
if (TransactionSynchronizationManager.isSynchronizationActive()) {
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCommit() {
deleteObjectQuietly(resultFileUrl);
}
});
} else {
ossStorageService.deleteObject(resultFileUrl);
}
}
private void deleteResultObjectNowIfUnreferenced(String resultFileUrl) {
if (blank(resultFileUrl) || !isResultObjectUnreferenced(resultFileUrl)) return;
ossStorageService.deleteObject(resultFileUrl);
}
private boolean isResultObjectUnreferenced(String resultFileUrl) {
Long references = fileResultMapper.selectCount(new LambdaQueryWrapper<FileResultEntity>()
.eq(FileResultEntity::getResultFileUrl, resultFileUrl));
if (references == null || references == 0L) ossStorageService.deleteObject(resultFileUrl);
long dailyReferences = dailyFileService.countObjectReferences(resultFileUrl);
return (references == null || references == 0L) && dailyReferences == 0L;
}
public void ensureTaskOwnedByCurrentInstance(FileTaskEntity task, String operation) {
@@ -78,4 +78,28 @@ public interface ShopManageGroupMapper extends BaseMapper<ShopManageGroupEntity>
ORDER BY gm.user_id ASC
""")
List<Long> selectManagedMemberUserIds(@Param("operatorId") Long operatorId);
@Select("""
<script>
SELECT DISTINCT visible_user_id
FROM (
SELECT COALESCE(g.created_by_id, g.user_id) AS visible_user_id
FROM biz_shop_manage_group g
WHERE g.id IN
<foreach collection='groupIds' item='groupId' open='(' separator=',' close=')'>
#{groupId}
</foreach>
UNION
SELECT gm.user_id AS visible_user_id
FROM biz_shop_manage_group_member gm
WHERE gm.group_id IN
<foreach collection='groupIds' item='groupId' open='(' separator=',' close=')'>
#{groupId}
</foreach>
) visible_users
WHERE visible_user_id IS NOT NULL
ORDER BY visible_user_id ASC
</script>
""")
List<Long> selectUserIdsByGroupIds(@Param("groupIds") List<Long> groupIds);
}
@@ -9,6 +9,8 @@ import com.nanri.aiimage.modules.shopkey.mapper.ShopManageGroupMapper;
import com.nanri.aiimage.modules.shopkey.mapper.ShopManageGroupMemberMapper;
import com.nanri.aiimage.modules.shopkey.mapper.ShopManageMapper;
import com.nanri.aiimage.modules.shopkey.mapper.SkipPriceAsinMapper;
import com.nanri.aiimage.modules.invalidasin.mapper.InvalidAsinDataMapper;
import com.nanri.aiimage.modules.dedupe.mapper.DedupeTotalDataMapper;
import com.nanri.aiimage.modules.shopkey.model.dto.ShopManageGroupCreateRequest;
import com.nanri.aiimage.modules.shopkey.model.dto.ShopManageGroupUpdateRequest;
import com.nanri.aiimage.modules.shopkey.model.entity.QueryAsinEntity;
@@ -16,6 +18,8 @@ import com.nanri.aiimage.modules.shopkey.model.entity.ShopManageEntity;
import com.nanri.aiimage.modules.shopkey.model.entity.ShopManageGroupEntity;
import com.nanri.aiimage.modules.shopkey.model.entity.ShopManageGroupMemberEntity;
import com.nanri.aiimage.modules.shopkey.model.entity.SkipPriceAsinEntity;
import com.nanri.aiimage.modules.invalidasin.model.entity.InvalidAsinDataEntity;
import com.nanri.aiimage.modules.dedupe.model.entity.DedupeTotalDataEntity;
import com.nanri.aiimage.modules.shopkey.model.vo.ShopManageGroupItemVo;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
@@ -39,6 +43,8 @@ public class ShopManageGroupService {
private final ShopManageMapper shopManageMapper;
private final QueryAsinMapper queryAsinMapper;
private final SkipPriceAsinMapper skipPriceAsinMapper;
private final InvalidAsinDataMapper invalidAsinDataMapper;
private final DedupeTotalDataMapper dedupeTotalDataMapper;
private final AdminUserMapper adminUserMapper;
public List<ShopManageGroupItemVo> list() {
@@ -160,6 +166,16 @@ public class ShopManageGroupService {
if (queryAsinCount != null && queryAsinCount > 0) {
throw new BusinessException("该分组下存在查询 ASIN,无法删除");
}
Long invalidAsinDataCount = invalidAsinDataMapper.selectCount(new LambdaQueryWrapper<InvalidAsinDataEntity>()
.eq(InvalidAsinDataEntity::getGroupId, entity.getId()));
if (invalidAsinDataCount != null && invalidAsinDataCount > 0) {
throw new BusinessException("该分组下存在品牌数据库数据,无法删除");
}
Long dedupeTotalDataCount = dedupeTotalDataMapper.selectCount(new LambdaQueryWrapper<DedupeTotalDataEntity>()
.eq(DedupeTotalDataEntity::getGroupId, entity.getId()));
if (dedupeTotalDataCount != null && dedupeTotalDataCount > 0) {
throw new BusinessException("该分组下存在数据去重总数据,无法删除");
}
groupMemberMapper.delete(new LambdaQueryWrapper<ShopManageGroupMemberEntity>()
.eq(ShopManageGroupMemberEntity::getGroupId, entity.getId()));
groupMapper.deleteById(entity.getId());
@@ -0,0 +1,33 @@
CREATE TABLE IF NOT EXISTS biz_shop_data_crawl_daily_file (
id BIGINT NOT NULL AUTO_INCREMENT,
user_id BIGINT NOT NULL,
shop_key_hash CHAR(64) NOT NULL,
shop_key VARCHAR(1000) NOT NULL,
business_date DATE NOT NULL,
latest_task_id BIGINT NULL,
latest_result_id BIGINT NULL,
result_filename VARCHAR(255) NULL,
result_file_url VARCHAR(1000) NULL,
result_file_size BIGINT NULL,
result_content_type VARCHAR(128) NULL,
row_count INT NULL,
version BIGINT NOT NULL DEFAULT 0,
last_success_at DATETIME NULL,
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (id),
UNIQUE KEY uk_shop_data_crawl_daily_file (user_id, shop_key_hash, business_date),
KEY idx_shop_data_crawl_daily_result (latest_result_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
CREATE TABLE IF NOT EXISTS biz_shop_data_crawl_daily_member (
id BIGINT NOT NULL AUTO_INCREMENT,
daily_file_id BIGINT NOT NULL,
task_id BIGINT NOT NULL,
result_id BIGINT NOT NULL,
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (id),
UNIQUE KEY uk_shop_data_crawl_daily_member_result (result_id),
KEY idx_shop_data_crawl_daily_member_file_result (daily_file_id, result_id),
KEY idx_shop_data_crawl_daily_member_task (task_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
@@ -0,0 +1,74 @@
SET @invalid_asin_group_id_col_exists := (
SELECT COUNT(*)
FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA = DATABASE()
AND TABLE_NAME = 'biz_invalid_asin_data'
AND COLUMN_NAME = 'group_id'
);
SET @sql_add_invalid_asin_group_id := IF(
@invalid_asin_group_id_col_exists = 0,
'ALTER TABLE biz_invalid_asin_data ADD COLUMN group_id BIGINT NULL AFTER brand',
'SELECT 1'
);
PREPARE stmt_add_invalid_asin_group_id FROM @sql_add_invalid_asin_group_id;
EXECUTE stmt_add_invalid_asin_group_id;
DEALLOCATE PREPARE stmt_add_invalid_asin_group_id;
SET @invalid_asin_record_source_col_exists := (
SELECT COUNT(*)
FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA = DATABASE()
AND TABLE_NAME = 'biz_invalid_asin_data'
AND COLUMN_NAME = 'record_source'
);
SET @sql_add_invalid_asin_record_source := IF(
@invalid_asin_record_source_col_exists = 0,
'ALTER TABLE biz_invalid_asin_data ADD COLUMN record_source VARCHAR(16) NOT NULL DEFAULT ''AUTO'' AFTER group_id',
'SELECT 1'
);
PREPARE stmt_add_invalid_asin_record_source FROM @sql_add_invalid_asin_record_source;
EXECUTE stmt_add_invalid_asin_record_source;
DEALLOCATE PREPARE stmt_add_invalid_asin_record_source;
-- Historical rows have no reliable origin marker. Treat them as automatic rows.
UPDATE biz_invalid_asin_data
SET record_source = 'AUTO'
WHERE record_source IS NULL
OR TRIM(record_source) = ''
OR UPPER(TRIM(record_source)) NOT IN ('AUTO', 'MANUAL');
SET @invalid_asin_source_group_index_exists := (
SELECT COUNT(*)
FROM information_schema.STATISTICS
WHERE TABLE_SCHEMA = DATABASE()
AND TABLE_NAME = 'biz_invalid_asin_data'
AND INDEX_NAME = 'idx_invalid_asin_source_group_id'
);
SET @sql_add_invalid_asin_source_group_index := IF(
@invalid_asin_source_group_index_exists = 0,
'ALTER TABLE biz_invalid_asin_data ADD INDEX idx_invalid_asin_source_group_id (record_source, group_id, id)',
'SELECT 1'
);
PREPARE stmt_add_invalid_asin_source_group_index FROM @sql_add_invalid_asin_source_group_index;
EXECUTE stmt_add_invalid_asin_source_group_index;
DEALLOCATE PREPARE stmt_add_invalid_asin_source_group_index;
SET @invalid_asin_group_index_exists := (
SELECT COUNT(*)
FROM information_schema.STATISTICS
WHERE TABLE_SCHEMA = DATABASE()
AND TABLE_NAME = 'biz_invalid_asin_data'
AND INDEX_NAME = 'idx_invalid_asin_group_id'
);
SET @sql_add_invalid_asin_group_index := IF(
@invalid_asin_group_index_exists = 0,
'ALTER TABLE biz_invalid_asin_data ADD INDEX idx_invalid_asin_group_id (group_id)',
'SELECT 1'
);
PREPARE stmt_add_invalid_asin_group_index FROM @sql_add_invalid_asin_group_index;
EXECUTE stmt_add_invalid_asin_group_index;
DEALLOCATE PREPARE stmt_add_invalid_asin_group_index;
@@ -0,0 +1,53 @@
SET @dedupe_group_id_col_exists := (
SELECT COUNT(*)
FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA = DATABASE()
AND TABLE_NAME = 'biz_dedupe_total_data'
AND COLUMN_NAME = 'group_id'
);
SET @sql_add_dedupe_group_id := IF(
@dedupe_group_id_col_exists = 0,
'ALTER TABLE biz_dedupe_total_data ADD COLUMN group_id BIGINT NULL AFTER data_value',
'SELECT 1'
);
PREPARE stmt_add_dedupe_group_id FROM @sql_add_dedupe_group_id;
EXECUTE stmt_add_dedupe_group_id;
DEALLOCATE PREPARE stmt_add_dedupe_group_id;
UPDATE biz_dedupe_total_data d
JOIN (
SELECT matched.id, MIN(matched.group_id) AS group_id
FROM (
SELECT d0.id, g.id AS group_id
FROM biz_dedupe_total_data d0
INNER JOIN biz_shop_manage_group g
ON COALESCE(g.created_by_id, g.user_id) = d0.uploader_user_id
UNION
SELECT d1.id, gm.group_id
FROM biz_dedupe_total_data d1
INNER JOIN biz_shop_manage_group_member gm
ON gm.user_id = d1.uploader_user_id
) matched
GROUP BY matched.id
HAVING COUNT(DISTINCT matched.group_id) = 1
) resolved ON resolved.id = d.id
SET d.group_id = resolved.group_id
WHERE d.group_id IS NULL;
SET @dedupe_group_idx_exists := (
SELECT COUNT(*)
FROM information_schema.STATISTICS
WHERE TABLE_SCHEMA = DATABASE()
AND TABLE_NAME = 'biz_dedupe_total_data'
AND INDEX_NAME = 'idx_dedupe_total_data_group_id'
);
SET @sql_add_dedupe_group_idx := IF(
@dedupe_group_idx_exists = 0,
'ALTER TABLE biz_dedupe_total_data ADD INDEX idx_dedupe_total_data_group_id (group_id, id)',
'SELECT 1'
);
PREPARE stmt_add_dedupe_group_idx FROM @sql_add_dedupe_group_idx;
EXECUTE stmt_add_dedupe_group_idx;
DEALLOCATE PREPARE stmt_add_dedupe_group_idx;