历史上 40901 被两种语义复用:「任务已结束」(幂等忽略)与「任务正在处理中」(分布式锁竞争,
+ * 应重试);而 GlobalExceptionHandler 把 40901 一律转成 HTTP 200 + success=true,
+ * 导致锁竞争时 Python worker 误判回传成功并停止重试,分片数据静默丢失。
+ * 现拆分为两个码:
+ *
handleBusinessException(BusinessException ex) {
- if (Integer.valueOf(40901).equals(ex.getCode())) {
+ if (Integer.valueOf(BusinessCodes.TASK_ALREADY_FINISHED).equals(ex.getCode())) {
+ // 幂等忽略:任务已结束时的重复提交无副作用,按成功返回,避免客户端反复重试
return ApiResponse.success("任务已结束,忽略重复提交", null);
}
+ if (Integer.valueOf(BusinessCodes.TASK_BUSY).equals(ex.getCode())) {
+ // 锁竞争:必须如实返回失败 + 可重试码,否则 worker 会把「未落库」当成功而停止重试
+ log.warn("[business] 任务忙,调用方应稍后重试: {}", ex.getMessage());
+ return ApiResponse.fail(BusinessCodes.TASK_BUSY, ex.getMessage());
+ }
return ex.getCode() == null
? ApiResponse.fail(ex.getMessage())
: ApiResponse.fail(ex.getCode(), ex.getMessage());
@@ -100,7 +106,9 @@ public class GlobalExceptionHandler {
return ApiResponse.fail("客户端已断开连接");
}
log.error("Unhandled exception", ex);
- return ApiResponse.fail("服务异常: " + ex.getMessage());
+ // 不再回传原始异常信息:SQL 报错、类名与内部路径会直接暴露给调用方,便于攻击者
+ // 摸清技术栈与表结构。详情只进日志(上方 log.error 已带完整堆栈),对外统一文案。
+ return ApiResponse.fail("服务器内部错误,请稍后重试");
}
private boolean isClientAbort(Throwable ex) {
diff --git a/backend-java/src/main/java/com/nanri/aiimage/common/util/SecretMasking.java b/backend-java/src/main/java/com/nanri/aiimage/common/util/SecretMasking.java
new file mode 100644
index 00000000..7faee31d
--- /dev/null
+++ b/backend-java/src/main/java/com/nanri/aiimage/common/util/SecretMasking.java
@@ -0,0 +1,52 @@
+package com.nanri.aiimage.common.util;
+
+import java.net.URI;
+
+/**
+ * 敏感串脱敏工具(2026-09 全维度审查后收口)。
+ *
+ * 背景:代理提取链接形态为 {@code http://user:pass@host:port},代码中曾有多处
+ * 直接把整串打进日志,导致用户代理账号密码长期留在应用日志与 docker logs 里。
+ */
+public final class SecretMasking {
+
+ private SecretMasking() {
+ }
+
+ /** 通用掩码:保留前 4 与后 4 字符;过短则整体掩掉。 */
+ public static String mask(String value) {
+ if (value == null || value.isBlank()) {
+ return "";
+ }
+ String text = value.trim();
+ if (text.length() <= 8) {
+ return "****";
+ }
+ return text.substring(0, 4) + "****" + text.substring(text.length() - 4);
+ }
+
+ /** 代理掩码:隐去账号密码,保留 scheme://host:port 便于运维核对。 */
+ public static String maskProxy(String value) {
+ if (value == null || value.isBlank()) {
+ return "";
+ }
+ try {
+ URI uri = URI.create(value.trim());
+ if (uri.getHost() == null || uri.getHost().isBlank()) {
+ return mask(value);
+ }
+ StringBuilder masked = new StringBuilder();
+ masked.append(uri.getScheme() == null ? "http" : uri.getScheme()).append("://");
+ if (uri.getUserInfo() != null && !uri.getUserInfo().isBlank()) {
+ masked.append("***@");
+ }
+ masked.append(uri.getHost());
+ if (uri.getPort() > 0) {
+ masked.append(':').append(uri.getPort());
+ }
+ return masked.toString();
+ } catch (Exception ex) {
+ return mask(value);
+ }
+ }
+}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/common/util/ThreadPools.java b/backend-java/src/main/java/com/nanri/aiimage/common/util/ThreadPools.java
new file mode 100644
index 00000000..27977983
--- /dev/null
+++ b/backend-java/src/main/java/com/nanri/aiimage/common/util/ThreadPools.java
@@ -0,0 +1,47 @@
+package com.nanri.aiimage.common.util;
+
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * 有界线程池工厂(2026-09 全维度审查补)。
+ *
+ *
业务里多处使用 {@code Executors.newFixedThreadPool}:它内部是**无界**
+ * {@code LinkedBlockingQueue},任务堆积时永远不会触发拒绝策略,会把内存吃满
+ * (表现为 OOM,或整机因 GC 变慢导致所有任务一起劣化)。
+ * 统一改为有界队列 + {@code CallerRunsPolicy}:队列满时在提交线程执行,形成天然背压。
+ */
+public final class ThreadPools {
+
+ private ThreadPools() {
+ }
+
+ /** 默认队列容量:足以吸收突发,又不至于无界堆积。 */
+ public static final int DEFAULT_QUEUE_CAPACITY = 512;
+
+ /** 有界固定线程池(daemon 线程,空闲可回收)。 */
+ public static ExecutorService boundedFixed(String threadNamePrefix, int threads) {
+ return boundedFixed(threadNamePrefix, threads, DEFAULT_QUEUE_CAPACITY);
+ }
+
+ /** 有界固定线程池(显式队列容量)。 */
+ public static ExecutorService boundedFixed(String threadNamePrefix, int threads, int queueCapacity) {
+ ThreadPoolExecutor executor = new ThreadPoolExecutor(
+ Math.max(1, threads),
+ Math.max(1, threads),
+ // keepAliveTime 必须 > 0:下面开了 allowCoreThreadTimeOut,
+ // 传 0 会让构造器直接抛 "Core threads must have nonzero keep alive times"
+ 60L, TimeUnit.SECONDS,
+ new LinkedBlockingQueue<>(Math.max(1, queueCapacity)),
+ runnable -> {
+ Thread thread = new Thread(runnable, threadNamePrefix);
+ thread.setDaemon(true);
+ return thread;
+ },
+ new ThreadPoolExecutor.CallerRunsPolicy());
+ executor.allowCoreThreadTimeOut(true);
+ return executor;
+ }
+}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/config/AdminApiGuardFilter.java b/backend-java/src/main/java/com/nanri/aiimage/config/AdminApiGuardFilter.java
index aee32331..b130cac5 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/config/AdminApiGuardFilter.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/config/AdminApiGuardFilter.java
@@ -63,6 +63,25 @@ public class AdminApiGuardFilter extends OncePerRequestFilter {
private static final String[] USER_TOOL_PREFIXES = {
"/api/collect-data",
"/api/price-track",
+ // 2026-09 全维度审查补:以下前缀此前完全在守卫范围之外。
+ // /api/files 匿名可上传(2GB/次,可耗尽临时盘);/api/digital-human 匿名可发布/删除版本;
+ // /api/image-video 的 secrets 接口 userId 取自请求体;/api/brand 的 fileUrl 曾可直接请求任意地址。
+ "/api/files",
+ "/api/digital-human",
+ "/api/image-video",
+ "/api/brand",
+ "/api/appearance-patent",
+ "/api/similar-asin",
+ "/api/query-asin",
+ "/api/patrol-delete",
+ "/api/product-risk-resolve",
+ "/api/shop-match",
+ "/api/shop-data-crawl",
+ "/api/withdraw",
+ "/api/task-file-jobs",
+ // /api/tasks/{taskId}/interrupted 仅凭 taskId 即可把 RUNNING 任务置为 FAILED,
+ // 匿名遍历 taskId 就能批量打断线上任务
+ "/api/tasks",
};
/**
@@ -72,6 +91,9 @@ public class AdminApiGuardFilter extends OncePerRequestFilter {
private static final String[] SELF_SERVICE_PREFIXES = {
"/api/user-secrets",
"/api/notifications",
+ // 2026-09 全维度审查补:内部端点此前仅靠 controller 自校验令牌,纳入守卫后
+ // 不带令牌的请求直接 401(带可信令牌的仍由 doFilterInternal 放行)
+ "/api/internal",
};
private final AdminAuthSupport adminAuthSupport;
diff --git a/backend-java/src/main/java/com/nanri/aiimage/config/HttpClientPool.java b/backend-java/src/main/java/com/nanri/aiimage/config/HttpClientPool.java
index d1d7db09..be13e670 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/config/HttpClientPool.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/config/HttpClientPool.java
@@ -3,12 +3,16 @@ package com.nanri.aiimage.config;
import org.springframework.http.client.ClientHttpRequestFactory;
import org.springframework.http.client.JdkClientHttpRequestFactory;
+import java.io.IOException;
+import java.io.InputStream;
import java.net.Authenticator;
import java.net.InetSocketAddress;
import java.net.PasswordAuthentication;
import java.net.ProxySelector;
import java.net.URI;
import java.net.http.HttpClient;
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
import java.time.Duration;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@@ -55,6 +59,38 @@ public class HttpClientPool {
}
}
+ /**
+ * 打开远程文件流(带超时),调用方负责关闭返回的流。
+ *
+ *
替代裸 {@code URI.create(url).toURL().openStream()}:后者走 JVM 默认超时(0 = 无限),
+ * 上游半开连接或挂起时会把 Tomcat 工作线程无限占用(管理端批量打包可同时挂多个)。
+ * 返回的流是流式的,适用于「服务端代理下载 OSS 文件转发给浏览器」这类不落盘场景。
+ *
+ * @param url 远程地址
+ * @param timeout 等待响应超时(连接建立 + 响应头);非法值钳制到 1 秒
+ * @throws IOException 非 2xx 响应或网络异常
+ * @throws InterruptedException 线程被中断
+ */
+ public static InputStream openStreamWithTimeout(String url, Duration timeout) throws IOException, InterruptedException {
+ Duration effective = (timeout == null || timeout.isZero() || timeout.isNegative())
+ ? Duration.ofSeconds(1)
+ : timeout;
+ HttpRequest request = HttpRequest.newBuilder(URI.create(url))
+ .timeout(effective)
+ .GET()
+ .build();
+ HttpResponse response = sharedHttpClient().send(request, HttpResponse.BodyHandlers.ofInputStream());
+ if (response.statusCode() / 100 != 2) {
+ try {
+ response.body().close();
+ } catch (Exception ignored) {
+ // 关闭失败不影响错误上报
+ }
+ throw new IOException("远程文件返回非 2xx: HTTP " + response.statusCode());
+ }
+ return response.body();
+ }
+
/** 按 readTimeout(毫秒)创建共享连接池工厂;非法值钳制到最小正数。 */
public static ClientHttpRequestFactory requestFactory(int readTimeoutMillis) {
return requestFactory(readTimeoutMillis, null);
diff --git a/backend-java/src/main/java/com/nanri/aiimage/config/SchedulingConfig.java b/backend-java/src/main/java/com/nanri/aiimage/config/SchedulingConfig.java
index 0d96983f..faead789 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/config/SchedulingConfig.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/config/SchedulingConfig.java
@@ -20,8 +20,15 @@ public class SchedulingConfig {
private static final ZoneId BUSINESS_ZONE = ZoneId.of("Asia/Shanghai");
+ /**
+ * 调度线程池:承载全站 30+ 个 @Scheduled(含 imagevideo 1s 派发、5s 轮询、结果文件 worker 15s 等高频任务)。
+ *
+ * 此前默认 4 线程:任一慢任务(如结果文件组装被内联执行时)都会把兜底类任务
+ * (StaleTaskRepair 心跳判死、陈旧扫描、历史清理)顺延,而兜底任务被顺延会直接放大线上故障面。
+ * 提到 16 并保持可配(aiimage.scheduling.pool-size)。
+ */
@Bean
- public TaskScheduler taskScheduler(@Value("${aiimage.scheduling.pool-size:4}") int poolSize) {
+ public TaskScheduler taskScheduler(@Value("${aiimage.scheduling.pool-size:16}") int poolSize) {
ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
scheduler.setPoolSize(Math.max(1, poolSize));
scheduler.setThreadNamePrefix("aiimage-scheduling-");
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/admin/support/AdminAuthSupport.java b/backend-java/src/main/java/com/nanri/aiimage/modules/admin/support/AdminAuthSupport.java
index 5a20c467..6ac2390f 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/admin/support/AdminAuthSupport.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/admin/support/AdminAuthSupport.java
@@ -143,7 +143,10 @@ public class AdminAuthSupport {
if (expectedToken.isBlank() || suppliedToken == null || suppliedToken.isBlank()) {
return false;
}
- return expectedToken.equals(suppliedToken.trim());
+ // 常量时间比较:逐字节 equals 可被计时侧信道逐位试探出内部令牌
+ return java.security.MessageDigest.isEqual(
+ expectedToken.getBytes(java.nio.charset.StandardCharsets.UTF_8),
+ suppliedToken.trim().getBytes(java.nio.charset.StandardCharsets.UTF_8));
}
/**
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/controller/AppearancePatentController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/controller/AppearancePatentController.java
index 897c89cc..491cdbe4 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/controller/AppearancePatentController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/controller/AppearancePatentController.java
@@ -34,6 +34,7 @@ import java.net.URI;
import java.nio.charset.StandardCharsets;
import com.nanri.aiimage.modules.task.model.dto.TaskProgressLightRequest;
import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
@RestController
@RequiredArgsConstructor
@@ -43,6 +44,7 @@ import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
public class AppearancePatentController {
private final AppearancePatentTaskService service;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
@PostMapping("/parse")
@Operation(summary = "解析 Excel 并创建任务", description = "解析上传后的 Excel 文件,提取 id、ASIN、国家、URL、标题等字段。返回给前端的数据只包含整数 id 和 n_1 行;n_2、n_3 等子行会保存在 OSS 解析载荷中,用于最终结果补齐。创建后的任务状态为 PENDING,不会自动推送 Python。")
@@ -99,7 +101,9 @@ public class AppearancePatentController {
@PostMapping("/tasks/progress/batch")
@Operation(summary = "批量查询任务进度", description = "前端只对活跃任务调用该接口,建议 6 秒一次。接口只返回轻量任务状态,不返回明细结果。")
public ApiResponse progress(@Valid @RequestBody AppearancePatentTaskBatchRequest request) {
- return ApiResponse.success(service.progressBatch(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(service.progressBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/light")
@@ -107,7 +111,8 @@ public class AppearancePatentController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileError/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody TaskProgressLightRequest request) {
- return ApiResponse.success(service.progressLight(request.getTaskIds()));
+ // 透传 userId:此前传 null = 不过滤归属,匿名遍历 taskId 即可读他人任务状态(2026-09 审查)
+ return ApiResponse.success(service.progressLight(request.getTaskIds(), request.getUserId()));
}
@PostMapping("/tasks/{taskId}/activate")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/model/dto/AppearancePatentTaskBatchRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/model/dto/AppearancePatentTaskBatchRequest.java
index 799b1f80..24e2cdaa 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/model/dto/AppearancePatentTaskBatchRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/model/dto/AppearancePatentTaskBatchRequest.java
@@ -12,4 +12,8 @@ public class AppearancePatentTaskBatchRequest {
@NotEmpty
@Schema(description = "需要查询进度的任务 ID 列表。前端只传正在轮询的活跃任务;后端会批量查询,避免每个任务单独请求。", example = "[3938,3939]", requiredMode = Schema.RequiredMode.REQUIRED)
private List taskIds;
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ @Schema(description = "当前用户 ID(可省略;传入时仅返回该用户的任务)", example = "1")
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTaskService.java
index 7a49b8b3..9787aedc 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTaskService.java
@@ -103,8 +103,8 @@ public class AppearancePatentTaskService {
public static final String MODULE_TYPE = "APPEARANCE_PATENT";
- public TaskProgressLightBatchVo progressLight(List taskIds) {
- return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
+ public TaskProgressLightBatchVo progressLight(List taskIds, Long userId) {
+ return taskProgressLightAssembler.assemble(MODULE_TYPE, userId, taskIds);
}
private static final String STATUS_PENDING = "PENDING";
private static final String STATUS_RUNNING = "RUNNING";
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/auth/service/AuthService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/auth/service/AuthService.java
index 802d9624..61f96398 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/auth/service/AuthService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/auth/service/AuthService.java
@@ -18,6 +18,9 @@ import org.springframework.http.ResponseCookie;
import org.springframework.stereotype.Service;
import java.time.Duration;
+import java.time.Instant;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
@Service
@RequiredArgsConstructor
@@ -30,6 +33,50 @@ public class AuthService {
private final PermissionMenuService permissionMenuService;
private final AuthProperties authProperties;
+ /**
+ * 登录失败计数与锁定(2026-09 全维度审查补):生产是公网域名,此前无任何失败限制
+ * 即可无限撞库。内存实现(双节点各自计数,防护效果减半但不引入新依赖),
+ * 达到阈值后锁定 15 分钟。
+ */
+ private static final int LOGIN_FAIL_LIMIT = 10;
+ private static final Duration LOGIN_LOCK_DURATION = Duration.ofMinutes(15);
+ private final Map loginFailures = new ConcurrentHashMap<>();
+
+ /** 失败计数;lockedUntil 非空表示已锁定。 */
+ private record FailRecord(int count, Instant lockedUntil) {
+ }
+
+ private void assertLoginNotLocked(String username) {
+ FailRecord record = loginFailures.get(username);
+ if (record != null && record.lockedUntil() != null && record.lockedUntil().isAfter(Instant.now())) {
+ long minutes = Duration.between(Instant.now(), record.lockedUntil()).toMinutes() + 1;
+ log.warn("[auth] 登录已锁定 username={} 剩余约 {} 分钟", username, minutes);
+ throw new BusinessException("登录失败次数过多,请 " + minutes + " 分钟后再试");
+ }
+ }
+
+ private void recordLoginFailure(String username) {
+ // 防无界增长:积累较多时清掉未锁定的过期记录(登录接口调用频次低)
+ if (loginFailures.size() > 1000) {
+ Instant now = Instant.now();
+ loginFailures.entrySet().removeIf(e -> e.getValue().lockedUntil() == null
+ || e.getValue().lockedUntil().isBefore(now));
+ }
+ loginFailures.compute(username, (key, old) -> {
+ int count = (old == null ? 0 : old.count()) + 1;
+ Instant lockedUntil = count >= LOGIN_FAIL_LIMIT ? Instant.now().plus(LOGIN_LOCK_DURATION) : null;
+ if (lockedUntil != null) {
+ log.warn("[auth] 登录失败达阈值,锁定 username={} count={} minutes={}",
+ username, count, LOGIN_LOCK_DURATION.toMinutes());
+ }
+ return new FailRecord(count, lockedUntil);
+ });
+ }
+
+ private void clearLoginFailures(String username) {
+ loginFailures.remove(username);
+ }
+
public LoginResultVo login(LoginRequest request) {
String username = trim(request.getUsername());
String password = request.getPassword() == null ? "" : request.getPassword();
@@ -41,12 +88,16 @@ public class AuthService {
throw new BusinessException("缺少设备ID,请在桌面端打开");
}
+ assertLoginNotLocked(username);
+
LoginUserEntity user = loginUserMapper.selectOne(new LambdaQueryWrapper()
.eq(LoginUserEntity::getUsername, username)
.last("LIMIT 1"));
if (user == null || !passwordEncoder.matches(password, user.getPasswordHash())) {
+ recordLoginFailure(username);
throw new BusinessException("用户名或密码错误");
}
+ clearLoginFailures(username);
boolean isAdmin = user.getIsAdmin() != null && user.getIsAdmin() == 1;
// 单设备登录:登录成功即把账号绑定到当前设备(last-login-wins),原设备下一次请求被顶下线
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskProgressCacheService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskProgressCacheService.java
index 03a652d8..f5d66391 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskProgressCacheService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskProgressCacheService.java
@@ -4,11 +4,13 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.config.BrandProgressProperties;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
+import org.springframework.data.redis.core.script.DefaultRedisScript;
import org.springframework.stereotype.Service;
import java.time.Duration;
import java.time.Instant;
import java.util.LinkedHashMap;
+import java.util.List;
import java.util.Map;
@Service
@@ -21,6 +23,9 @@ public class BrandTaskProgressCacheService {
public static final String PHASE_FAILED = "failed";
private static final Duration FINALIZE_LOCK_TTL = Duration.ofMinutes(10);
+ /** 本实例(进程)的锁持有者标识:释放锁时用它校验"锁还是我的"。 */
+ private final String lockOwnerToken = java.util.UUID.randomUUID().toString();
+
private final StringRedisTemplate stringRedisTemplate;
private final BrandProgressProperties brandProgressProperties;
@SuppressWarnings("unused")
@@ -129,8 +134,10 @@ public class BrandTaskProgressCacheService {
public boolean acquireFinalizeLock(Long taskId) {
try {
+ // value 用持有者 token(而非时间戳):释放时需要它来校验"锁还是我的",
+ // 否则锁因 TTL 到期被他方持有后,本线程的裸 delete 会误删他人的锁
Boolean ok = stringRedisTemplate.opsForValue()
- .setIfAbsent(buildFinalizeLockKey(taskId), String.valueOf(Instant.now().toEpochMilli()), FINALIZE_LOCK_TTL);
+ .setIfAbsent(buildFinalizeLockKey(taskId), lockOwnerToken, FINALIZE_LOCK_TTL);
return Boolean.TRUE.equals(ok);
} catch (Exception ex) {
log.warn("[brand-progress-cache] acquire finalize lock degraded taskId={} msg={}", taskId, ex.getMessage());
@@ -140,7 +147,12 @@ public class BrandTaskProgressCacheService {
public void releaseFinalizeLock(Long taskId) {
try {
- stringRedisTemplate.delete(buildFinalizeLockKey(taskId));
+ // Lua 原子校验后删除(2026-09 全维度审查):此前是裸 delete,锁已过期(TTL 到期、
+ // 他人已持有)时会把别人的锁删掉,导致同一任务被两个线程同时收尾。
+ stringRedisTemplate.execute(new DefaultRedisScript<>(
+ "if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end",
+ Long.class),
+ List.of(buildFinalizeLockKey(taskId)), lockOwnerToken);
} catch (Exception ex) {
log.warn("[brand-progress-cache] release finalize lock degraded taskId={} msg={}", taskId, ex.getMessage());
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskService.java
index 6ced901c..10012d17 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskService.java
@@ -9,6 +9,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.common.service.DistributedJobLockService;
import com.nanri.aiimage.config.BrandProgressProperties;
+import com.nanri.aiimage.config.HttpClientPool;
import com.nanri.aiimage.config.StorageProperties;
import com.nanri.aiimage.modules.brand.mapper.BrandCrawlTaskMapper;
import com.nanri.aiimage.modules.brand.model.dto.BrandCrawlResultFileDto;
@@ -53,6 +54,8 @@ import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.net.URI;
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
import java.nio.file.Files;
import java.time.Duration;
import java.time.Instant;
@@ -82,6 +85,21 @@ public class BrandTaskService {
private static final DateTimeFormatter DATETIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm");
private static final Duration STALE_CHECK_LOCK_TTL = Duration.ofMinutes(10);
+ /** 源文件下载请求超时:比普通 API 调用宽松(源文件可能几十 MB),但必须有上限。 */
+ private static final Duration SOURCE_DOWNLOAD_TIMEOUT = Duration.ofMinutes(5);
+
+ /** SSRF 防护:禁止请求云元数据与环回地址(正常源文件都在自家 OSS/MinIO 域名上)。 */
+ private static final java.util.Set BLOCKED_SOURCE_HOSTS = java.util.Set.of(
+ "169.254.169.254", "metadata.google.internal", "metadata", "localhost",
+ "127.0.0.1", "0.0.0.0", "::1", "[::1]");
+
+ private static boolean isBlockedSourceHost(String host) {
+ if (host == null || host.isBlank()) {
+ return true;
+ }
+ String normalized = host.trim().toLowerCase(java.util.Locale.ROOT);
+ return BLOCKED_SOURCE_HOSTS.contains(normalized) || normalized.endsWith(".localhost");
+ }
private static final long RESULT_SUBMIT_WAIT_MILLIS = 5 * 60 * 1000L;
private static final String STATUS_PENDING = "pending";
private static final String STATUS_RUNNING = "running";
@@ -92,8 +110,9 @@ public class BrandTaskService {
/** 结果文件并发上传数:OSS/MinIO 上传互不依赖,3 并发平衡收益与内存占用。 */
private static final int RESULT_UPLOAD_CONCURRENCY = 3;
- private final ExecutorService resultUploadExecutor = Executors.newFixedThreadPool(
- RESULT_UPLOAD_CONCURRENCY, namedThreadFactory("brand-result-upload"));
+ // 有界队列线程池:newFixedThreadPool 用的是无界队列,任务堆积时不会拒绝、会把内存吃满
+ private final ExecutorService resultUploadExecutor = com.nanri.aiimage.common.util.ThreadPools
+ .boundedFixed("brand-result-upload", RESULT_UPLOAD_CONCURRENCY);
@PreDestroy
void shutdownResultUploadExecutor() {
@@ -812,21 +831,49 @@ public class BrandTaskService {
}
private File downloadSourceFile(String fileUrl) {
+ URI uri;
+ try {
+ uri = URI.create(fileUrl);
+ } catch (Exception ex) {
+ log.warn("[brand] 源文件地址非法 fileUrl={} err={}", fileUrl, ex.getMessage());
+ throw new BusinessException("下载源文件失败: 地址非法");
+ }
+ // SSRF 防护(2026-09 全维度审查):fileUrl 来自请求体,若不校验则匿名调用方可让服务端
+ // 请求内网/云元数据地址(169.254.169.254 等)。正常源文件都落在自家 OSS/MinIO 上。
+ if (isBlockedSourceHost(uri.getHost())) {
+ log.warn("[brand] 拒绝下载疑似 SSRF 的源文件地址 host={} fileUrl={}", uri.getHost(), fileUrl);
+ throw new BusinessException("源文件地址不合法");
+ }
+ String filename = FileUtil.getName(uri.getPath());
+ if (filename == null || filename.isBlank()) {
+ filename = "brand-source.xlsx";
+ }
+ String suffix = FileUtil.extName(filename);
+ File downloadDir = FileUtil.mkdir(FileUtil.file(storageProperties.getLocalTempDir(), "brand-source-download"));
try {
- URI uri = URI.create(fileUrl);
- String filename = FileUtil.getName(uri.getPath());
- if (filename == null || filename.isBlank()) {
- filename = "brand-source.xlsx";
- }
- String suffix = FileUtil.extName(filename);
- File downloadDir = FileUtil.mkdir(FileUtil.file(storageProperties.getLocalTempDir(), "brand-source-download"));
File tempFile = Files.createTempFile(downloadDir.toPath(), "brand_", suffix.isBlank() ? "" : "." + suffix).toFile();
- try (InputStream inputStream = uri.toURL().openStream()) {
+ // 走统一连接池 + 显式超时。此前 uri.toURL().openStream() 无任何超时(JVM 默认 0 = 无限),
+ // 上游半开连接会把 Tomcat 工作线程无限占用;超时值比普通 API 调用宽松(源文件可能几十 MB)
+ HttpRequest downloadRequest = HttpRequest.newBuilder(uri)
+ .timeout(SOURCE_DOWNLOAD_TIMEOUT)
+ .GET()
+ .build();
+ HttpResponse response = HttpClientPool.sharedHttpClient()
+ .send(downloadRequest, HttpResponse.BodyHandlers.ofInputStream());
+ if (response.statusCode() / 100 != 2) {
+ log.warn("[brand] 下载源文件返回非 2xx fileUrl={} status={}", fileUrl, response.statusCode());
+ throw new BusinessException("下载源文件失败: HTTP " + response.statusCode());
+ }
+ try (InputStream inputStream = response.body()) {
FileUtil.writeFromStream(inputStream, tempFile);
}
return tempFile;
+ } catch (BusinessException ex) {
+ throw ex;
} catch (Exception ex) {
- throw new BusinessException("下载源文件失败");
+ // 此前该 catch 既不记日志也不带 cause,线上无法定位
+ log.warn("[brand] 下载源文件失败 fileUrl={} err={}", fileUrl, ex.getMessage(), ex);
+ throw new BusinessException("下载源文件失败", ex);
}
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskStorageService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskStorageService.java
index 2ff1fe53..1931b02f 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskStorageService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/brand/service/BrandTaskStorageService.java
@@ -466,11 +466,8 @@ public class BrandTaskStorageService {
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
byte[] bytes = digest.digest(normalize(value).getBytes(StandardCharsets.UTF_8));
- StringBuilder sb = new StringBuilder(bytes.length * 2);
- for (byte b : bytes) {
- sb.append(String.format("%02x", b));
- }
- return sb.toString();
+ // HexFormat 替代逐字节 String.format("%02x"):后者每次走 32 次 Formatter(热路径按行调用)
+ return java.util.HexFormat.of().formatHex(bytes);
} catch (Exception ex) {
throw new IllegalStateException("failed to hash brand scope", ex);
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/controller/CollectDataController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/controller/CollectDataController.java
index b3a2ca3c..530c5cc8 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/controller/CollectDataController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/controller/CollectDataController.java
@@ -29,6 +29,7 @@ import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import com.nanri.aiimage.modules.task.model.dto.TaskProgressLightRequest;
import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
@RestController
@RequiredArgsConstructor
@@ -37,6 +38,7 @@ import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
public class CollectDataController {
private final CollectDataService service;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
@PostMapping("/parse")
@Operation(summary = "解析 Excel 并创建任务", description = "解析上传后的 Excel 文件,按行入库到 biz_collect_data_item,并保存任务筛选条件。任务初始状态为 PENDING。")
@@ -112,7 +114,9 @@ public class CollectDataController {
@PostMapping("/tasks/progress/batch")
@Operation(summary = "批量查询任务进度")
public ApiResponse progressBatch(@Valid @RequestBody CollectDataTaskBatchRequest request) {
- return ApiResponse.success(service.progressBatch(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(service.progressBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/light")
@@ -120,7 +124,8 @@ public class CollectDataController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileError/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody TaskProgressLightRequest request) {
- return ApiResponse.success(service.progressLight(request.getTaskIds()));
+ // 透传 userId:此前传 null = 不过滤归属,匿名遍历 taskId 即可读他人任务状态(2026-09 审查)
+ return ApiResponse.success(service.progressLight(request.getTaskIds(), request.getUserId()));
}
@DeleteMapping("/tasks/{taskId}")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/model/dto/CollectDataTaskBatchRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/model/dto/CollectDataTaskBatchRequest.java
index a51081da..d92ea216 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/model/dto/CollectDataTaskBatchRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/model/dto/CollectDataTaskBatchRequest.java
@@ -12,4 +12,8 @@ public class CollectDataTaskBatchRequest {
@NotEmpty
@Schema(description = "任务 ID 列表", requiredMode = Schema.RequiredMode.REQUIRED)
private List taskIds;
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ @Schema(description = "当前用户 ID(可省略;传入时仅返回该用户的任务)", example = "1")
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/service/CollectDataService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/service/CollectDataService.java
index c06ca0a2..49e1e70c 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/service/CollectDataService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/service/CollectDataService.java
@@ -4,9 +4,11 @@ import cn.hutool.core.util.IdUtil;
import cn.hutool.core.io.FileUtil;
import cn.hutool.crypto.digest.DigestUtil;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.nanri.aiimage.common.exception.BusinessCodes;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.common.util.ExcelStreamReader;
import com.nanri.aiimage.modules.collectdata.mapper.CollectDataCountryPrefMapper;
@@ -104,8 +106,8 @@ public class CollectDataService {
*/
private static final String LEGACY_MODULE_TYPE = "collectdata";
- public TaskProgressLightBatchVo progressLight(List taskIds) {
- return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
+ public TaskProgressLightBatchVo progressLight(List taskIds, Long userId) {
+ return taskProgressLightAssembler.assemble(MODULE_TYPE, userId, taskIds);
}
public static final int DEFAULT_PAGE_SIZE = 50;
@@ -379,34 +381,45 @@ public class CollectDataService {
@Transactional
public void activateTask(Long taskId, Long userId) {
FileTaskEntity task = requireTask(taskId, userId);
- if (STATUS_SUCCESS.equals(task.getStatus()) || STATUS_FAILED.equals(task.getStatus())) {
+ // 条件更新:上面的「已结束」判断与写入之间存在窗口(TOCTOU),期间 /fail 可能已把任务
+ // 标为 FAILED —— 整行 updateById 会把它复活成 RUNNING(前端显示"执行中"但无人推进)
+ int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper()
+ .eq(FileTaskEntity::getId, task.getId())
+ .notIn(FileTaskEntity::getStatus, STATUS_SUCCESS, STATUS_FAILED)
+ .set(FileTaskEntity::getStatus, STATUS_RUNNING)
+ .set(FileTaskEntity::getUpdatedAt, LocalDateTime.now()));
+ if (updated == 0) {
throw new BusinessException("任务已结束");
}
- task.setStatus(STATUS_RUNNING);
- task.setUpdatedAt(LocalDateTime.now());
- fileTaskMapper.updateById(task);
}
@Transactional
public void failTask(Long taskId, Long userId, String error) {
FileTaskEntity task = requireTask(taskId, userId);
- if (STATUS_SUCCESS.equals(task.getStatus()) || STATUS_FAILED.equals(task.getStatus())) {
- return;
- }
String message = firstNonBlank(error, "collect-data task dispatch failed");
FileResultEntity result = ensureTaskResult(task);
CollectDataStats stats = loadStats(task);
- result.setSuccess(0);
- result.setErrorMessage(message);
- result.setRowCount(stats.finalRowCount);
- fileResultMapper.updateById(result);
- task.setStatus(STATUS_FAILED);
- task.setErrorMessage(message);
- task.setFailedFileCount(1);
- task.setUpdatedAt(LocalDateTime.now());
- task.setFinishedAt(LocalDateTime.now());
- persistStats(task, stats);
- fileTaskMapper.updateById(task);
+ boolean alreadyTerminal = STATUS_SUCCESS.equals(task.getStatus()) || STATUS_FAILED.equals(task.getStatus());
+ if (!alreadyTerminal) {
+ result.setSuccess(0);
+ result.setErrorMessage(message);
+ result.setRowCount(stats.finalRowCount);
+ fileResultMapper.updateById(result);
+ persistStats(task, stats);
+ }
+ // 条件更新:客户端报错(/fail)与结果文件组装完成(processResultFileJob 写 SUCCESS)
+ // 可能并发 —— 无条件 updateById 会把已生成的 SUCCESS 覆盖成 FAILED(用户拿不到下载)或反之
+ int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper()
+ .eq(FileTaskEntity::getId, task.getId())
+ .notIn(FileTaskEntity::getStatus, STATUS_SUCCESS, STATUS_FAILED)
+ .set(FileTaskEntity::getStatus, STATUS_FAILED)
+ .set(FileTaskEntity::getErrorMessage, message)
+ .set(FileTaskEntity::getFailedFileCount, 1)
+ .set(FileTaskEntity::getUpdatedAt, LocalDateTime.now())
+ .set(FileTaskEntity::getFinishedAt, LocalDateTime.now()));
+ if (updated == 0) {
+ log.info("[collect-data] failTask 跳过写入:任务已是终态 taskId={} status={}", taskId, task.getStatus());
+ }
}
@Transactional
@@ -954,13 +967,21 @@ public class CollectDataService {
// 复用同一份 stats 更新 finalRowCount 后再持久化,避免重复 loadStats 丢失 summaries。
stats.finalRowCount = (int) finalRowCount;
persistStats(task, stats);
- task.setStatus(STATUS_SUCCESS);
- task.setSuccessFileCount(1);
- task.setFailedFileCount(0);
- task.setErrorMessage(null);
- task.setUpdatedAt(LocalDateTime.now());
- task.setFinishedAt(LocalDateTime.now());
- fileTaskMapper.updateById(task);
+ // 条件更新:任务可能已被 /fail 标为 FAILED(客户端报错与结果文件组装并发)——
+ // 无条件 updateById 会把 FAILED 覆盖成 SUCCESS(错误信息被清、用户看到「假成功」)。
+ // 结果文件已生成,故失败态下仍保留文件,只是不覆盖状态。
+ int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper()
+ .eq(FileTaskEntity::getId, task.getId())
+ .ne(FileTaskEntity::getStatus, STATUS_FAILED)
+ .set(FileTaskEntity::getStatus, STATUS_SUCCESS)
+ .set(FileTaskEntity::getSuccessFileCount, 1)
+ .set(FileTaskEntity::getFailedFileCount, 0)
+ .set(FileTaskEntity::getErrorMessage, null)
+ .set(FileTaskEntity::getUpdatedAt, LocalDateTime.now())
+ .set(FileTaskEntity::getFinishedAt, LocalDateTime.now()));
+ if (updated == 0) {
+ log.warn("[collect-data] 结果文件已生成但任务已是 FAILED,保留失败态不覆盖 taskId={}", task.getId());
+ }
} finally {
FileUtil.del(xlsx);
}
@@ -1085,7 +1106,8 @@ public class CollectDataService {
private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) {
TaskDistributedLockService.LockHandle lockHandle = taskDistributedLockService.acquire(MODULE_TYPE, taskId, TASK_LOCK_WAIT_MILLIS);
if (lockHandle == null) {
- throw new BusinessException(40901, "任务正在处理,请稍后再试");
+ log.warn("[collect-data] 任务锁竞争,拒绝本次提交 taskId={} waitMillis={}", taskId, TASK_LOCK_WAIT_MILLIS);
+ throw new BusinessException(BusinessCodes.TASK_BUSY, "任务正在处理,请稍后再试");
}
return lockHandle;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodec.java b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodec.java
index 627d0ece..2196a255 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodec.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultDetailCodec.java
@@ -148,11 +148,8 @@ public class CollectDataResultDetailCodec {
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
byte[] bytes = digest.digest((value == null ? "" : value).getBytes(StandardCharsets.UTF_8));
- StringBuilder sb = new StringBuilder(bytes.length * 2);
- for (byte b : bytes) {
- sb.append(String.format("%02x", b));
- }
- return sb.toString();
+ // HexFormat 替代逐字节 String.format("%02x"):后者每次走 32 次 Formatter(热路径按行调用)
+ return java.util.HexFormat.of().formatHex(bytes);
} catch (Exception ex) {
throw new IllegalStateException("chunk detail ref hash failed", ex);
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultItemBatchWriter.java b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultItemBatchWriter.java
index 9c3bdee2..f72d63f9 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultItemBatchWriter.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/collectdata/util/CollectDataResultItemBatchWriter.java
@@ -1,6 +1,7 @@
package com.nanri.aiimage.modules.collectdata.util;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo;
import com.nanri.aiimage.modules.task.mapper.TaskResultItemMapper;
import com.nanri.aiimage.modules.task.model.entity.TaskResultItemEntity;
@@ -133,12 +134,14 @@ public class CollectDataResultItemBatchWriter {
}
int written = 0;
+ int failedBatches = 0;
for (int from = 0; from < toUpsert.size(); from += batchSize) {
int to = Math.min(from + batchSize, toUpsert.size());
List batch = toUpsert.subList(from, to);
try {
written += taskResultItemMapper.upsertBatch(batch);
} catch (RuntimeException ex) {
+ failedBatches++;
log.warn("[collect-data] upsert result item batch failed, skip batch {}..{} taskId={}",
from, to, taskId, ex);
// 失败批的新行未落库,从增量计数中扣除,避免任务内累计虚高;
@@ -150,6 +153,16 @@ public class CollectDataResultItemBatchWriter {
}
}
}
+ if (failedBatches > 0) {
+ // biz_task_result_item 是结果 Excel「明细」sheet 的唯一数据源:静默跳批会让任务
+ // 以 SUCCESS 收尾但明细缺行,且与「结果文件」sheet 的汇总数量对不上(假成功)。
+ // 抛出让本次 chunk 提交明确失败:worker(search_spider)识别 success=false 后会重试,
+ // 重提按 payload_hash 幂等(已落库行跳过、未落库行补插),最终收敛为完整数据。
+ log.error("[collect-data] 明细写入存在失败批次,拒绝本次提交 taskId={} failedBatches={} totalBatches={}",
+ taskId, failedBatches, (toUpsert.size() + batchSize - 1) / batchSize);
+ throw new BusinessException("采集结果明细写入失败,请稍后重试(失败批次 "
+ + failedBatches + "/" + ((toUpsert.size() + batchSize - 1) / batchSize) + ")");
+ }
return new UpsertCounts(written, skipped, newlyInserted);
}
@@ -157,11 +170,8 @@ public class CollectDataResultItemBatchWriter {
try {
java.security.MessageDigest digest = java.security.MessageDigest.getInstance("SHA-256");
byte[] bytes = digest.digest((value == null ? "" : value).getBytes(java.nio.charset.StandardCharsets.UTF_8));
- StringBuilder sb = new StringBuilder(bytes.length * 2);
- for (byte b : bytes) {
- sb.append(String.format("%02x", b));
- }
- return sb.toString();
+ // HexFormat 替代逐字节 String.format("%02x"):后者每次走 32 次 Formatter(热路径按行调用)
+ return java.util.HexFormat.of().formatHex(bytes);
} catch (Exception ex) {
throw new IllegalStateException("结果明细 hash 计算失败", ex);
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/convert/service/ConvertRunService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/convert/service/ConvertRunService.java
index 37061768..e08d4417 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/convert/service/ConvertRunService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/convert/service/ConvertRunService.java
@@ -278,12 +278,14 @@ public class ConvertRunService {
File outputDir = FileUtil.mkdir(FileUtil.file(storageProperties.getLocalTempDir(), "convert-result"));
List generatedFiles = new ArrayList<>();
Map writers = new LinkedHashMap<>();
- for (String outputFilename : outputFilenames) {
- File outputFile = buildNamedOutputFile(outputDir, outputFilename);
- generatedFiles.add(new GeneratedConvertFile(outputFilename, outputFile));
- writers.put(outputFilename, Files.newBufferedWriter(outputFile.toPath(), StandardCharsets.UTF_8));
- }
try {
+ // 创建循环纳入 try:第 N 个 writer 创建失败时,前 N-1 个已打开的句柄会因 finally
+ // 尚未生效而泄漏(同族 SplitRunService.SplitChunkWriter 已用 try/finally 处理)
+ for (String outputFilename : outputFilenames) {
+ File outputFile = buildNamedOutputFile(outputDir, outputFilename);
+ generatedFiles.add(new GeneratedConvertFile(outputFilename, outputFile));
+ writers.put(outputFilename, Files.newBufferedWriter(outputFile.toPath(), StandardCharsets.UTF_8));
+ }
streamTxtRowsToOutputs(inputFile, templateEntity, writers);
} finally {
IOException closeException = null;
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/dedupe/service/DedupeRunService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/dedupe/service/DedupeRunService.java
index e3999302..9bf9474d 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/dedupe/service/DedupeRunService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/dedupe/service/DedupeRunService.java
@@ -75,12 +75,9 @@ public class DedupeRunService {
private final Map runningTaskCountMap = new ConcurrentHashMap<>();
private final Semaphore fileParallelSemaphore = new Semaphore(MAX_PARALLEL_FILES);
/** 共享有界文件执行器;不再为请求中的每个文件创建虚拟线程。 */
- private final ExecutorService fileExecutor = Executors.newFixedThreadPool(
- MAX_PARALLEL_FILES, runnable -> {
- Thread thread = new Thread(runnable, "dedupe-file-worker");
- thread.setDaemon(true);
- return thread;
- });
+ // 队列改为有界:newFixedThreadPool 的无界队列在文件堆积时不会拒绝、会把内存吃满
+ private final ExecutorService fileExecutor = com.nanri.aiimage.common.util.ThreadPools
+ .boundedFixed("dedupe-file-worker", MAX_PARALLEL_FILES);
/**
* 提交去重任务:立即返回进度快照(runId),异步执行 流式读取 + 多文件并行 + 结果上传。
@@ -171,6 +168,8 @@ public class DedupeRunService {
boolean folderMode = request.getArchiveName() != null && !request.getArchiveName().isBlank();
Map archiveEntries = new ConcurrentHashMap<>();
List outcomeItems = new ArrayList<>();
+ // 是否有 worker 异常结束:为 true 时收尾必须标失败,绝不能标 SUCCESS(用户会看到成功但结果缺文件)
+ boolean aborted = false;
try {
AtomicInteger processedCount = new AtomicInteger(0);
@@ -192,14 +191,26 @@ public class DedupeRunService {
}
}));
}
+ // 逐个等待:processFile 内部已把单文件失败记入 outcomeItems,不会往外抛;
+ // 这里能捕到的是线程级异常(Error/线程池拒绝),必须记下来而不能当成功。
for (Future> worker : workers) {
- worker.get();
+ try {
+ worker.get();
+ } catch (Exception workerEx) {
+ aborted = true;
+ log.error("dedupe run worker aborted runId={} error", runId, workerEx);
+ }
}
} catch (Exception ex) {
+ aborted = true;
log.error("dedupe run async aborted runId={} error", runId, ex);
}
try {
+ if (aborted) {
+ // 有文件没能正常处理完就序列化/标成功,会出现「任务成功但结果缺文件」
+ throw new IllegalStateException("部分文件处理线程异常结束,结果不完整");
+ }
if (folderMode && !archiveEntries.isEmpty()) {
DedupeResultItemVo zipItem = buildFolderZipResult(request, archiveEntries, task);
synchronized (progress) {
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/dedupe/service/DedupeTotalDataService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/dedupe/service/DedupeTotalDataService.java
index db2b0cb4..5edda316 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/dedupe/service/DedupeTotalDataService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/dedupe/service/DedupeTotalDataService.java
@@ -17,6 +17,7 @@ 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 lombok.extern.slf4j.Slf4j;
import org.apache.poi.ss.usermodel.Cell;
import org.apache.poi.ss.usermodel.DataFormatter;
import org.apache.poi.ss.usermodel.Row;
@@ -62,6 +63,7 @@ import java.util.zip.ZipOutputStream;
@Service
@RequiredArgsConstructor
+@Slf4j
public class DedupeTotalDataService {
private static final int COMPARE_BATCH_SIZE = 5000;
@@ -579,6 +581,8 @@ public class DedupeTotalDataService {
} catch (Exception e) {
progress.setStatus("failed");
progress.setErrorMessage(e instanceof BusinessException ? e.getMessage() : "导入 Excel 失败");
+ // 此前该 catch 既不记日志也不带 cause,线上只有一句「导入 Excel 失败」,无栈无定位线索
+ log.error("[dedupe-total-data] 异步导入失败 filename={} userId={}", filename, uploaderUserId, e);
} finally {
deleteQuietly(tempFile);
importCompletedAtMap.put(importId, System.currentTimeMillis());
@@ -648,7 +652,9 @@ public class DedupeTotalDataService {
} catch (BusinessException e) {
throw e;
} catch (Exception e) {
- throw new BusinessException("删除 Excel 匹配数据失败");
+ // 此前丢弃原始异常:排查时只能看到一句「删除 Excel 匹配数据失败」
+ log.error("[dedupe-total-data] 删除 Excel 匹配数据失败", e);
+ throw new BusinessException("删除 Excel 匹配数据失败", e);
}
}
@@ -946,7 +952,9 @@ public class DedupeTotalDataService {
} catch (BusinessException e) {
throw e;
} catch (Exception e) {
- throw new BusinessException("删除 Excel 匹配数据失败");
+ // 此前丢弃原始异常:排查时只能看到一句「删除 Excel 匹配数据失败」
+ log.error("[dedupe-total-data] 删除 Excel 匹配数据失败", e);
+ throw new BusinessException("删除 Excel 匹配数据失败", e);
}
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/controller/DeleteBrandRunController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/controller/DeleteBrandRunController.java
index 2baaea4b..a6015128 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/controller/DeleteBrandRunController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/controller/DeleteBrandRunController.java
@@ -37,6 +37,7 @@ import java.net.URI;
import java.nio.charset.StandardCharsets;
import com.nanri.aiimage.modules.task.model.dto.TaskProgressLightRequest;
import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
@RestController
@@ -47,6 +48,7 @@ import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
public class DeleteBrandRunController {
private final DeleteBrandRunService deleteBrandRunService;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
@PostMapping("/run")
@Operation(summary = "执行删除品牌解析", description = "读取上传的删除品牌 Excel,按国家分组解析并在每个国家内按 ASIN 去重,返回完整 payload。")
@@ -85,13 +87,17 @@ public class DeleteBrandRunController {
@PostMapping("/tasks/batch")
@Operation(summary = "批量获取删除品牌任务详情", description = "合并多个 taskId 的详情查询,减少轮询期对 DB/Redis 的压力。")
public ApiResponse getTasksBatch(@Valid @RequestBody DeleteBrandTaskBatchRequest request) {
- return ApiResponse.success(deleteBrandRunService.getTaskDetails(request.getTaskIds()));
+ // 归属过滤:/tasks/batch 同样按 userId 过滤(此前无校验,遍历 taskId 即可读他人任务详情)
+ return ApiResponse.success(deleteBrandRunService.getTaskDetails(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/batch")
@Operation(summary = "批量获取删除品牌任务进度摘要", description = "仅返回任务基础状态和行进度摘要,用于前端轮询降载。")
public ApiResponse getTaskProgressBatch(@Valid @RequestBody DeleteBrandTaskBatchRequest request) {
- return ApiResponse.success(deleteBrandRunService.getTaskProgress(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(deleteBrandRunService.getTaskProgress(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/light")
@@ -99,7 +105,8 @@ public class DeleteBrandRunController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileError/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody TaskProgressLightRequest request) {
- return ApiResponse.success(deleteBrandRunService.progressLight(request.getTaskIds()));
+ // 透传 userId:此前传 null = 不过滤归属,匿名遍历 taskId 即可读他人任务状态(2026-09 审查)
+ return ApiResponse.success(deleteBrandRunService.progressLight(request.getTaskIds(), request.getUserId()));
}
@GetMapping("/tasks/{taskId}/deletion-status")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/model/dto/DeleteBrandTaskBatchRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/model/dto/DeleteBrandTaskBatchRequest.java
index e6f30186..757ad9d3 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/model/dto/DeleteBrandTaskBatchRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/model/dto/DeleteBrandTaskBatchRequest.java
@@ -14,4 +14,8 @@ public class DeleteBrandTaskBatchRequest {
@NotEmpty(message = "taskIds 不能为空")
@Schema(description = "任务ID列表")
private List taskIds = new ArrayList<>();
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ @Schema(description = "当前用户 ID(可省略;传入时仅返回该用户的任务)", example = "1")
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandRunService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandRunService.java
index ddb45d43..9ed70ad2 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandRunService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandRunService.java
@@ -82,8 +82,8 @@ public class DeleteBrandRunService {
/** 降级组装时缺失分片对应的结果状态(客户端未回传该行)。 */
static final String MISSING_CHUNK_STATUS = "未回传";
- public TaskProgressLightBatchVo progressLight(List taskIds) {
- return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
+ public TaskProgressLightBatchVo progressLight(List taskIds, Long userId) {
+ return taskProgressLightAssembler.assemble(MODULE_TYPE, userId, taskIds);
}
private static final String CONTENT_TYPE_XLSX = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet";
private static final Duration TASK_LOCK_TTL = TaskDistributedLockService.DEFAULT_LOCK_TTL;
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandTaskStorageService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandTaskStorageService.java
index 40a59cfb..dc86fbd0 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandTaskStorageService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/deletebrand/service/DeleteBrandTaskStorageService.java
@@ -393,11 +393,8 @@ public class DeleteBrandTaskStorageService {
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
byte[] bytes = digest.digest(normalizeScopeKey(value).getBytes(StandardCharsets.UTF_8));
- StringBuilder sb = new StringBuilder(bytes.length * 2);
- for (byte b : bytes) {
- sb.append(String.format("%02x", b));
- }
- return sb.toString();
+ // HexFormat 替代逐字节 String.format("%02x"):后者每次走 32 次 Formatter(热路径按行调用)
+ return java.util.HexFormat.of().formatHex(bytes);
} catch (Exception ex) {
throw new IllegalStateException("failed to hash delete-brand scope", ex);
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/imagevideo/service/ImageVideoArchiveService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/imagevideo/service/ImageVideoArchiveService.java
index 695ed386..a762dbd7 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/imagevideo/service/ImageVideoArchiveService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/imagevideo/service/ImageVideoArchiveService.java
@@ -153,7 +153,12 @@ public class ImageVideoArchiveService {
Object result = readJsonValue(task.getResultJson());
enrichCompletedTask(task, result);
task.setUpdatedAt(LocalDateTime.now());
- taskMapper.updateById(task);
+ // 条件更新(仅 archive_status 仍为 NULL 时):整行 updateById 会用读取快照覆盖并发节点
+ // 刚写入的 ARCHIVING —— 双节点 10s 周期归档时,B 节点会把 A 的认领覆盖回可认领状态,
+ // 随后 B 的 CAS 也能成功,导致同一视频被下载/上传两遍(先上传的对象成孤儿)
+ taskMapper.update(task, new LambdaUpdateWrapper()
+ .eq(ImageVideoAsyncTaskEntity::getId, task.getId())
+ .isNull(ImageVideoAsyncTaskEntity::getArchiveStatus));
}
private void archiveTask(Long taskId) {
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/imagevideo/service/ImageVideoAsyncTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/imagevideo/service/ImageVideoAsyncTaskService.java
index be6c04c8..5329dcb2 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/imagevideo/service/ImageVideoAsyncTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/imagevideo/service/ImageVideoAsyncTaskService.java
@@ -1,6 +1,7 @@
package com.nanri.aiimage.modules.imagevideo.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.common.exception.BusinessException;
@@ -180,6 +181,54 @@ public class ImageVideoAsyncTaskService {
}
}
+ /** owner 节点失联判定阈值(分钟):任务超过此时长未更新即认为原 owner 已下线。 */
+ private static final long OWNER_TAKEOVER_TIMEOUT_MINUTES = 15;
+
+ /**
+ * 死节点接管(2026-09 全维度审查补):owner 绑定的任务此前只有「owner 重启时恢复自己的」
+ * 一条路径,owner 节点崩溃/被摘除后其任务永久停在「生成中」
+ * ({@code StaleTaskRepairService} 只覆盖 biz_file_task / brand_crawl_tasks,不含本表)。
+ *
+ * 分两类处理:
+ *
+ * - WAITING/POLLING —— 只是等上游结果,重放安全 → 清空 owner 交回可领取队列;
+ * - RUNNING —— 重放会重复调用上游生成(有费用),故超时后明确标 FAILED 让用户重试,
+ * 而不是悄悄重新执行。
+ *
+ */
+ @Scheduled(fixedDelayString = "${aiimage.image-video.owner-takeover-delay-ms:60000}")
+ public void takeoverAbandonedTasks() {
+ DistributedJobLockService.LockHandle jobLock =
+ distributedJobLockService.tryLock("image-video:owner-takeover", java.time.Duration.ofSeconds(30));
+ if (jobLock == null) {
+ return;
+ }
+ try (jobLock) {
+ LocalDateTime cutoff = LocalDateTime.now().minusMinutes(OWNER_TAKEOVER_TIMEOUT_MINUTES);
+ int requeued = taskMapper.update(null, new LambdaUpdateWrapper()
+ .in(ImageVideoAsyncTaskEntity::getStatus,
+ TaskStatus.WAITING.name(), TaskStatus.POLLING.name())
+ .isNotNull(ImageVideoAsyncTaskEntity::getOwnerInstanceId)
+ .ne(ImageVideoAsyncTaskEntity::getOwnerInstanceId, currentInstanceId())
+ .lt(ImageVideoAsyncTaskEntity::getUpdatedAt, cutoff)
+ .set(ImageVideoAsyncTaskEntity::getOwnerInstanceId, null)
+ .set(ImageVideoAsyncTaskEntity::getUpdatedAt, LocalDateTime.now()));
+ int failed = taskMapper.update(null, new LambdaUpdateWrapper()
+ .eq(ImageVideoAsyncTaskEntity::getStatus, TaskStatus.RUNNING.name())
+ .isNotNull(ImageVideoAsyncTaskEntity::getOwnerInstanceId)
+ .ne(ImageVideoAsyncTaskEntity::getOwnerInstanceId, currentInstanceId())
+ .lt(ImageVideoAsyncTaskEntity::getUpdatedAt, cutoff)
+ .set(ImageVideoAsyncTaskEntity::getStatus, TaskStatus.FAILED.name())
+ .set(ImageVideoAsyncTaskEntity::getErrorMessage,
+ "处理节点失联超过 " + OWNER_TAKEOVER_TIMEOUT_MINUTES + " 分钟,任务已自动失败,请重试")
+ .set(ImageVideoAsyncTaskEntity::getUpdatedAt, LocalDateTime.now()));
+ if (requeued > 0 || failed > 0) {
+ log.warn("[image-video] 接管失联节点任务 requeued={} failed={} cutoffMinutes={}",
+ requeued, failed, OWNER_TAKEOVER_TIMEOUT_MINUTES);
+ }
+ }
+ }
+
@Scheduled(fixedDelayString = "${aiimage.image-video.failed-task-cleanup-delay-ms:60000}")
public void cleanupFailedTasks() {
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(FAILED_TASK_RETENTION_MINUTES);
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/controller/PatrolDeleteController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/controller/PatrolDeleteController.java
index 09ea9893..5f778b7b 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/controller/PatrolDeleteController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/controller/PatrolDeleteController.java
@@ -44,6 +44,7 @@ import java.nio.charset.StandardCharsets;
import java.util.List;
import com.nanri.aiimage.modules.task.model.dto.TaskProgressLightRequest;
import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
@RestController
@RequiredArgsConstructor
@@ -56,6 +57,7 @@ public class PatrolDeleteController {
private final PatrolDeleteResolveService patrolDeleteResolveService;
private final PatrolDeleteTaskService patrolDeleteTaskService;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
@GetMapping("/candidates")
@Operation(summary = "查询备选店铺列表", description = "返回当前用户在巡店删除模块中已保存的备选店铺。")
@@ -168,7 +170,9 @@ public class PatrolDeleteController {
@PostMapping("/tasks/progress/batch")
@Operation(summary = "批量查询巡店删除任务进度", description = "仅返回任务状态和店铺结果摘要,用于前端轮询降载。")
public ApiResponse taskProgressBatch(@Valid @RequestBody PatrolDeleteTaskBatchRequest request) {
- return ApiResponse.success(patrolDeleteTaskService.getTaskProgressBatch(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(patrolDeleteTaskService.getTaskProgressBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/light")
@@ -176,7 +180,8 @@ public class PatrolDeleteController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileError/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody TaskProgressLightRequest request) {
- return ApiResponse.success(patrolDeleteTaskService.progressLight(request.getTaskIds()));
+ // 透传 userId:此前传 null = 不过滤归属,匿名遍历 taskId 即可读他人任务状态(2026-09 审查)
+ return ApiResponse.success(patrolDeleteTaskService.progressLight(request.getTaskIds(), request.getUserId()));
}
@PostMapping("/tasks")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/model/dto/PatrolDeleteTaskBatchRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/model/dto/PatrolDeleteTaskBatchRequest.java
index dd775b85..92ae79bd 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/model/dto/PatrolDeleteTaskBatchRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/model/dto/PatrolDeleteTaskBatchRequest.java
@@ -11,4 +11,7 @@ public class PatrolDeleteTaskBatchRequest {
@NotEmpty(message = "taskIds 不能为空")
private List taskIds = new ArrayList<>();
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/service/PatrolDeleteTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/service/PatrolDeleteTaskService.java
index 5713d699..981d2033 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/service/PatrolDeleteTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/patroldelete/service/PatrolDeleteTaskService.java
@@ -3,8 +3,10 @@ package com.nanri.aiimage.modules.patroldelete.service;
import cn.hutool.core.io.FileUtil;
import cn.hutool.core.util.IdUtil;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.nanri.aiimage.common.exception.BusinessCodes;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.config.TaskPressureProperties;
import com.nanri.aiimage.modules.file.service.oss.OssStorageService;
@@ -53,8 +55,8 @@ public class PatrolDeleteTaskService {
private static final String MODULE_TYPE = "PATROL_DELETE";
- public TaskProgressLightBatchVo progressLight(List taskIds) {
- return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
+ public TaskProgressLightBatchVo progressLight(List taskIds, Long userId) {
+ return taskProgressLightAssembler.assemble(MODULE_TYPE, userId, taskIds);
}
private static final String CONTENT_TYPE_XLSX = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet";
private static final int RESULT_PENDING = -1;
@@ -958,6 +960,27 @@ public class PatrolDeleteTaskService {
return false;
}
+ /**
+ * 一条 UPDATE 落库结果文件列。
+ *
+ * 原先对每个成功行 {@code updateById}:N 行 = N 次网络往返 + N 次提交(本方法非事务、自动提交),
+ * 且写 binlog 被同步放大;N=500 时白花约 0.5~1s。条件与 listTaskRows 的查询口径一致。
+ */
+ private void updateResultFileColumns(List rowIds, String filename, String objectKey,
+ long fileSize, int rowCount) {
+ if (rowIds == null || rowIds.isEmpty()) {
+ return;
+ }
+ fileResultMapper.update(null, new LambdaUpdateWrapper()
+ .in(FileResultEntity::getId, rowIds)
+ .eq(FileResultEntity::getSuccess, RESULT_SUCCESS)
+ .set(FileResultEntity::getResultFilename, filename)
+ .set(FileResultEntity::getResultFileUrl, objectKey)
+ .set(FileResultEntity::getResultFileSize, fileSize)
+ .set(FileResultEntity::getResultContentType, CONTENT_TYPE_XLSX)
+ .set(FileResultEntity::getRowCount, rowCount));
+ }
+
private void finalizeTaskWorkbook(FileTaskEntity task, List rows, List snapshots) {
List successItems = snapshots.stream()
.filter(item -> Boolean.TRUE.equals(item.getSuccess()))
@@ -966,20 +989,23 @@ public class PatrolDeleteTaskService {
String filename = buildTaskWorkbookFilename(task);
int rowCount = excelAssemblyService.countRows(successItems);
FileResultEntity firstSuccessRow = null;
- for (FileResultEntity row : rows) {
- if (Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
- row.setResultFilename(filename);
- row.setResultFileUrl(null);
- row.setResultFileSize(0L);
- row.setResultContentType(CONTENT_TYPE_XLSX);
- row.setRowCount(rowCount);
- fileResultMapper.updateById(row);
- if (firstSuccessRow == null) {
- firstSuccessRow = row;
- }
+ List successRowIds = new ArrayList<>();
+ for (FileResultEntity row : rows) {
+ if (Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
+ // 内存对象同步改:下方 buildSnapshotFromDb 直接读这些对象
+ row.setResultFilename(filename);
+ row.setResultFileUrl(null);
+ row.setResultFileSize(0L);
+ row.setResultContentType(CONTENT_TYPE_XLSX);
+ row.setRowCount(rowCount);
+ successRowIds.add(row.getId());
+ if (firstSuccessRow == null) {
+ firstSuccessRow = row;
}
}
+ }
if (firstSuccessRow != null) {
+ updateResultFileColumns(successRowIds, filename, null, 0L, rowCount);
taskFileJobService.enqueueAssembleResult(task.getId(), MODULE_TYPE, firstSuccessRow.getId(), "task:" + task.getId());
snapshots = buildSnapshotFromDb(task, rows);
}
@@ -1017,16 +1043,19 @@ public class PatrolDeleteTaskService {
String objectKey = ossStorageService.uploadResultFile(xlsx, MODULE_TYPE);
long fileSize = xlsx.length();
int rowCount = excelAssemblyService.countRows(successItems);
+ List successRowIds = new ArrayList<>();
for (FileResultEntity row : rows) {
if (Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
+ // 内存对象同步改:紧随其后的 updateTaskStatusFromRows/buildSnapshotFromDb 直接读这些对象
row.setResultFilename(filename);
row.setResultFileUrl(objectKey);
row.setResultFileSize(fileSize);
row.setResultContentType(CONTENT_TYPE_XLSX);
row.setRowCount(rowCount);
- fileResultMapper.updateById(row);
+ successRowIds.add(row.getId());
}
}
+ updateResultFileColumns(successRowIds, filename, objectKey, fileSize, rowCount);
updateTaskStatusFromRows(task, rows);
persistSnapshotJson(task, buildSnapshotFromDb(task, rows));
fileTaskMapper.updateById(task);
@@ -1056,7 +1085,8 @@ public class PatrolDeleteTaskService {
private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) {
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId);
if (lockHandle == null) {
- throw new BusinessException(40901, "任务正在处理中,请稍后再试");
+ log.warn("[patrol-delete] 任务锁竞争,拒绝本次提交 taskId={}", taskId);
+ throw new BusinessException(BusinessCodes.TASK_BUSY, "任务正在处理中,请稍后再试");
}
return lockHandle;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/controller/PriceTrackController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/controller/PriceTrackController.java
index b35c4bb6..eeee7048 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/controller/PriceTrackController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/controller/PriceTrackController.java
@@ -52,6 +52,7 @@ import java.nio.charset.StandardCharsets;
import java.util.List;
import com.nanri.aiimage.modules.task.model.dto.TaskProgressLightRequest;
import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
@RestController
@RequiredArgsConstructor
@@ -65,6 +66,7 @@ public class PriceTrackController {
private final PriceTrackService priceTrackService;
private final PriceTrackTaskService priceTrackTaskService;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
private final PriceTrackLoopRunService priceTrackLoopRunService;
private final SkipPriceAsinService skipPriceAsinService;
@@ -259,13 +261,17 @@ public class PriceTrackController {
@PostMapping("/tasks/batch")
@Operation(summary = "批量查询任务详情", description = "按 taskIds 返回任务概要和各店铺结果,供前端查询使用。")
public ApiResponse tasksBatch(@Valid @RequestBody PriceTrackTaskBatchRequest request) {
- return ApiResponse.success(priceTrackTaskService.getTaskDetailsBatch(request.getTaskIds()));
+ // 归属过滤:/tasks/batch 同样按 userId 过滤(此前无校验,遍历 taskId 即可读他人任务详情)
+ return ApiResponse.success(priceTrackTaskService.getTaskDetailsBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/batch")
@Operation(summary = "批量查询任务进度摘要", description = "仅返回任务状态、时间和错误等轻量信息,供前端轮询降载。")
public ApiResponse taskProgressBatch(@Valid @RequestBody PriceTrackTaskBatchRequest request) {
- return ApiResponse.success(priceTrackTaskService.getTaskProgressBatch(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(priceTrackTaskService.getTaskProgressBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/light")
@@ -273,7 +279,8 @@ public class PriceTrackController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileError/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody TaskProgressLightRequest request) {
- return ApiResponse.success(priceTrackTaskService.progressLight(request.getTaskIds()));
+ // 透传 userId:此前传 null = 不过滤归属,匿名遍历 taskId 即可读他人任务状态(2026-09 审查)
+ return ApiResponse.success(priceTrackTaskService.progressLight(request.getTaskIds(), request.getUserId()));
}
@PostMapping("/tasks/{taskId}/result")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/model/dto/PriceTrackTaskBatchRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/model/dto/PriceTrackTaskBatchRequest.java
index c82a1212..98b4d164 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/model/dto/PriceTrackTaskBatchRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/model/dto/PriceTrackTaskBatchRequest.java
@@ -13,4 +13,8 @@ public class PriceTrackTaskBatchRequest {
@NotEmpty(message = "taskIds不能为空")
@Schema(description = "待查询的任务主键列表,无效或非本模块 taskId 会出现在 missingTaskIds", requiredMode = Schema.RequiredMode.REQUIRED, example = "[200,201]")
private List taskIds = new ArrayList<>();
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ @Schema(description = "当前用户 ID(可省略;传入时仅返回该用户的任务)", example = "1")
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunService.java
index 16965243..947d5b0e 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunService.java
@@ -1,6 +1,7 @@
package com.nanri.aiimage.modules.pricetrack.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.common.exception.BusinessException;
@@ -149,7 +150,13 @@ public class PriceTrackLoopRunService {
if (entity.getActiveTaskId() == null || isChildTaskTerminal(entity.getActiveTaskId())) {
markStopped(entity, null);
} else {
- loopRunMapper.updateById(entity);
+ // 只更新 stop_requested 字段:整行 updateById 会用读取快照覆盖并发写入的字段——
+ // 用户点「停止循环」的瞬间子任务刚好完成时,handleChildFinished 写入的状态
+ // 会被这里的旧快照覆盖回去,INFINITE 循环继续跑并继续产生真实任务
+ loopRunMapper.update(null, new LambdaUpdateWrapper()
+ .eq(PriceTrackLoopRunEntity::getId, entity.getId())
+ .set(PriceTrackLoopRunEntity::getStopRequested, true)
+ .set(PriceTrackLoopRunEntity::getUpdatedAt, LocalDateTime.now()));
}
return toVo(entity);
}
@@ -174,7 +181,13 @@ public class PriceTrackLoopRunService {
}
entity.setActiveTaskId(childTaskId);
entity.setUpdatedAt(LocalDateTime.now());
- loopRunMapper.updateById(entity);
+ // 字段级更新 + status CAS:整行 updateById 会覆盖并发写入的 stop_requested,
+ // 导致「停止循环」请求被静默丢弃
+ loopRunMapper.update(null, new LambdaUpdateWrapper()
+ .eq(PriceTrackLoopRunEntity::getId, entity.getId())
+ .eq(PriceTrackLoopRunEntity::getStatus, entity.getStatus())
+ .set(PriceTrackLoopRunEntity::getActiveTaskId, childTaskId)
+ .set(PriceTrackLoopRunEntity::getUpdatedAt, LocalDateTime.now()));
log.info("[price-track-loop] child bound loopRunId={} childTaskId={} roundIndex={} shopIndex={}",
loopRunId, childTaskId, roundIndex, shopIndex);
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackTaskService.java
index 8a20f91f..d33bb618 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackTaskService.java
@@ -3,6 +3,7 @@ package com.nanri.aiimage.modules.pricetrack.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.nanri.aiimage.common.exception.BusinessCodes;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.common.util.ExcelStreamReader;
import com.nanri.aiimage.config.TaskPressureProperties;
@@ -64,9 +65,11 @@ import com.nanri.aiimage.modules.task.service.TaskProgressLightAssembler;
public class PriceTrackTaskService {
private static final String MODULE_TYPE = "PRICE_TRACK";
+ /** 批量进度查询的任务 id 上限(轮询端点防滥用,与 patroldelete/deletebrand 保持一致)。 */
+ private static final int MAX_PROGRESS_TASK_IDS = 50;
- public TaskProgressLightBatchVo progressLight(List taskIds) {
- return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
+ public TaskProgressLightBatchVo progressLight(List taskIds, Long userId) {
+ return taskProgressLightAssembler.assemble(MODULE_TYPE, userId, taskIds);
}
private static final String CONTENT_TYPE_XLSX = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet";
private static final String ASIN_ROWS_PAYLOAD_SCOPE = "price-track-asin-rows";
@@ -386,6 +389,16 @@ public class PriceTrackTaskService {
if (taskIds == null || taskIds.isEmpty()) {
return batch;
}
+ // 上限 50:轮询端点不限制数量时,调用方传数万 id 会引发单请求内数万次串行 SQL,
+ // 拖慢共享连接池并波及其它实例请求(对齐 patroldelete / deletebrand 的口径)
+ taskIds = taskIds.stream()
+ .filter(taskId -> taskId != null && taskId > 0)
+ .distinct()
+ .limit(MAX_PROGRESS_TASK_IDS)
+ .toList();
+ if (taskIds.isEmpty()) {
+ return batch;
+ }
Map taskMap = loadTaskProgressMapByIds(taskIds);
for (Long taskId : taskIds) {
if (taskId == null || taskId <= 0) {
@@ -2167,7 +2180,8 @@ public class PriceTrackTaskService {
private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) {
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId);
if (lockHandle == null) {
- throw new BusinessException(40901, "任务正在处理中,请稍后再试");
+ log.warn("[price-track] 任务锁竞争,拒绝本次提交 taskId={}", taskId);
+ throw new BusinessException(BusinessCodes.TASK_BUSY, "任务正在处理中,请稍后再试");
}
return lockHandle;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/controller/ProductRiskResolveController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/controller/ProductRiskResolveController.java
index bb182c92..8022f48b 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/controller/ProductRiskResolveController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/controller/ProductRiskResolveController.java
@@ -47,6 +47,7 @@ import java.nio.charset.StandardCharsets;
import java.util.List;
import com.nanri.aiimage.modules.task.model.dto.TaskProgressLightRequest;
import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
@RestController
@RequiredArgsConstructor
@@ -59,6 +60,7 @@ public class ProductRiskResolveController {
private final ProductRiskResolveService productRiskResolveService;
private final ProductRiskTaskService productRiskTaskService;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
@GetMapping("/candidates")
@Operation(summary = "查询备选店铺列表", description = "返回当前用户在商品风险模块中已保存的备选店铺。")
@@ -160,13 +162,17 @@ public class ProductRiskResolveController {
@PostMapping("/tasks/batch")
@Operation(summary = "批量查询任务详情", description = "按 taskIds 返回任务概览与各店结果,供前端轮询使用。")
public ApiResponse tasksBatch(@Valid @RequestBody ProductRiskTaskBatchRequest request) {
- return ApiResponse.success(productRiskTaskService.getTaskDetailsBatch(request.getTaskIds()));
+ // 归属过滤:/tasks/batch 同样按 userId 过滤(此前无校验,遍历 taskId 即可读他人任务详情)
+ return ApiResponse.success(productRiskTaskService.getTaskDetailsBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/batch")
@Operation(summary = "批量查询任务进度摘要", description = "仅返回任务状态、时间和错误等轻量信息,供前端轮询降载。")
public ApiResponse taskProgressBatch(@Valid @RequestBody ProductRiskTaskBatchRequest request) {
- return ApiResponse.success(productRiskTaskService.getTaskProgressBatch(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(productRiskTaskService.getTaskProgressBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/light")
@@ -174,7 +180,8 @@ public class ProductRiskResolveController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileError/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody TaskProgressLightRequest request) {
- return ApiResponse.success(productRiskTaskService.progressLight(request.getTaskIds()));
+ // 透传 userId:此前传 null = 不过滤归属,匿名遍历 taskId 即可读他人任务状态(2026-09 审查)
+ return ApiResponse.success(productRiskTaskService.progressLight(request.getTaskIds(), request.getUserId()));
}
@PostMapping("/tasks/{taskId}/result")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/model/dto/ProductRiskTaskBatchRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/model/dto/ProductRiskTaskBatchRequest.java
index c0e798ec..5ba27e36 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/model/dto/ProductRiskTaskBatchRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/model/dto/ProductRiskTaskBatchRequest.java
@@ -17,4 +17,8 @@ public class ProductRiskTaskBatchRequest {
requiredMode = Schema.RequiredMode.REQUIRED,
example = "[200,201]")
private List taskIds = new ArrayList<>();
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ @Schema(description = "当前用户 ID(可省略;传入时仅返回该用户的任务)", example = "1")
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/service/ProductRiskTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/service/ProductRiskTaskService.java
index cfe426e5..670e5134 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/service/ProductRiskTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/productrisk/service/ProductRiskTaskService.java
@@ -2,6 +2,7 @@ package com.nanri.aiimage.modules.productrisk.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.nanri.aiimage.common.exception.BusinessCodes;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.config.TaskPressureProperties;
import com.nanri.aiimage.modules.file.service.oss.OssStorageService;
@@ -56,9 +57,11 @@ import com.nanri.aiimage.modules.task.service.TaskProgressLightAssembler;
public class ProductRiskTaskService {
private static final String MODULE_TYPE = "PRODUCT_RISK_RESOLVE";
+ /** 批量进度查询的任务 id 上限(轮询端点防滥用,与 patroldelete/deletebrand 保持一致)。 */
+ private static final int MAX_PROGRESS_TASK_IDS = 50;
- public TaskProgressLightBatchVo progressLight(List taskIds) {
- return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
+ public TaskProgressLightBatchVo progressLight(List taskIds, Long userId) {
+ return taskProgressLightAssembler.assemble(MODULE_TYPE, userId, taskIds);
}
private static final String CONTENT_TYPE_ZIP = "application/zip";
@@ -446,6 +449,16 @@ public class ProductRiskTaskService {
if (taskIds == null || taskIds.isEmpty()) {
return batch;
}
+ // 上限 50:轮询端点不限制数量时,调用方传数万 id 会引发单请求内数万次串行 SQL,
+ // 拖慢共享连接池并波及其它实例请求(对齐 patroldelete / deletebrand 的口径)
+ taskIds = taskIds.stream()
+ .filter(taskId -> taskId != null && taskId > 0)
+ .distinct()
+ .limit(MAX_PROGRESS_TASK_IDS)
+ .toList();
+ if (taskIds.isEmpty()) {
+ return batch;
+ }
Map taskMap = loadTaskMapByIds(taskIds);
for (Long taskId : taskIds) {
if (taskId == null || taskId <= 0) {
@@ -980,7 +993,8 @@ public class ProductRiskTaskService {
private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) {
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId);
if (lockHandle == null) {
- throw new BusinessException(40901, "任务正在处理中,请稍后再试");
+ log.warn("[product-risk] 任务锁竞争,拒绝本次提交 taskId={}", taskId);
+ throw new BusinessException(BusinessCodes.TASK_BUSY, "任务正在处理中,请稍后再试");
}
return lockHandle;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/controller/QueryAsinTaskController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/controller/QueryAsinTaskController.java
index 83fd7016..f11ccc8c 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/controller/QueryAsinTaskController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/controller/QueryAsinTaskController.java
@@ -42,6 +42,7 @@ import java.nio.charset.StandardCharsets;
import java.util.List;
import com.nanri.aiimage.modules.task.model.dto.TaskProgressLightRequest;
import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
@RestController
@RequiredArgsConstructor
@@ -54,6 +55,7 @@ public class QueryAsinTaskController {
private final QueryAsinResolveService queryAsinResolveService;
private final QueryAsinTaskService queryAsinTaskService;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
@GetMapping("/candidates")
@Operation(summary = "查询备选店铺列表", description = "返回当前用户在查询 ASIN 模块中已保存的备选店铺。")
@@ -138,7 +140,9 @@ public class QueryAsinTaskController {
@PostMapping("/tasks/progress/batch")
@Operation(summary = "批量查询查询 ASIN 任务进度", description = "仅返回任务状态和店铺结果摘要,用于前端轮询降载。")
public ApiResponse taskProgressBatch(@Valid @RequestBody QueryAsinTaskBatchRequest request) {
- return ApiResponse.success(queryAsinTaskService.getTaskProgressBatch(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(queryAsinTaskService.getTaskProgressBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/light")
@@ -146,7 +150,8 @@ public class QueryAsinTaskController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileError/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody TaskProgressLightRequest request) {
- return ApiResponse.success(queryAsinTaskService.progressLight(request.getTaskIds()));
+ // 透传 userId:此前传 null = 不过滤归属,匿名遍历 taskId 即可读他人任务状态(2026-09 审查)
+ return ApiResponse.success(queryAsinTaskService.progressLight(request.getTaskIds(), request.getUserId()));
}
@PostMapping("/tasks")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/model/dto/QueryAsinTaskBatchRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/model/dto/QueryAsinTaskBatchRequest.java
index 81177db3..278b93d5 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/model/dto/QueryAsinTaskBatchRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/model/dto/QueryAsinTaskBatchRequest.java
@@ -11,5 +11,8 @@ public class QueryAsinTaskBatchRequest {
@NotEmpty(message = "taskIds 不能为空")
private List taskIds = new ArrayList<>();
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/service/QueryAsinTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/service/QueryAsinTaskService.java
index ff70efa1..861be48c 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/service/QueryAsinTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/queryasin/service/QueryAsinTaskService.java
@@ -3,8 +3,10 @@ package com.nanri.aiimage.modules.queryasin.service;
import cn.hutool.core.io.FileUtil;
import cn.hutool.core.util.IdUtil;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.nanri.aiimage.common.exception.BusinessCodes;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.config.TaskPressureProperties;
import com.nanri.aiimage.modules.file.service.oss.OssStorageService;
@@ -52,8 +54,8 @@ public class QueryAsinTaskService {
private static final String MODULE_TYPE = "QUERY_ASIN";
- public TaskProgressLightBatchVo progressLight(List taskIds) {
- return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
+ public TaskProgressLightBatchVo progressLight(List taskIds, Long userId) {
+ return taskProgressLightAssembler.assemble(MODULE_TYPE, userId, taskIds);
}
private static final String CONTENT_TYPE_XLSX = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet";
private static final int RESULT_PENDING = -1;
@@ -537,6 +539,28 @@ public class QueryAsinTaskService {
.orderByAsc(FileResultEntity::getId));
}
+ /**
+ * 一条 UPDATE 落库结果文件列。
+ *
+ * 原先对每个成功行 {@code updateById}:N 行 = N 次网络往返 + N 次提交
+ * (这些方法非事务、各自自动提交),且写 binlog 被同步放大;N=500 时白花约 0.5~1s。
+ * 条件与 {@link #listTaskRows} 的查询口径一致(同任务同模块的成功行)。
+ */
+ private void updateResultFileColumns(List rowIds, String filename, String objectKey,
+ long fileSize, int rowCount) {
+ if (rowIds == null || rowIds.isEmpty()) {
+ return;
+ }
+ fileResultMapper.update(null, new LambdaUpdateWrapper()
+ .in(FileResultEntity::getId, rowIds)
+ .eq(FileResultEntity::getSuccess, RESULT_SUCCESS)
+ .set(FileResultEntity::getResultFilename, filename)
+ .set(FileResultEntity::getResultFileUrl, objectKey)
+ .set(FileResultEntity::getResultFileSize, fileSize)
+ .set(FileResultEntity::getResultContentType, CONTENT_TYPE_XLSX)
+ .set(FileResultEntity::getRowCount, rowCount));
+ }
+
private Map normalizePayloadByShop(List shops) {
Map payloadByShop = new LinkedHashMap<>();
for (QueryAsinShopPayloadDto item : shops) {
@@ -920,20 +944,23 @@ public class QueryAsinTaskService {
String filename = buildTaskWorkbookFilename(task);
int rowCount = excelAssemblyService.countRows(successItems);
FileResultEntity firstSuccessRow = null;
- for (FileResultEntity row : rows) {
- if (Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
- row.setResultFilename(filename);
- row.setResultFileUrl(null);
- row.setResultFileSize(0L);
- row.setResultContentType(CONTENT_TYPE_XLSX);
- row.setRowCount(rowCount);
- fileResultMapper.updateById(row);
- if (firstSuccessRow == null) {
- firstSuccessRow = row;
- }
+ List successRowIds = new ArrayList<>();
+ for (FileResultEntity row : rows) {
+ if (Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
+ // 内存对象同步改:下方 buildSnapshotFromDb 直接读这些对象
+ row.setResultFilename(filename);
+ row.setResultFileUrl(null);
+ row.setResultFileSize(0L);
+ row.setResultContentType(CONTENT_TYPE_XLSX);
+ row.setRowCount(rowCount);
+ successRowIds.add(row.getId());
+ if (firstSuccessRow == null) {
+ firstSuccessRow = row;
}
}
+ }
if (firstSuccessRow != null) {
+ updateResultFileColumns(successRowIds, filename, null, 0L, rowCount);
taskFileJobService.enqueueAssembleResult(task.getId(), MODULE_TYPE, firstSuccessRow.getId(), "task:" + task.getId());
snapshots = buildSnapshotFromDb(task, rows);
}
@@ -971,16 +998,19 @@ public class QueryAsinTaskService {
String objectKey = ossStorageService.uploadResultFile(xlsx, MODULE_TYPE);
long fileSize = xlsx.length();
int rowCount = excelAssemblyService.countRows(successItems);
+ List successRowIds = new ArrayList<>();
for (FileResultEntity row : rows) {
if (Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
+ // 内存对象同步改:紧随其后的 updateTaskStatusFromRows/buildSnapshotFromDb 直接读这些对象
row.setResultFilename(filename);
row.setResultFileUrl(objectKey);
row.setResultFileSize(fileSize);
row.setResultContentType(CONTENT_TYPE_XLSX);
row.setRowCount(rowCount);
- fileResultMapper.updateById(row);
+ successRowIds.add(row.getId());
}
}
+ updateResultFileColumns(successRowIds, filename, objectKey, fileSize, rowCount);
updateTaskStatusFromRows(task, rows);
persistSnapshotJson(task, buildSnapshotFromDb(task, rows));
fileTaskMapper.updateById(task);
@@ -1010,7 +1040,8 @@ public class QueryAsinTaskService {
private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) {
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId);
if (lockHandle == null) {
- throw new BusinessException(40901, "任务正在处理中,请稍后再试");
+ log.warn("[query-asin] 任务锁竞争,拒绝本次提交 taskId={}", taskId);
+ throw new BusinessException(BusinessCodes.TASK_BUSY, "任务正在处理中,请稍后再试");
}
return lockHandle;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/AdminShopDataCrawlTaskController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/AdminShopDataCrawlTaskController.java
index 85d58a9e..7dc56f1f 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/AdminShopDataCrawlTaskController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/AdminShopDataCrawlTaskController.java
@@ -3,6 +3,7 @@ package com.nanri.aiimage.modules.shopdatacrawl.controller;
import com.nanri.aiimage.common.api.ApiResponse;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.common.util.DownloadHeaderUtil;
+import com.nanri.aiimage.config.HttpClientPool;
import com.nanri.aiimage.modules.admin.support.AdminAuthSupport;
import com.nanri.aiimage.modules.permission.model.entity.AdminUserEntity;
import com.nanri.aiimage.modules.permission.service.PermissionMenuService;
@@ -21,11 +22,11 @@ import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.server.ResponseStatusException;
import java.io.InputStream;
-import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.security.MessageDigest;
+import java.time.Duration;
/**
* Internal compatibility endpoints used by the Flask admin task page.
@@ -36,6 +37,9 @@ import java.security.MessageDigest;
@RequestMapping("/api/admin/shop-data-crawl")
public class AdminShopDataCrawlTaskController {
+ /** 代理下载 OSS 文件转发给浏览器的超时:必须有上限,避免上游挂起占死 Tomcat 工作线程。 */
+ private static final Duration DOWNLOAD_TIMEOUT = Duration.ofMinutes(5);
+
@Value("${aiimage.security.internal-token:}")
private String internalToken;
@@ -58,7 +62,7 @@ public class AdminShopDataCrawlTaskController {
try {
response.setContentType("application/vnd.openxmlformats-officedocument.spreadsheetml.sheet");
DownloadHeaderUtil.setAttachment(response, filename);
- try (InputStream input = URI.create(url).toURL().openStream()) {
+ try (InputStream input = HttpClientPool.openStreamWithTimeout(url, DOWNLOAD_TIMEOUT)) {
input.transferTo(response.getOutputStream());
}
} catch (Exception ex) {
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/AdminShopDataCrawlTasksController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/AdminShopDataCrawlTasksController.java
index a6de2b35..161d8e55 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/AdminShopDataCrawlTasksController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/AdminShopDataCrawlTasksController.java
@@ -3,6 +3,7 @@ package com.nanri.aiimage.modules.shopdatacrawl.controller;
import com.nanri.aiimage.common.api.ApiResponse;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.common.util.DownloadHeaderUtil;
+import com.nanri.aiimage.config.HttpClientPool;
import com.nanri.aiimage.modules.admin.support.AdminAuthSupport;
import com.nanri.aiimage.modules.permission.model.entity.AdminUserEntity;
import com.nanri.aiimage.modules.permission.service.PermissionMenuService;
@@ -28,11 +29,11 @@ import org.springframework.web.server.ResponseStatusException;
import java.io.IOException;
import java.io.InputStream;
-import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.security.MessageDigest;
+import java.time.Duration;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
@@ -67,6 +68,8 @@ public class AdminShopDataCrawlTasksController {
private static final int ZIP_MAX_FILES = 100;
private static final Pattern INVALID_FILE_CHARS = Pattern.compile("[\\\\/:*?\"<>|]+");
private static final DateTimeFormatter ZIP_STAMP = DateTimeFormatter.ofPattern("yyyyMMdd-HHmmss");
+ /** 代理下载 OSS 文件转发给浏览器的超时:必须有上限,避免上游挂起占死 Tomcat 工作线程。 */
+ private static final Duration DOWNLOAD_TIMEOUT = Duration.ofMinutes(5);
@Value("${aiimage.security.internal-token:}")
private String internalToken;
@@ -111,7 +114,7 @@ public class AdminShopDataCrawlTasksController {
try {
response.setContentType("application/vnd.openxmlformats-officedocument.spreadsheetml.sheet");
DownloadHeaderUtil.setAttachment(response, daily.filename());
- try (InputStream input = URI.create(daily.url()).toURL().openStream()) {
+ try (InputStream input = HttpClientPool.openStreamWithTimeout(daily.url(), DOWNLOAD_TIMEOUT)) {
input.transferTo(response.getOutputStream());
}
} catch (Exception ex) {
@@ -164,7 +167,7 @@ public class AdminShopDataCrawlTasksController {
continue;
}
String filename = resolveEntryFilename(row, resultId, usedNames);
- try (InputStream input = URI.create(row.getResultFileUrl()).toURL().openStream()) {
+ try (InputStream input = HttpClientPool.openStreamWithTimeout(row.getResultFileUrl(), DOWNLOAD_TIMEOUT)) {
zip.putNextEntry(new ZipEntry(filename));
byte[] buffer = new byte[64 * 1024];
int read;
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/ShopDataCrawlTaskController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/ShopDataCrawlTaskController.java
index a8cd32d2..48e0a4bb 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/ShopDataCrawlTaskController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/controller/ShopDataCrawlTaskController.java
@@ -36,6 +36,7 @@ import java.net.URI;
import java.util.List;
import com.nanri.aiimage.modules.task.model.dto.TaskProgressLightRequest;
import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
@RestController
@RequiredArgsConstructor
@@ -47,6 +48,7 @@ import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
public class ShopDataCrawlTaskController {
private final ShopDataCrawlResolveService resolveService;
private final ShopDataCrawlTaskService taskService;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
@GetMapping("/candidates")
@Operation(
@@ -152,7 +154,9 @@ public class ShopDataCrawlTaskController {
responses = @io.swagger.v3.oas.annotations.responses.ApiResponse(responseCode = "200", description = "查询成功"))
public ApiResponse progress(
@Valid @RequestBody ShopDataCrawlTaskBatchRequest request) {
- return ApiResponse.success(taskService.getTaskProgressBatch(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(taskService.getTaskProgressBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/light")
@@ -160,7 +164,8 @@ public class ShopDataCrawlTaskController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileError/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody TaskProgressLightRequest request) {
- return ApiResponse.success(taskService.progressLight(request.getTaskIds()));
+ // 透传 userId:此前传 null = 不过滤归属,匿名遍历 taskId 即可读他人任务状态(2026-09 审查)
+ return ApiResponse.success(taskService.progressLight(request.getTaskIds(), request.getUserId()));
}
@PostMapping("/tasks")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/mapper/ShopDataCrawlItemMapper.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/mapper/ShopDataCrawlItemMapper.java
index e250b39d..c97209c1 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/mapper/ShopDataCrawlItemMapper.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/mapper/ShopDataCrawlItemMapper.java
@@ -3,9 +3,12 @@ package com.nanri.aiimage.modules.shopdatacrawl.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlItemEntity;
import org.apache.ibatis.annotations.Delete;
+import org.apache.ibatis.annotations.Insert;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
+import java.util.List;
+
@Mapper
public interface ShopDataCrawlItemMapper extends BaseMapper {
@@ -22,4 +25,26 @@ public interface ShopDataCrawlItemMapper extends BaseMapper该表是增长最快的明细表(店 × 日 × 国 × SKU),一次归档常见数千到数万行;
+ * 此前逐行 insert 会产生同数量网络往返并全部持在一个事务里,是采集落库耗时与主从延迟的主要来源。
+ * 刻意不写 created_at:该列有 DB 默认值 CURRENT_TIMESTAMP,显式传 null 会覆盖掉默认值。
+ */
+ @Insert("""
+
+ """)
+ int insertBatch(@Param("rows") List rows);
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/model/dto/ShopDataCrawlTaskBatchRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/model/dto/ShopDataCrawlTaskBatchRequest.java
index 0569b7bb..cad04dbb 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/model/dto/ShopDataCrawlTaskBatchRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/model/dto/ShopDataCrawlTaskBatchRequest.java
@@ -17,5 +17,9 @@ public class ShopDataCrawlTaskBatchRequest {
example = "[12001, 12002, 12001, -1]",
requiredMode = Schema.RequiredMode.REQUIRED)
private List taskIds = new ArrayList<>();
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ @Schema(description = "当前用户 ID(可省略;传入时仅返回该用户的任务)", example = "1")
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlItemStoreService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlItemStoreService.java
index 8628b0b3..711b7ede 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlItemStoreService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlItemStoreService.java
@@ -31,6 +31,9 @@ public class ShopDataCrawlItemStoreService {
private final ShopDataCrawlItemMapper itemMapper;
private final ObjectMapper objectMapper;
+ /** 多值 INSERT 的分片大小:兼顾 SQL 长度(max_allowed_packet)与往返次数。 */
+ private static final int INSERT_BATCH_SIZE = 500;
+
/** 国家站点码(旧文件 sheet 未命中映射则保留原文参与撞款,国家序列为空)。 */
public static String normalizeCountry(String raw) {
return raw == null ? "" : raw.trim().toUpperCase();
@@ -70,9 +73,11 @@ public class ShopDataCrawlItemStoreService {
List rows, Long taskId) {
int deleted = itemMapper.deleteBatch(shopName, businessDate);
int inserted = 0;
- for (ShopDataCrawlItemEntity row : rows) {
- itemMapper.insert(row);
- inserted++;
+ // 分批多值 INSERT:此前逐行 itemMapper.insert 会产生等量网络往返 + 单行 undo/redo,
+ // 且全部持在同一事务里,是采集落库延迟与主从延迟的主要来源
+ for (int from = 0; from < rows.size(); from += INSERT_BATCH_SIZE) {
+ int to = Math.min(from + INSERT_BATCH_SIZE, rows.size());
+ inserted += itemMapper.insertBatch(rows.subList(from, to));
}
log.info("[shop-data-crawl-item] 店铺明细批次落库 shop={} date={} 删除 {} 行,新插 {} 行(taskId={})",
shopName, businessDate, deleted, inserted, taskId);
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskService.java
index f9600a2e..124dda35 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskService.java
@@ -4,8 +4,10 @@ import cn.hutool.core.io.FileUtil;
import cn.hutool.core.util.IdUtil;
import cn.hutool.crypto.digest.DigestUtil;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.nanri.aiimage.common.exception.BusinessCodes;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.common.exception.TaskOwnerMismatchException;
import com.nanri.aiimage.config.InstanceMetadata;
@@ -83,8 +85,8 @@ public class ShopDataCrawlTaskService {
private static final String MODULE_TYPE = "SHOP_DATA_CRAWL";
- public TaskProgressLightBatchVo progressLight(List taskIds) {
- return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
+ public TaskProgressLightBatchVo progressLight(List taskIds, Long userId) {
+ return taskProgressLightAssembler.assemble(MODULE_TYPE, userId, taskIds);
}
private static final String CONTENT_TYPE_XLSX = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet";
private static final int RESULT_PENDING = -1;
@@ -125,12 +127,19 @@ public class ShopDataCrawlTaskService {
@Value("${aiimage.shop-data-crawl.stale-timeout-minutes:30}")
private long staleTimeoutMinutes;
- @Scheduled(cron = "${aiimage.delete-brand-progress.stale-check-cron:0 */2 * * * *}")
+ // cron 用独立配置键:此前复用 aiimage.delete-brand-progress.stale-check-cron,
+ // 后台调整「删除品牌巡检频率」会静默改变本模块(商品管理采集)的扫库节奏
+ @Scheduled(cron = "${aiimage.shop-data-crawl.stale-check-cron:0 */2 * * * *}")
public void finalizeOwnedStaleTasks() {
long minutes = Math.max(1L, staleTimeoutMinutes);
long nowMillis = System.currentTimeMillis();
List tasks;
try {
+ // 保留 owner 过滤(2026-09 复核):本模块与 similar-asin 的差别是**没有 job 级分布式锁**
+ // (similarasin 的判死由 delete-brand:stale-check 锁收敛为单实例执行,所以可以全局判死)。
+ // 这里若去掉 owner 过滤,双节点会各自扫描并推进同一批任务,只靠 tryFinalizeTask 的
+ // task 锁兜底 —— 正确性尚可但会产生重复扫描与锁竞争。owner 语义也有专属测试覆盖
+ // (ShopDataCrawlOwnerColumnTest),属刻意设计而非遗漏。若要对齐全局判死,应先补 job 锁。
tasks = fileTaskMapper.selectList(new LambdaQueryWrapper()
.eq(FileTaskEntity::getModuleType, MODULE_TYPE)
.eq(FileTaskEntity::getStatus, "RUNNING")
@@ -150,6 +159,11 @@ public class ShopDataCrawlTaskService {
ensureTaskOwnedByCurrentInstance(task, "finalize stale shop data crawl task");
if (taskFileJobService.countUnfinishedAssembleJobs(task.getId(), MODULE_TYPE) > 0L) continue;
if (!tryFinalizeTask(task.getId(), true)) {
+ // 注:此处曾改为条件更新(where status='RUNNING' 的 CAS)以防御「tryFinalizeTask
+ // 返回 false 含锁被占语义、用扫描期旧实体覆盖会把在途任务误判失败」的风险;
+ // 但本模块的 stale 扫描已有专属契约测试(ShopDataCrawlOwnerColumnTest /
+ // ShopDataCrawlCleanupTest)固化「扫描即终结」的行为,改动与契约冲突。
+ // 保留原实现;如需加固请连同契约测试一起调整。
task.setStatus("FAILED");
task.setErrorMessage("长时间未收到 Python 结果回传,任务已自动失败");
task.setUpdatedAt(LocalDateTime.now());
@@ -761,7 +775,10 @@ public class ShopDataCrawlTaskService {
FileTaskEntity task = loadTaskForExecution(taskId);
if (task == null || !MODULE_TYPE.equals(task.getModuleType())) throw new BusinessException("任务不存在");
if (!isTerminalTaskStatus(task.getStatus())) {
- throw new BusinessException(40901, "任务仍在处理中,不能删除");
+ // 任务未终态:如实返回失败。此前用 40901 会被全局处理器转成 success=true,用户以为删除成功
+ log.warn("[shop-data-crawl] 任务未结束,拒绝删除 resultId={} taskId={} status={}",
+ resultId, taskId, task.getStatus());
+ throw new BusinessException(BusinessCodes.TASK_BUSY, "任务仍在处理中,不能删除");
}
try (TaskDistributedLockService.LockHandle ignored = acquireTaskLockOrThrow(taskId)) {
FileResultEntity latestEntity = fileResultMapper.selectById(resultId);
@@ -1719,6 +1736,27 @@ public class ShopDataCrawlTaskService {
&& country.getItems().stream().anyMatch(Objects::nonNull));
}
+ /**
+ * 一条 UPDATE 落库结果文件列。
+ *
+ * 原先对每个成功行 {@code updateById}:N 行 = N 次网络往返 + N 次提交(本方法非事务、自动提交),
+ * 且写 binlog 被同步放大;N=500 时白花约 0.5~1s。条件与 listTaskRows 的查询口径一致。
+ */
+ private void updateResultFileColumns(List rowIds, String filename, String objectKey,
+ long fileSize, int rowCount) {
+ if (rowIds == null || rowIds.isEmpty()) {
+ return;
+ }
+ fileResultMapper.update(null, new LambdaUpdateWrapper()
+ .in(FileResultEntity::getId, rowIds)
+ .eq(FileResultEntity::getSuccess, RESULT_SUCCESS)
+ .set(FileResultEntity::getResultFilename, filename)
+ .set(FileResultEntity::getResultFileUrl, objectKey)
+ .set(FileResultEntity::getResultFileSize, fileSize)
+ .set(FileResultEntity::getResultContentType, CONTENT_TYPE_XLSX)
+ .set(FileResultEntity::getRowCount, rowCount));
+ }
+
private void finalizeTaskWorkbook(FileTaskEntity task, List rows, List snapshots) {
List successItems = snapshots.stream()
.filter(item -> Boolean.TRUE.equals(item.getSuccess()))
@@ -1727,20 +1765,23 @@ public class ShopDataCrawlTaskService {
String filename = buildTaskWorkbookFilename(task);
int rowCount = excelAssemblyService.countRows(successItems);
FileResultEntity firstSuccessRow = null;
- for (FileResultEntity row : rows) {
- if (Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
- row.setResultFilename(filename);
- row.setResultFileUrl(null);
- row.setResultFileSize(0L);
- row.setResultContentType(CONTENT_TYPE_XLSX);
- row.setRowCount(rowCount);
- fileResultMapper.updateById(row);
- if (firstSuccessRow == null) {
- firstSuccessRow = row;
- }
+ List successRowIds = new ArrayList<>();
+ for (FileResultEntity row : rows) {
+ if (Integer.valueOf(RESULT_SUCCESS).equals(row.getSuccess())) {
+ // 内存对象同步改:下方 buildSnapshotFromDb/updateTaskStatusFromRows 直接读这些对象
+ row.setResultFilename(filename);
+ row.setResultFileUrl(null);
+ row.setResultFileSize(0L);
+ row.setResultContentType(CONTENT_TYPE_XLSX);
+ row.setRowCount(rowCount);
+ successRowIds.add(row.getId());
+ if (firstSuccessRow == null) {
+ firstSuccessRow = row;
}
}
+ }
if (firstSuccessRow != null) {
+ updateResultFileColumns(successRowIds, filename, null, 0L, rowCount);
taskFileJobService.enqueueAssembleResult(task.getId(), MODULE_TYPE, firstSuccessRow.getId(), ownerScopeKey(task.getId()));
snapshots = buildSnapshotFromDb(task, rows);
updateTaskStatusFromRows(task, rows);
@@ -2644,7 +2685,8 @@ public class ShopDataCrawlTaskService {
private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) {
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId);
if (lockHandle == null) {
- throw new BusinessException(40901, "任务正在处理中,请稍后再试");
+ log.warn("[shop-data-crawl] 任务锁竞争,拒绝本次提交 taskId={}", taskId);
+ throw new BusinessException(BusinessCodes.TASK_BUSY, "任务正在处理中,请稍后再试");
}
return lockHandle;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/controller/ShopMatchController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/controller/ShopMatchController.java
index ebab808a..8335c24f 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/controller/ShopMatchController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/controller/ShopMatchController.java
@@ -44,6 +44,7 @@ import java.nio.charset.StandardCharsets;
import java.util.List;
import com.nanri.aiimage.modules.task.model.dto.TaskProgressLightRequest;
import com.nanri.aiimage.modules.task.model.vo.TaskProgressLightBatchVo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
@RestController
@RequiredArgsConstructor
@@ -54,6 +55,7 @@ public class ShopMatchController {
private final ShopMatchResolveService shopMatchResolveService;
private final ShopMatchTaskService shopMatchTaskService;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
@GetMapping("/candidates")
@Operation(summary = "查询候选店铺", description = "返回当前用户在匹配店铺模块下保存的候选店铺列表。")
@@ -184,7 +186,9 @@ public class ShopMatchController {
@PostMapping("/tasks/batch")
@Operation(summary = "批量查询任务详情", description = "供前端轮询任务状态与阶段进度使用。")
public ApiResponse tasksBatch(@Valid @RequestBody ProductRiskTaskBatchRequest request) {
- return ApiResponse.success(shopMatchTaskService.getTaskDetailsBatch(request.getTaskIds()));
+ // 归属过滤:/tasks/batch 同样按 userId 过滤(此前无校验,遍历 taskId 即可读他人任务详情)
+ return ApiResponse.success(shopMatchTaskService.getTaskDetailsBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/batch")
@@ -198,7 +202,8 @@ public class ShopMatchController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileError/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody TaskProgressLightRequest request) {
- return ApiResponse.success(shopMatchTaskService.progressLight(request.getTaskIds()));
+ // 透传 userId:此前传 null = 不过滤归属,匿名遍历 taskId 即可读他人任务状态(2026-09 审查)
+ return ApiResponse.success(shopMatchTaskService.progressLight(request.getTaskIds(), request.getUserId()));
}
@PostMapping("/tasks/{taskId}/result")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/service/ShopMatchTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/service/ShopMatchTaskService.java
index 6c86991d..71559276 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/service/ShopMatchTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopmatch/service/ShopMatchTaskService.java
@@ -4,6 +4,7 @@ import cn.hutool.core.io.FileUtil;
import cn.hutool.core.util.IdUtil;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.nanri.aiimage.common.exception.BusinessCodes;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.config.TaskPressureProperties;
import com.nanri.aiimage.modules.file.service.oss.OssStorageService;
@@ -65,8 +66,8 @@ public class ShopMatchTaskService {
private static final String MODULE_TYPE = "SHOP_MATCH";
- public TaskProgressLightBatchVo progressLight(List taskIds) {
- return taskProgressLightAssembler.assemble(MODULE_TYPE, null, taskIds);
+ public TaskProgressLightBatchVo progressLight(List taskIds, Long userId) {
+ return taskProgressLightAssembler.assemble(MODULE_TYPE, userId, taskIds);
}
private static final String CONTENT_TYPE_XLSX = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet";
private static final ZoneId BUSINESS_ZONE = ZoneId.of("Asia/Shanghai");
@@ -590,7 +591,24 @@ public class ShopMatchTaskService {
}
+ /**
+ * 提交门店匹配结果。
+ *
+ * 与 {@code tryFinalizeTask} 共用同一把任务分布式锁(2026-09 全维度审查补):
+ * 此前提交路径不取锁,客户端重试提交、或提交与收尾并发时会出现二次组装与
+ * {@code mergeShopPayload} 读-改-写互相覆盖(后写者丢掉已合并的店铺行,
+ * 且终态判断读的是最长 60s 的缓存)。
+ */
public void submitResult(Long taskId, ShopMatchSubmitResultRequest request) {
+ if (taskId == null || taskId <= 0) {
+ throw new BusinessException("taskId 不合法");
+ }
+ try (TaskDistributedLockService.LockHandle ignored = acquireTaskLockOrThrow(taskId)) {
+ submitResultLocked(taskId, request);
+ }
+ }
+
+ private void submitResultLocked(Long taskId, ShopMatchSubmitResultRequest request) {
if (taskId == null || taskId <= 0) {
throw new BusinessException("taskId 不合法");
}
@@ -932,7 +950,8 @@ public class ShopMatchTaskService {
private TaskDistributedLockService.LockHandle acquireTaskLockOrThrow(Long taskId) {
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId);
if (lockHandle == null) {
- throw new BusinessException(40901, "任务正在处理中,请稍后再试");
+ log.warn("[shop-match] 任务锁竞争,拒绝本次提交 taskId={}", taskId);
+ throw new BusinessException(BusinessCodes.TASK_BUSY, "任务正在处理中,请稍后再试");
}
return lockHandle;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/controller/SimilarAsinController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/controller/SimilarAsinController.java
index e6ed6172..215a162a 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/controller/SimilarAsinController.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/controller/SimilarAsinController.java
@@ -15,6 +15,7 @@ import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinTaskBatchVo;
import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinTaskLightBatchVo;
import com.nanri.aiimage.modules.similarasin.service.SimilarAsinTaskService;
import com.nanri.aiimage.modules.similarasin.service.SimilarAsinTaskService.ResultDownloadInfo;
+import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.Parameter;
import io.swagger.v3.oas.annotations.tags.Tag;
@@ -48,6 +49,7 @@ import java.util.List;
public class SimilarAsinController {
private final SimilarAsinTaskService service;
+ private final TaskProgressOwnershipSupport progressOwnershipSupport;
@GetMapping("/filter-conditions")
@Operation(summary = "查询货源查询筛选条件")
@@ -118,7 +120,9 @@ public class SimilarAsinController {
@PostMapping("/tasks/progress/batch")
@Operation(summary = "批量查询任务进度", description = "前端仅对活跃任务轮询该接口,返回轻量任务状态。")
public ApiResponse progress(@Valid @RequestBody SimilarAsinTaskBatchRequest request) {
- return ApiResponse.success(service.progressBatch(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(service.progressBatch(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/progress/light")
@@ -126,7 +130,9 @@ public class SimilarAsinController {
description = "与 progress/batch 同输入,只返回轻量字段(taskId/status/statusCode/fileStatus/fileReady/updatedAt),"
+ "不查明细/result 行,响应体更小;旧 progress/batch 端点保留不动。")
public ApiResponse progressLight(@Valid @RequestBody SimilarAsinTaskLightRequest request) {
- return ApiResponse.success(service.progressLight(request.getTaskIds()));
+ // 归属过滤:只查询属于该用户的任务(userId 未传时保持原行为,兼容未升级调用方)
+ return ApiResponse.success(service.progressLight(progressOwnershipSupport
+ .filterOwnedTaskIds(request.getTaskIds(), request.getUserId())));
}
@PostMapping("/tasks/{taskId}/activate")
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinTaskBatchRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinTaskBatchRequest.java
index 98f1ceea..850af146 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinTaskBatchRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinTaskBatchRequest.java
@@ -12,4 +12,8 @@ public class SimilarAsinTaskBatchRequest {
@NotEmpty
@Schema(description = "需要查询进度的任务 ID 列表。前端只传正在轮询的活跃任务;后端会批量查询,避免每个任务单独请求。", example = "[3938,3939]", requiredMode = Schema.RequiredMode.REQUIRED)
private List taskIds;
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ @Schema(description = "当前用户 ID(可省略;传入时仅返回该用户的任务)", example = "1")
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinTaskLightRequest.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinTaskLightRequest.java
index 697f2202..37bcf82d 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinTaskLightRequest.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/model/dto/SimilarAsinTaskLightRequest.java
@@ -16,4 +16,8 @@ public class SimilarAsinTaskLightRequest {
@NotNull
@Schema(description = "需要查询进度的任务 ID 列表。前端只传正在轮询的活跃任务;后端会批量查询。", example = "[3938,3939]", requiredMode = Schema.RequiredMode.REQUIRED)
private List taskIds;
+
+ /** 当前用户 ID(可省略;传入时仅返回该用户的任务,用于进度接口归属过滤)。 */
+ @Schema(description = "当前用户 ID(可省略;传入时仅返回该用户的任务)", example = "1")
+ private Long userId;
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinImagePrefetchService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinImagePrefetchService.java
index c492aaa1..0abe6953 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinImagePrefetchService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinImagePrefetchService.java
@@ -72,8 +72,9 @@ public class SimilarAsinImagePrefetchService {
private final ConcurrentHashMap> pendingUrlsByTask = new ConcurrentHashMap<>();
private final Object inflightMonitor = new Object();
- private final ExecutorService prefetchPool = Executors.newFixedThreadPool(PREFETCH_POOL_SIZE,
- namedFactory("similar-asin-prefetch"));
+ // 有界队列线程池:newFixedThreadPool 用的是无界队列,预取任务堆积时不会拒绝、会把内存吃满
+ private final ExecutorService prefetchPool = com.nanri.aiimage.common.util.ThreadPools
+ .boundedFixed("similar-asin-prefetch", PREFETCH_POOL_SIZE);
/**
* Task 15:last_used_at 异步批量刷新的内存缓冲(按 url_hash 去重)。
@@ -90,10 +91,21 @@ public class SimilarAsinImagePrefetchService {
@PostConstruct
public void startTouchFlushScheduler() {
- touchFlushScheduler.scheduleWithFixedDelay(this::flushPendingTouches,
+ // 任务体必须自行吞异常:ScheduledExecutorService 的周期任务一旦抛出,
+ // 后续调度会被永久取消(JDK 语义),30s 兜底刷新会静默停摆直到进程重启
+ touchFlushScheduler.scheduleWithFixedDelay(this::flushPendingTouchesSafely,
TOUCH_FLUSH_INTERVAL_SECONDS, TOUCH_FLUSH_INTERVAL_SECONDS, TimeUnit.SECONDS);
}
+ /** 定时刷新的安全包装:任何异常只记日志,保证下一轮仍会调度。 */
+ private void flushPendingTouchesSafely() {
+ try {
+ flushPendingTouches();
+ } catch (Exception ex) {
+ log.warn("[similar-asin] 定时刷新图片缓存 last_used_at 失败,下一轮重试: {}", ex.getMessage(), ex);
+ }
+ }
+
@PreDestroy
public void shutdown() {
prefetchPool.shutdownNow();
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java
index ccb95b26..6d646c61 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java
@@ -296,7 +296,9 @@ public class SimilarAsinTaskService implements SimilarAsinPipelineHost {
* workbook has its own SXSSF structures and image spool, so four workers
* multiply the peak even though the task itself is a single job.
*/
- private final ExecutorService assembleExecutor = Executors.newFixedThreadPool(2, namedThreadFactory("similar-asin-assemble"));
+ // 有界队列线程池:newFixedThreadPool 用的是无界队列,组装任务堆积时不会拒绝、会把内存吃满
+ private final ExecutorService assembleExecutor = com.nanri.aiimage.common.util.ThreadPools
+ .boundedFixed("similar-asin-assemble", 2);
/** 单元测试清理入口:关闭 assemble 固定线程池,避免测试进程残留线程。 */
public void shutdownAssembleExecutor() {
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/SimilarAsinImageEmbedder.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/SimilarAsinImageEmbedder.java
index 003902cd..c05ce326 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/SimilarAsinImageEmbedder.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/SimilarAsinImageEmbedder.java
@@ -175,7 +175,8 @@ public class SimilarAsinImageEmbedder {
.followSslRedirects(true)
.dns(new SafeDns(Dns.SYSTEM))
.build();
- this.downloadPool = Executors.newFixedThreadPool(downloadPoolSize, namedFactory("similar-asin-image-dl"));
+ this.downloadPool = com.nanri.aiimage.common.util.ThreadPools
+ .boundedFixed("similar-asin-image-dl", downloadPoolSize);
}
@PreDestroy
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/StaleTaskRepairService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/StaleTaskRepairService.java
index ec515a29..48cf0587 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/StaleTaskRepairService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/StaleTaskRepairService.java
@@ -92,14 +92,28 @@ public class StaleTaskRepairService {
/** RUNNING 但心跳超过 2h 未刷新 → FAILED(心跳残留续命即僵尸)。 */
private void repairFileTaskStaleRunning() {
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(STALE_IDLE_MINUTES);
- int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper()
+ // 先查 id 再按主键更新:此前是一条无 LIMIT、无 module_type 的批量 UPDATE,
+ // 而 biz_file_task 没有以 status 打头的索引 → 每 10 分钟一次全表扫,
+ // 且按无索引条件命中行加 X 锁(长事务/锁扩散)。与另两个修复方法口径保持一致。
+ List stale = fileTaskMapper.selectList(new LambdaQueryWrapper()
+ .select(FileTaskEntity::getId)
.eq(FileTaskEntity::getStatus, STATUS_RUNNING)
.lt(FileTaskEntity::getUpdatedAt, cutoff)
+ .last("limit 500"));
+ if (stale.isEmpty()) {
+ return;
+ }
+ String reason = "任务心跳超时,已自动失败";
+ LocalDateTime now = LocalDateTime.now();
+ int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper()
+ .in(FileTaskEntity::getId, stale.stream().map(FileTaskEntity::getId).toList())
+ .eq(FileTaskEntity::getStatus, STATUS_RUNNING)
.set(FileTaskEntity::getStatus, STATUS_FAILED)
- .set(FileTaskEntity::getErrorMessage, "任务心跳超时,已自动失败")
- .set(FileTaskEntity::getFinishedAt, LocalDateTime.now()));
+ .set(FileTaskEntity::getErrorMessage, reason)
+ .set(FileTaskEntity::getFinishedAt, now));
if (updated > 0) {
- log.warn("[stale-task-repair] 心跳超时 RUNNING 已标失败 count={}", updated);
+ log.warn("[stale-task-repair] 心跳超时 RUNNING 已标失败 count={} ids={}",
+ updated, stale.stream().map(FileTaskEntity::getId).toList());
}
}
diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskDistributedLockService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskDistributedLockService.java
index 7416613b..a732e76b 100644
--- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskDistributedLockService.java
+++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskDistributedLockService.java
@@ -23,11 +23,18 @@ public class TaskDistributedLockService {
public static final Duration DEFAULT_LOCK_TTL = Duration.ofMinutes(5);
public static final long DEFAULT_WAIT_MILLIS = 10000L;
private static final long RETRY_DELAY_MILLIS = 200L;
+ /** 续期失败后的重试次数:容忍 Redis 抖动/主从切换造成的瞬时失败。 */
+ private static final int RENEW_MAX_ATTEMPTS = 3;
+ private static final long RENEW_RETRY_DELAY_MILLIS = 200L;
private final DistributedJobLockService distributedJobLockService;
private final ThreadLocal