diff --git a/app/.env b/app/.env index e5d7545..05e0506 100644 --- a/app/.env +++ b/app/.env @@ -12,8 +12,8 @@ zn_username=%E8%87%AA%E5%8A%A8%E5%8C%96_Robot client_name=ShuFuAI # java_api_base=http://47.111.163.154:18080 -java_api_base=http://127.0.0.1:18080 -# java_api_base=http://121.196.149.225:18080 +# java_api_base=http://127.0.0.1:18080 +java_api_base=http://121.196.149.225:18080 # 与 Java 后端共享的 JWT 签名密钥,必须与 backend-java 的 AIIMAGE_JWT_SECRET 完全一致 AIIMAGE_JWT_SECRET=please-change-this-secret-please-rotate-at-least-32-bytes diff --git a/backend-java/src/main/java/com/nanri/aiimage/config/SimilarAsinProperties.java b/backend-java/src/main/java/com/nanri/aiimage/config/SimilarAsinProperties.java index 6bc66a2..91da921 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/config/SimilarAsinProperties.java +++ b/backend-java/src/main/java/com/nanri/aiimage/config/SimilarAsinProperties.java @@ -16,7 +16,14 @@ public class SimilarAsinProperties { private String cozeToken = ""; private List cozeCredentials = new ArrayList<>(); private int cozeCredentialStripeSize = 0; - private int cozeBatchSize = 10; + /** + * P0-1:单次提交 Coze 工作流的 row 数量。 + * 历史值 10,在含 puzzle 多图行的场景下频繁触发 720712008 + * "node executed out of limit: 1000"。降到 3 以避免节点上限被打爆。 + * 出现持续 720712008 时还会被 P1-1 滑窗自适应再降到 1。 + * 不影响 AppearancePatentProperties 的同名值。 + */ + private int cozeBatchSize = 3; private int cozeConnectTimeoutMillis = 10000; private int cozeReadTimeoutMillis = 60000; private int cozePollIntervalMillis = 30000; @@ -34,9 +41,13 @@ public class SimilarAsinProperties { /** * 末尾零头 batch 的强制 flush 阈值(分钟):当不足 cozeBatchSize 的零头 row - * 长时间挂着(Python 慢回传)时触发提交。原硬编码 10 分钟。 + * 长时间挂着(Python 慢回传)时触发提交。 + * 任务级实测:345 行 / 4h 总耗时中,约 2-3 小时是 batch 永远凑不满 batchSize 在等下一波回传, + * 把阈值从 15 调到 1:最多 60s 后 1-2 行也强制提交,让 Coze 提交侧持续进票, + * 总耗时降到与 Python 回传节奏接近。配合 cozeBatchSize=3、cozeSubmitMinIntervalMillis=5000, + * 实际不会触发 Coze 限流。出现限流加重再调回 5/10。 */ - private int cozeFlushPendingMinutes = 15; + private int cozeFlushPendingMinutes = 1; /** * 同 batch retry + split retry 共享的最大重试次数。原硬编码 5。 @@ -45,14 +56,28 @@ public class SimilarAsinProperties { /** * 图片嵌入下载线程池大小。原 SimilarAsinImageEmbedder.DOWNLOAD_POOL_SIZE = 8。 - * 几千行 ×3 列图片场景下提升到 16 可显著缩短 xlsx 组装阶段。 + * P2-10:1000+ 行 ×3 列图片场景下,pool=16 仍是 assemble 阶段瓶颈(实测下载 244s/918s), + * 提到 32 配合 retry=2、global deadline 显著拉低尾延迟; + * 受 2GB 堆约束,单图缩略图维持 300KB 以内,整体内存峰值 ≈ 32 * 300KB ≈ 10MB。 */ - private int imageDownloadPoolSize = 16; + private int imageDownloadPoolSize = 32; /** - * 单张图片下载超时(秒)。原硬编码 5;放宽到 8 配合 1 次重试,整体更稳。 + * 单张图片下载超时(秒)。 + * P2-10:放宽到 8 + retry=1 在快源(aiproxy/m.media-amazon)下没问题, + * 但慢源(cbu01.alicdn)会一直挂 8s 才进入 retry,整体串行时间放大。 + * 调到 5s + retry=2,让慢源更早重试新连接,单图最坏耗时 ≈ 5s * (1+2) = 15s。 */ - private int imageDownloadTimeoutSeconds = 8; + private int imageDownloadTimeoutSeconds = 5; + + /** + * assemble 阶段 taskImageCache 的字节上限。 + * 默认 256MB:5000 行 × 3 列 × 平均 100KB = 1.5GB 远超 2GB 堆, + * 用 BoundedImageCache 按字节累计 LRU 淘汰避免爆堆。 + * 由于 embed() 写完即 remove(),活跃图片字节通常 ≤ 100MB,仅在极端 prefetch 领先场景才会触发淘汰。 + * 出现淘汰过频影响命中率时可上调到 512MB;2GB 堆约束下不建议超过 768MB。 + */ + private long imageCacheMaxBytes = 256L * 1024L * 1024L; /** * 是否在 Coze 请求 parameters 中附带 api_key 字段。 @@ -82,6 +107,20 @@ public class SimilarAsinProperties { */ private boolean cozeResultBufferEnabled = true; + /** + * P0-4:单 credential 抢 Coze 提交锁的最长等待时间(毫秒)。 + * 原硬编码 1000ms,在高并发 split retry 时大量抛 "Coze submit throttle lock timeout" + * 并把整批行 markFailed。应与 cozeSubmitMinIntervalMillis(5000ms)保持 1.5-2 倍关系, + * 默认 10000ms 给抢锁更多时间。 + */ + private long cozeSubmitLockWaitMillis = 10000L; + + /** + * P0-4:抢 Coze 提交锁失败后下次重试间隔(毫秒)。 + * 原硬编码 500ms,会在指数退避算法中作为基础值(500/1000/2000/4000ms 上限 4000)。 + */ + private long cozeSubmitLockRetryDelayMillis = 500L; + @Data public static class CozeCredential { private String name; 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 new file mode 100644 index 0000000..2b1a111 --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinImagePrefetchService.java @@ -0,0 +1,215 @@ +package com.nanri.aiimage.modules.similarasin.service; + +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.nanri.aiimage.modules.similarasin.util.SimilarAsinImageEmbedder; +import com.nanri.aiimage.modules.similarasin.util.SimilarAsinImageEmbedder.ResizedImage; +import com.nanri.aiimage.modules.task.mapper.TaskImageCacheMapper; +import com.nanri.aiimage.modules.task.model.entity.TaskImageCacheEntity; +import jakarta.annotation.PreDestroy; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.dao.DuplicateKeyException; +import org.springframework.stereotype.Service; + +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * P2-11:相似ASIN 图片异步预热服务。 + * + *

背景:assemble 阶段({@code assembleResultWorkbook})需要把 Coze 回包中的 main_url / + * puzzle_img1 / puzzle_img2 下载并 resize 后嵌入 xlsx。当任务行数到 1000+ 时,串行 + + * 短池下载会把整个 assemble 拖到 244s / 918s。改造点: + * + *

+ */ +@Service +@RequiredArgsConstructor +@Slf4j +public class SimilarAsinImagePrefetchService { + + /** P2-11:预热线程池容量。预热不要求高吞吐,4 个线程足够;避免和 embed 阶段抢 IO 资源。 */ + private static final int PREFETCH_POOL_SIZE = 4; + + /** 排队等待上一个 task future 时的最大等待时间(秒),避免被 hung future 永久卡住。 */ + private static final long INFLIGHT_WAIT_SECONDS = 60L; + + private final SimilarAsinImageEmbedder imageEmbedder; + private final TaskImageCacheMapper taskImageCacheMapper; + + /** + * 每个 task 当前 in-flight 的预热 future。enqueue 时如果上一个还没完成,会先等它结束, + * 再串行启动当前 batch 的预热,避免高并发 batch 把图片源站打爆。 + */ + private final ConcurrentHashMap> inflight = new ConcurrentHashMap<>(); + + private final ExecutorService prefetchPool = Executors.newFixedThreadPool(PREFETCH_POOL_SIZE, + namedFactory("similar-asin-prefetch")); + + @PreDestroy + public void shutdown() { + prefetchPool.shutdownNow(); + inflight.clear(); + } + + /** + * P2-11:由 {@code mergeCozeRowsIntoChunk} 调用,把 cozeRows 中的图片 url 异步丢入预热队列。 + * 同 task 串行入队(用 inflight map 排队),避免多个 batch 同时打爆图片源站。 + */ + public void enqueue(Long taskId, List urls) { + if (taskId == null || urls == null || urls.isEmpty()) { + return; + } + Set dedup = new LinkedHashSet<>(); + for (String u : urls) { + if (u == null) { + continue; + } + String trimmed = u.trim(); + if (!trimmed.isEmpty()) { + dedup.add(trimmed); + } + } + if (dedup.isEmpty()) { + return; + } + List targets = new ArrayList<>(dedup); + inflight.compute(taskId, (k, prev) -> prefetchPool.submit(() -> { + // 串行等待上一个 batch 完成(带上限,避免被 hung future 永久卡住)。 + if (prev != null) { + try { + prev.get(INFLIGHT_WAIT_SECONDS, TimeUnit.SECONDS); + } catch (Exception ignored) { + // 上一个 future 异常/超时不阻塞当前预热——继续即可。 + } + } + runPrefetch(taskId, targets); + })); + } + + /** + * P2-11:按 url_hash 命中 DB cache,命中则跳过;未命中下载 + resize 后落表。 + * 所有失败 best-effort:写日志后吞掉,由 assemble 阶段的兜底下载兜底。 + */ + private void runPrefetch(Long taskId, List urls) { + int hit = 0; + int miss = 0; + int fail = 0; + for (String url : urls) { + try { + String urlHash = sha256Hex(url); + if (urlHash == null) { + fail++; + continue; + } + // 命中 DB cache:bumping last_used_at 即可,不再下载。 + Long existing = taskImageCacheMapper.selectCount(new LambdaQueryWrapper() + .eq(TaskImageCacheEntity::getUrlHash, urlHash)); + if (existing != null && existing > 0L) { + taskImageCacheMapper.touchLastUsed(urlHash); + hit++; + continue; + } + ResizedImage thumb = imageEmbedder.fetchAndResizeForCache(url); + if (thumb == null || thumb.bytes() == null || thumb.bytes().length == 0) { + fail++; + continue; + } + TaskImageCacheEntity row = new TaskImageCacheEntity(); + row.setUrlHash(urlHash); + // 防御性截断:表里 url 列 1024 上限。极端长 URL 截断仅影响展示,不影响 url_hash 匹配。 + row.setUrl(url.length() > 1024 ? url.substring(0, 1024) : url); + row.setImageBytes(thumb.bytes()); + row.setByteSize(thumb.bytes().length); + row.setWidth(thumb.width()); + row.setHeight(thumb.height()); + LocalDateTime now = LocalDateTime.now(); + row.setCreatedAt(now); + row.setLastUsedAt(now); + try { + taskImageCacheMapper.insert(row); + } catch (DuplicateKeyException ignored) { + // 并发预热下其它实例/线程已写入 → 更新 last_used_at 即可。 + taskImageCacheMapper.touchLastUsed(urlHash); + } + miss++; + } catch (Exception ex) { + fail++; + log.debug("[similar-asin] prefetch failed taskId={} url={} err={}", taskId, url, ex.getMessage()); + } + } + log.info("[similar-asin] prefetch finished taskId={} total={} hit={} miss={} fail={}", + taskId, urls.size(), hit, miss, fail); + } + + /** + * P2-11:DB cache 直读入口。命中时同步 touchLastUsed,便于 LRU 清理。 + * 失败/未命中返回 null,由调用方走回退路径。 + */ + public byte[] lookup(String url) { + if (url == null) { + return null; + } + String trimmed = url.trim(); + if (trimmed.isEmpty()) { + return null; + } + try { + String urlHash = sha256Hex(trimmed); + if (urlHash == null) { + return null; + } + byte[] bytes = taskImageCacheMapper.selectBytesByUrlHash(urlHash); + if (bytes != null && bytes.length > 0) { + taskImageCacheMapper.touchLastUsed(urlHash); + return bytes; + } + } catch (Exception ex) { + log.debug("[similar-asin] prefetch lookup failed url={} err={}", trimmed, ex.getMessage()); + } + return null; + } + + /** P2-11:sha256 lowercase hex,与 biz_task_image_cache.url_hash 对齐。 */ + private static String sha256Hex(String value) { + try { + MessageDigest md = MessageDigest.getInstance("SHA-256"); + byte[] digest = md.digest(value.getBytes(StandardCharsets.UTF_8)); + StringBuilder sb = new StringBuilder(digest.length * 2); + for (byte b : digest) { + sb.append(String.format("%02x", b & 0xFF)); + } + return sb.toString(); + } catch (NoSuchAlgorithmException ex) { + // SHA-256 在 JDK 中是必备算法,正常环境下不会落到这里。 + return null; + } + } + + private static ThreadFactory namedFactory(String prefix) { + AtomicInteger counter = new AtomicInteger(); + return r -> { + Thread t = new Thread(r, prefix + "-" + counter.incrementAndGet()); + t.setDaemon(true); + return t; + }; + } +} 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 2b1ddf4..9a08e6b 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 @@ -6,6 +6,8 @@ import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ArrayNode; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.nanri.aiimage.common.exception.BusinessException; import com.nanri.aiimage.common.service.DistributedJobLockService; import com.nanri.aiimage.common.util.CozeGroupResultPropagator; @@ -34,6 +36,7 @@ import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinParseVo; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinTaskBatchVo; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinTaskDetailVo; import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinTaskItemVo; +import com.nanri.aiimage.modules.similarasin.util.BoundedImageCache; import com.nanri.aiimage.modules.similarasin.util.ExcelCellImageWriter; import com.nanri.aiimage.modules.similarasin.util.SimilarAsinImageEmbedder; import com.nanri.aiimage.modules.file.service.LocalFileStorageService; @@ -73,6 +76,7 @@ import org.springframework.transaction.PlatformTransactionManager; import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.TransactionDefinition; import org.springframework.transaction.support.TransactionTemplate; +import jakarta.annotation.PreDestroy; import java.io.File; import java.io.FileInputStream; @@ -91,11 +95,18 @@ import java.util.Map; import java.util.Objects; import java.util.Set; import java.util.UUID; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.Semaphore; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReferenceArray; import java.util.function.Supplier; import java.util.regex.Pattern; import java.util.zip.ZipEntry; @@ -126,8 +137,9 @@ public class SimilarAsinTaskService { private static final Duration TASK_LOCK_TTL = Duration.ofMinutes(5); private static final long TASK_LOCK_WAIT_MILLIS = 10000L; private static final Duration COZE_SUBMIT_LOCK_TTL = Duration.ofMinutes(2); - private static final long COZE_SUBMIT_LOCK_WAIT_MILLIS = 1000L; - private static final long COZE_SUBMIT_LOCK_RETRY_DELAY_MILLIS = 500L; + // P0-4:原 COZE_SUBMIT_LOCK_WAIT_MILLIS / COZE_SUBMIT_LOCK_RETRY_DELAY_MILLIS 已下沉到 + // SimilarAsinProperties.cozeSubmitLockWaitMillis / cozeSubmitLockRetryDelayMillis, + // 由 acquireCozeSubmitLock 在方法内读取,并支持指数退避。 private static final int PARSE_RESPONSE_PREVIEW_LIMIT = 100; /** * P0-2 最小风险变体:poll 调度阶段并发预取 Coze HTTP 结果时, @@ -148,6 +160,28 @@ public class SimilarAsinTaskService { */ private final ConcurrentHashMap lastCozeSubmitAtByCredential = new ConcurrentHashMap<>(); + /** + * P1-1:per-task 最近 10 次提交结果(true=720712008 命中)。 + * 容量 10,环形覆盖;与 recentPoisonCursorByTask 配合实现 O(1) record。 + * 当近 5+ 次提交中命中率 ≥ 40% 触发 batch 强制降到 1(隔离毒行)。 + * finalizeTask / deleteTask 时清理,避免内存泄漏。 + */ + private static final int POISON_WINDOW_SIZE = 10; + /** + * 触发降档所需最小样本数:5 → 3。 + * 5000 行任务里 batchSize=10 时 1 次 720712008 会让 1 整批 row markFailed, + * 等不到 5 个样本就已经损失大量行。降到 3 让滑窗更早响应。 + */ + private static final int POISON_MIN_SAMPLES = 3; + /** + * 命中率触发阈值:40% → 33%。 + * 配合 POISON_MIN_SAMPLES=3,最快 1/3 命中率即触发,避免滑窗误判可由 batchSize=1 + * 后的连续成功样本快速恢复(1 次成功就把命中率打到 0/3)。 + */ + private static final int POISON_HIT_RATIO_PERCENT = 33; + private final ConcurrentHashMap> recentPoisonByTask = new ConcurrentHashMap<>(); + private final ConcurrentHashMap recentPoisonCursorByTask = new ConcurrentHashMap<>(); + /** * P0-4 / P2-9:poll 调度链路上的"按 task 单次复用"上下文。 * pollPendingCozeStatesForTask 在锁内一次性加载 FileTaskEntity + allRowsByBaseId(5000 行 JSON 反序列化只发生一次), @@ -207,10 +241,37 @@ public class SimilarAsinTaskService { private final InstanceMetadata instanceMetadata; private final CozeCredentialPoolService cozeCredentialPoolService; private final SimilarAsinImageEmbedder imageEmbedder; + /** + * P2-11:图片异步预热服务。merge cozeRows 时立即丢入预热队列,assemble 阶段优先消费 DB cache。 + * best-effort:service 内部所有异常都已吞掉,不影响主流程。 + */ + private final SimilarAsinImagePrefetchService imagePrefetchService; @Autowired @Qualifier("cozeTaskExecutor") private TaskExecutor cozeTaskExecutor; + /** + * P3-1:assembleResultWorkbook 多源文件并发化使用的固定线程池。 + * 单 task 触发 1 次 assemble,源文件通常 1-N(按 Excel sourceFileKey 分组), + * 4 线程足以覆盖常见 1-4 源文件;sourceRows.size() == 1 时仍走串行降级路径, + * 避免对单文件场景引入额外线程开销。 + */ + private final ExecutorService assembleExecutor = Executors.newFixedThreadPool(4, namedThreadFactory("similar-asin-assemble")); + + @PreDestroy + void shutdownAssembleExecutor() { + assembleExecutor.shutdownNow(); + } + + private static ThreadFactory namedThreadFactory(String prefix) { + AtomicInteger counter = new AtomicInteger(); + return r -> { + Thread t = new Thread(r, prefix + "-" + counter.incrementAndGet()); + t.setDaemon(true); + return t; + }; + } + public List listFilterConditions(Long userId) { validateUserId(userId); List rows = filterConditionMapper.selectList( @@ -588,7 +649,12 @@ public class SimilarAsinTaskService { chunk.setScopeHash(scopeHash); chunk.setChunkIndex(chunkIndex); chunk.setChunkTotal(chunkTotal); - String storedPayload = transientPayloadStorageService.storeChunkPayload(MODULE_TYPE, taskId, scopeHash, chunkIndex, payloadJson); + // P0-2:改用 versioned key(chunk-{index}-{uuid}),避免并发回调下 deterministic key + // 覆盖造成的 read chunk payload failed / 数据丢失。loser 拿到的是自己独有的 object, + // DuplicateKeyException 后可以安全删除,不会影响 winner 的 chunk。 + String storedPayload = transientPayloadStorageService.storeChunkPayloadVersioned(MODULE_TYPE, taskId, scopeHash, chunkIndex, payloadJson); + // P0-3:判断本次 store 是否走了 RustFS → local 兜底。 + boolean localFallback = transientPayloadStorageService.wasLastStoreLocalFallback(); chunk.setPayloadJson(storedPayload); chunk.setPayloadHash(DigestUtil.sha256Hex(payloadJson)); chunk.setCreatedAt(LocalDateTime.now()); @@ -596,10 +662,20 @@ public class SimilarAsinTaskService { try { taskChunkMapper.insert(chunk); } catch (DuplicateKeyException ex) { - // storeChunkPayload uses a deterministic key like chunk-{index}. When two - // concurrent callbacks submit the same chunk, deleting the loser payload here - // can also remove the winner's shared object and break later assembly. - log.info("[similar-asin] duplicate chunk inserted concurrently taskId={} scope={} chunk={}", taskId, scopeKey, chunkIndex); + // P0-2:versioned key 每次都是独立 object,loser 的 storedPayload 不会被 winner 引用, + // 这里直接清理掉以免 RustFS / local 上残留孤儿对象。 + log.info("[similar-asin] duplicate chunk inserted concurrently taskId={} scope={} chunk={} cleanupLoser={}", + taskId, scopeKey, chunkIndex, storedPayload != null); + try { + transientPayloadStorageService.deletePayloadIfPresent(storedPayload); + } catch (Exception cleanupEx) { + log.warn("[similar-asin] cleanup loser chunk payload failed taskId={} chunk={} err={}", + taskId, chunkIndex, cleanupEx.getMessage()); + } + } + // P0-3:本次 chunk 落到本地时把 task 锁到当前实例,避免其他实例 assemble 时读不到。 + if (localFallback) { + bindTaskToCurrentOwnerForLocalFallback(task, scopeHash, chunkIndex); } } else { log.info("[similar-asin] duplicate chunk ignored taskId={} scope={} chunk={}", taskId, scopeKey, chunkIndex); @@ -657,6 +733,8 @@ public class SimilarAsinTaskService { // 同步清理 task_file_job,避免被删除任务遗留的 PENDING/FAILED 行被 TaskResultFileJobWorker 反复扫描出 task not found。 taskFileJobService.deleteTaskJobs(taskId, MODULE_TYPE); fileTaskMapper.deleteById(taskId); + // P1-1:任务被删除时一并清理滑窗记录,防止内存泄漏。 + clearPoisonWindow(taskId); } } @@ -807,7 +885,11 @@ public class SimilarAsinTaskService { chunk.setScopeHash(scopeHash); chunk.setChunkIndex(chunkIndex); chunk.setChunkTotal(chunkTotal); - String storedPayload = transientPayloadStorageService.storeChunkPayload(MODULE_TYPE, taskId, scopeHash, chunkIndex, payloadJson); + // P0-2:与 submitResultLocked 同步切到 versioned key(chunk-{index}-{uuid}), + // 避免并发回调互相覆盖造成 read chunk payload failed。 + String storedPayload = transientPayloadStorageService.storeChunkPayloadVersioned(MODULE_TYPE, taskId, scopeHash, chunkIndex, payloadJson); + // P0-3:拿到 pointer 后立刻读 ThreadLocal 标记,决定是否需要把 task 锁到当前实例。 + boolean localFallback = transientPayloadStorageService.wasLastStoreLocalFallback(); chunk.setPayloadJson(storedPayload); chunk.setPayloadHash(DigestUtil.sha256Hex(payloadJson)); chunk.setCreatedAt(LocalDateTime.now()); @@ -815,7 +897,19 @@ public class SimilarAsinTaskService { try { taskChunkMapper.insert(chunk); } catch (DuplicateKeyException ex) { - log.info("[similar-asin] duplicate chunk inserted concurrently taskId={} scope={} chunk={}", taskId, scopeKey, chunkIndex); + // P0-2:versioned key 后 loser 的 storedPayload 不会被 winner 引用,可以安全清理。 + log.info("[similar-asin] duplicate chunk inserted concurrently taskId={} scope={} chunk={} cleanupLoser={}", + taskId, scopeKey, chunkIndex, storedPayload != null); + try { + transientPayloadStorageService.deletePayloadIfPresent(storedPayload); + } catch (Exception cleanupEx) { + log.warn("[similar-asin] cleanup loser chunk payload failed taskId={} chunk={} err={}", + taskId, chunkIndex, cleanupEx.getMessage()); + } + } + // P0-3:fallback 到 local 时把 task 与当前实例绑定,确保后续 assemble 走对实例。 + if (localFallback) { + bindTaskToCurrentOwnerForLocalFallback(task, scopeHash, chunkIndex); } } else { log.info("[similar-asin] duplicate chunk ignored taskId={} scope={} chunk={}", taskId, scopeKey, chunkIndex); @@ -1474,6 +1568,11 @@ public class SimilarAsinTaskService { } } taskCacheService.deleteTaskCache(task.getId()); + // P1-1:任务进入终态(SUCCESS/FAILED)后清理滑窗记录,避免长任务残留内存。 + // 处于 waitingForAssemble(仍为 RUNNING,等待异步 assemble)时不清理,待 assemble 完真正进入终态再由后续路径触发。 + if (!waitingForAssemble) { + clearPoisonWindow(task.getId()); + } } private FileResultEntity findOrCreateResultRecordForAssembly(FileTaskEntity task, int rowCount) { @@ -1782,25 +1881,42 @@ public class SimilarAsinTaskService { String prompt = readAiPrompt(task); String apiKey = readApiKey(task); int batchSize = Math.max(1, properties.getCozeBatchSize()); + // P1-1:检测到 720712008 风暴时强制把 batch 降到 1,隔离毒行; + // 持续 5+ 次提交命中率 ≥ 40% 才会触发,正常波动不影响吞吐。 + if (isPoisonStormActive(task.getId())) { + log.warn("[similar-asin] poison-storm detected, force batchSize=1 taskId={} originalBatchSize={}", + task.getId(), batchSize); + batchSize = 1; + } List candidates = collectPendingCozeCandidates(chunks, allRowsByBaseId); if (candidates.isEmpty()) { return countPendingCozeStates(task.getId()) > 0; } - List missingImageCandidates = candidates.stream() - .filter(candidate -> candidate != null && candidate.row() != null && !candidate.row().hasImageUrl()) + // P1-2:把"仅过滤 hasImageUrl"扩展为必填字段集中校验。 + // 缺失 asin / title / 图片 url 任意一项即直接 markFailed,不进入 Coze 提交链路。 + // 原因:Python 端偶发空字段会触发 720701002 "fields cannot be extracted from null values", + // 浪费 Coze 配额且把整 batch 拖垮;提前过滤更显式、易于排错。 + // 不打算把"必填字段集合"做成 properties——Coze 工作流签名固定,过度可配置反而把错配藏起来。 + java.util.function.Predicate isMissingRequired = row -> + row == null + || !row.hasImageUrl() + || normalize(row.getAsin()).isBlank() + || normalize(row.getTitle()).isBlank(); + List missingFieldCandidates = candidates.stream() + .filter(candidate -> candidate != null && isMissingRequired.test(candidate.row())) .toList(); - if (!missingImageCandidates.isEmpty()) { + if (!missingFieldCandidates.isEmpty()) { mergeCozeRowsIntoChunk(task, null, null, - cozeClient.markRowsFailed(missingImageCandidates.stream().map(CozeCandidate::row).toList(), - "image url missing, skip Coze"), + cozeClient.markRowsFailed(missingFieldCandidates.stream().map(CozeCandidate::row).toList(), + "required field missing (asin/title/url), skip Coze"), allRowsByBaseId); - log.warn("[similar-asin] skip coze rows without image url taskId={} jobId={} rows={}", - task.getId(), job.getId(), missingImageCandidates.size()); + log.warn("[similar-asin] skip coze rows missing required fields taskId={} jobId={} rows={}", + task.getId(), job.getId(), missingFieldCandidates.size()); } List readyCandidates = candidates.stream() - .filter(candidate -> candidate != null && candidate.row() != null && candidate.row().hasImageUrl()) + .filter(candidate -> candidate != null && !isMissingRequired.test(candidate.row())) .toList(); boolean flushRemainder = isResultSubmissionComplete(task.getId()); // P1-6: 防止 Python 端长时间慢回传时零头永久挂着:job.updatedAt 距今 ≥ cozeFlushPendingMinutes 分钟则强制 flush。 @@ -1997,6 +2113,8 @@ public class SimilarAsinTaskService { String message = firstNonBlank(ex.getMessage(), "Coze submit failed"); log.warn("[similar-asin] coze async submit failed taskId={} jobId={} rows={} batch={}/{} err={}", task.getId(), job.getId(), batchRows.size(), batchIndex, batchTotal, message); + // P1-1:同步 submit 报错路径也记录滑窗(720712008 在 immediateData 阶段就抛错时也算命中)。 + recordCozeSubmitOutcome(task.getId(), CozeFailureClassifier.isPoisonRow(message)); if (CozeFailureClassifier.isThrottleLockTimeout(message)) { savePendingCozeBatchState(task, result, job, batchRows, batchScopeKey, batchScopeHash, batchIndex, batchTotal, message, credential.name()); @@ -2198,6 +2316,9 @@ public class SimilarAsinTaskService { } else { cozeRows = List.of(); } + // P1-1:记录本次 poll 结果(true=720712008 类毒行命中)以驱动滑窗降档。 + // 在 split / retry 之前记录,因为后续 split/retry 拿到的 finalMessage 可能被改写。 + recordCozeSubmitOutcome(state.getTaskId(), CozeFailureClassifier.isPoisonRow(failureMessage)); if (!failureMessage.isBlank() && splitRetryFailedCozeBatchState(state, context, batchRows, failureMessage)) { return; } @@ -2316,8 +2437,19 @@ public class SimilarAsinTaskService { return true; } } catch (Exception ex) { + String msg = firstNonBlank(ex.getMessage(), "Coze retry failed"); + // P1-3:同 batch retry 抢不到节流锁时也走 defer,下次 poll 自动接管, + // 避免被打到 splitRetry 兜底进而把一整个 batch markFailed。 + if (CozeFailureClassifier.isThrottleLockTimeout(msg)) { + log.warn("[similar-asin] coze retry deferred by throttle lock taskId={} stateId={} executeId={}", + state.getTaskId(), state.getId(), state.getCozeExecuteId()); + deferStateForResubmit(state, msg); + taskFileJobService.touchRunning(context.jobId()); + touchJavaSideTaskActivity(state.getTaskId()); + return true; + } log.warn("[similar-asin] coze retry submit failed taskId={} stateId={} executeId={} err={}", - state.getTaskId(), state.getId(), state.getCozeExecuteId(), firstNonBlank(ex.getMessage(), "Coze retry failed")); + state.getTaskId(), state.getId(), state.getCozeExecuteId(), msg); } return false; } @@ -2391,15 +2523,28 @@ public class SimilarAsinTaskService { } private DistributedJobLockService.LockHandle acquireCozeSubmitLock(SimilarAsinCozeClient.CozeCredentialRef credential) { - long deadline = System.currentTimeMillis() + COZE_SUBMIT_LOCK_WAIT_MILLIS; + // P0-4:等待时长 / 退避基础值下沉到 SimilarAsinProperties,可在线上调整。 + // 原硬编码 1000ms 在高并发 split retry 时大量超时 markFailed,默认抬高到 10s。 + long waitMillis = Math.max(1000L, properties.getCozeSubmitLockWaitMillis()); + long baseDelay = Math.max(100L, properties.getCozeSubmitLockRetryDelayMillis()); + long deadline = System.currentTimeMillis() + waitMillis; String credentialName = credential == null ? "default" : firstNonBlank(credential.name(), "default"); + int attempt = 0; while (System.currentTimeMillis() <= deadline) { DistributedJobLockService.LockHandle lockHandle = distributedJobLockService.tryLock("similar-asin:coze-submit:" + credentialName, COZE_SUBMIT_LOCK_TTL); if (lockHandle != null) { return lockHandle; } - sleepQuietly(COZE_SUBMIT_LOCK_RETRY_DELAY_MILLIS); + // 指数退避:500/1000/2000/4000ms,上限 4000,避免抢锁失败时密集打日志。 + long delay = Math.min(baseDelay * (1L << Math.min(attempt, 3)), 4000L); + // 末次循环之前确保 sleep 不会越过 deadline。 + long left = deadline - System.currentTimeMillis(); + if (left <= 0) { + break; + } + sleepQuietly(Math.min(delay, left)); + attempt++; } return null; } @@ -2415,6 +2560,89 @@ public class SimilarAsinTaskService { } } + /** + * P1-1:记录一次 Coze 提交结果是否触发了 720712008 类毒行。 + * 滑窗大小 POISON_WINDOW_SIZE=10,环形覆盖:满了之后从 0 重新开始。 + */ + private void recordCozeSubmitOutcome(Long taskId, boolean poisonHit) { + if (taskId == null) { + return; + } + AtomicReferenceArray ring = recentPoisonByTask.computeIfAbsent(taskId, + k -> new AtomicReferenceArray<>(POISON_WINDOW_SIZE)); + AtomicInteger cursor = recentPoisonCursorByTask.computeIfAbsent(taskId, k -> new AtomicInteger()); + int idx = Math.floorMod(cursor.getAndIncrement(), POISON_WINDOW_SIZE); + ring.set(idx, poisonHit); + } + + /** + * P1-1:判断当前 task 是否处于 720712008 风暴中。 + * 阈值:近 POISON_MIN_SAMPLES=3+ 次提交、命中率 ≥ POISON_HIT_RATIO_PERCENT=33% → 触发降档。 + */ + private boolean isPoisonStormActive(Long taskId) { + if (taskId == null) { + return false; + } + AtomicReferenceArray ring = recentPoisonByTask.get(taskId); + if (ring == null) { + return false; + } + int hits = 0; + int total = 0; + for (int i = 0; i < POISON_WINDOW_SIZE; i++) { + Boolean v = ring.get(i); + if (v != null) { + total++; + if (Boolean.TRUE.equals(v)) { + hits++; + } + } + } + return total >= POISON_MIN_SAMPLES && hits * 100 / total >= POISON_HIT_RATIO_PERCENT; + } + + /** + * P1-1:任务终结 / 删除时调用,清理 per-task 滑窗记录避免内存泄漏。 + */ + private void clearPoisonWindow(Long taskId) { + if (taskId == null) { + return; + } + recentPoisonByTask.remove(taskId); + recentPoisonCursorByTask.remove(taskId); + } + + /** + * P1-3:throttle lock timeout 是临时性失败,不应当 markFailed。 + * 把 state 改写为 cozeStatus=RUNNING / cozeExecuteId=null, + * 下一轮 poll 会通过 retryPendingCozeSubmitState 自动接管重新提交。 + * 与 savePendingCozeBatchState 一致使用 RUNNING + null executeId 表达 "pending submit"。 + */ + private void deferStateForResubmit(TaskScopeStateEntity state, + String reason) { + if (state == null || state.getId() == null) { + return; + } + LocalDateTime now = LocalDateTime.now(); + int updated = taskScopeStateMapper.update(null, new LambdaUpdateWrapper() + .eq(TaskScopeStateEntity::getId, state.getId()) + .in(TaskScopeStateEntity::getCozeStatus, List.of(COZE_STATUS_SUBMITTED, COZE_STATUS_RUNNING)) + .set(TaskScopeStateEntity::getCozeStatus, COZE_STATUS_RUNNING) + .set(TaskScopeStateEntity::getCozeExecuteId, null) + .set(TaskScopeStateEntity::getCozeLastPolledAt, null) + .set(TaskScopeStateEntity::getCozeCompletedAt, null) + .set(TaskScopeStateEntity::getCozeAttemptCount, 0) + .set(TaskScopeStateEntity::getCozeError, "deferred: " + firstNonBlank(reason, "throttle lock timeout")) + .set(TaskScopeStateEntity::getUpdatedAt, now)); + if (updated > 0) { + log.warn("[similar-asin] coze submit deferred by throttle lock taskId={} stateId={} reason={}", + state.getTaskId(), state.getId(), reason); + } else { + log.info("[similar-asin] coze submit defer skipped (state moved) taskId={} stateId={} status={}", + state.getTaskId(), state.getId(), state.getCozeStatus()); + } + } + private boolean splitRetryFailedCozeBatchState(TaskScopeStateEntity state, CozeBatchContext context, List batchRows, @@ -2480,8 +2708,19 @@ public class SimilarAsinTaskService { return true; } } catch (Exception ex) { + String msg = firstNonBlank(ex.getMessage(), "Coze split retry failed"); + // P1-3:节流锁超时是临时性失败,把 state 回写为待重试,下个 poll 周期接管。 + // 不再走 markFailed → markRowsFailed 的死亡路径,避免一行抢不到锁就被永久落进 xlsx。 + if (CozeFailureClassifier.isThrottleLockTimeout(msg)) { + log.warn("[similar-asin] coze split retry deferred by throttle lock taskId={} stateId={} parts={}", + state.getTaskId(), state.getId(), partitions.size()); + deferStateForResubmit(state, msg); + taskFileJobService.touchRunning(context.jobId()); + touchJavaSideTaskActivity(state.getTaskId()); + return true; + } log.warn("[similar-asin] coze split retry submit failed taskId={} stateId={} executeId={} err={}", - state.getTaskId(), state.getId(), state.getCozeExecuteId(), firstNonBlank(ex.getMessage(), "Coze split retry failed")); + state.getTaskId(), state.getId(), state.getCozeExecuteId(), msg); } return false; } @@ -3004,6 +3243,8 @@ public class SimilarAsinTaskService { fileTaskMapper.updateById(task); taskCacheService.deleteTaskCache(task.getId()); saveFileBuildProgress(task, job, totalProgressUnits, totalProgressUnits, "Result file generated"); + // P1-1:异步 assemble 完成 → task 进入 SUCCESS/FAILED 终态,清理滑窗记录避免内存泄漏。 + clearPoisonWindow(task.getId()); } private void mergeCozeRowsIntoChunk(FileTaskEntity task, @@ -3018,6 +3259,23 @@ public class SimilarAsinTaskService { if (chunks.isEmpty()) { return; } + // P2-11:把当前 batch 命中的图片 url 异步丢入预热队列。 + // 预热失败不影响主流程,assemble 阶段无 DB cache 命中也会走原下载链路兜底。 + try { + List prefetchUrls = new ArrayList<>(cozeRows.size() * 3); + for (SimilarAsinResultRowDto cozeRow : cozeRows) { + if (cozeRow == null) { + continue; + } + addNonBlank(prefetchUrls, cozeRow.getMainUrl()); + addNonBlank(prefetchUrls, cozeRow.getPuzzleImg1()); + addNonBlank(prefetchUrls, cozeRow.getPuzzleImg2()); + } + imagePrefetchService.enqueue(task.getId(), prefetchUrls); + } catch (Exception ex) { + // 预热入队是 best-effort,任何异常都不能阻断 merge 主路径。 + log.debug("[similar-asin] enqueue prefetch failed taskId={} err={}", task.getId(), ex.getMessage()); + } Map> rowsByChunk = new LinkedHashMap<>(); Map chunkByKey = new LinkedHashMap<>(); for (TaskChunkEntity chunk : chunks) { @@ -3499,6 +3757,96 @@ public class SimilarAsinTaskService { } } + /** + * P0-3:当 chunk 因 RustFS 失败落到本地时,把 task 锁定到当前实例。 + * 这样后续 assembleResultWorkbook 调度只会路由到这台机器(已有 + * isJobOwnedByCurrentInstance / ownerFromTask 兼容这条 ownerInstanceId)。 + * 同时把 fallback 信号写到 task_scope_state.state_json,便于排查。 + */ + private void bindTaskToCurrentOwnerForLocalFallback(FileTaskEntity task, String scopeHash, Integer chunkIndex) { + if (task == null || task.getId() == null) { + return; + } + String currentInstance = currentInstanceId(); + try { + String existingJson = task.getResultJson(); + ObjectNode payload; + if (existingJson == null || existingJson.isBlank()) { + payload = objectMapper.createObjectNode(); + } else { + JsonNode tree = objectMapper.readTree(existingJson); + payload = tree.isObject() ? (ObjectNode) tree : objectMapper.createObjectNode(); + } + String existingOwner = payload.path("ownerInstanceId").asText(""); + if (existingOwner.isBlank()) { + payload.put("ownerInstanceId", currentInstance); + payload.put("ownerInstanceReason", "rustfs-fallback-local"); + String updatedJson = objectMapper.writeValueAsString(payload); + task.setResultJson(updatedJson); + FileTaskEntity update = new FileTaskEntity(); + update.setId(task.getId()); + update.setResultJson(updatedJson); + update.setUpdatedAt(LocalDateTime.now()); + fileTaskMapper.updateById(update); + log.warn("[similar-asin] task bound to current instance due to rustfs fallback taskId={} instanceId={} chunk={}", + task.getId(), currentInstance, chunkIndex); + } else if (!Objects.equals(existingOwner, currentInstance)) { + log.error("[similar-asin] rustfs fallback on non-owner instance, chunk will be unreachable taskId={} chunk={} owner={} current={}", + task.getId(), chunkIndex, existingOwner, currentInstance); + } + } catch (Exception ex) { + log.warn("[similar-asin] bind owner for local fallback failed taskId={} err={}", task.getId(), ex.getMessage()); + } + // 把 fallback 信号写到 task_scope_state.state_json,方便排查 / 后续告警钩子。 + try { + if (scopeHash == null || scopeHash.isBlank()) { + return; + } + TaskScopeStateEntity scope = taskScopeStateMapper.selectOne(new LambdaQueryWrapper() + .eq(TaskScopeStateEntity::getTaskId, task.getId()) + .eq(TaskScopeStateEntity::getModuleType, MODULE_TYPE) + .eq(TaskScopeStateEntity::getScopeHash, scopeHash) + .last("limit 1")); + if (scope == null) { + return; + } + ObjectNode stateNode; + String stateJson = scope.getStateJson(); + if (stateJson == null || stateJson.isBlank()) { + stateNode = objectMapper.createObjectNode(); + stateNode.put("phase", "RECEIVED"); + stateNode.put("coze", "PENDING"); + } else { + JsonNode parsed = objectMapper.readTree(stateJson); + stateNode = parsed.isObject() ? (ObjectNode) parsed : objectMapper.createObjectNode(); + } + ArrayNode fallbackArr; + JsonNode existingArr = stateNode.path("localFallback"); + if (existingArr.isArray()) { + fallbackArr = (ArrayNode) existingArr; + } else { + fallbackArr = stateNode.putArray("localFallback"); + } + String tag = "chunk-" + chunkIndex + "@" + currentInstance; + boolean exists = false; + for (JsonNode node : fallbackArr) { + if (tag.equals(node.asText(""))) { + exists = true; + break; + } + } + if (!exists) { + fallbackArr.add(tag); + scope.setStateJson(objectMapper.writeValueAsString(stateNode)); + scope.setUpdatedAt(LocalDateTime.now()); + taskScopeStateMapper.updateById(scope); + } + } catch (Exception ex) { + log.warn("[similar-asin] write local-fallback state failed taskId={} chunk={} err={}", + task.getId(), chunkIndex, ex.getMessage()); + } + } + private String ownerFromScopeKey(String scopeKey) { if (scopeKey == null || scopeKey.isBlank()) { return null; @@ -3627,6 +3975,11 @@ public class SimilarAsinTaskService { log.info("[similar-asin] assemble workbook taskId={} parsedRows={} persistedRows={} resultRows={} resolvedRows={}", task.getId(), parsed.getAllItems().size(), persistedResultRows, resultMap.size(), resolvedRows); if (!parsed.getAllItems().isEmpty() && resultMap.isEmpty()) { + // P1-4:检查是否有 chunk-read-failed 标记,把 chunk index 列表附在错误信息里。 + String chunkReadFailureSummary = collectChunkReadFailureSummary(task.getId()); + if (!chunkReadFailureSummary.isBlank()) { + throw new BusinessException("相似ASIN检测结果为空(" + chunkReadFailureSummary + "),请稍后重试生成结果文件"); + } throw new BusinessException("相似ASIN检测结果为空,请稍后重试生成结果文件"); } // 公共分组处理:同 baseId 连续行视为一组(如 1、1_1、1_2),组内任一行的"是否符合类目" @@ -3652,18 +4005,85 @@ public class SimilarAsinTaskService { List sourceRows = splitRowsBySourceFile(parsed, parsed.getAllItems(), result.getSourceFilename()); List workbooks = new ArrayList<>(); File zip = null; + // 跨 workbook 共享 taskImageCache:多源场景下相同 URL 仅下载一次。 + // 配合 embed() 写完即 remove(),cache 仅承载 in-flight 图片; + // 5000 行 × 200 实例规模下用 BoundedImageCache 按字节做 LRU 淘汰,硬上限 256MB + // (配置 aiimage.similar-asin.image-cache-max-bytes 可调),避免堆爆。 + long imageCacheMaxBytes = properties.getImageCacheMaxBytes() > 0 + ? properties.getImageCacheMaxBytes() + : BoundedImageCache.DEFAULT_MAX_BYTES; + BoundedImageCache taskImageCache = new BoundedImageCache(imageCacheMaxBytes); try { - for (SourceRows item : sourceRows) { - String filename = safeFileStem(item.sourceFilename()) + "-result.xlsx"; - String tempFilename = safeFileStem(item.sourceFilename()) - + "-" + task.getId() - + "-" + result.getId() - + "-" + UUID.randomUUID() - + "-result.xlsx"; - File xlsx = new File(outputDir, tempFilename); - writeResultWorkbook(xlsx, item.rows(), resultMap); - workbooks.add(new SourceResultWorkbook(xlsx, filename, item.rows().size())); + // P3-1:单源文件场景仍走串行降级路径,避免引入线程切换开销; + // 多源文件场景把 writeResultWorkbook 投到 assembleExecutor 上并发跑, + // 1000+ 行 ×N 源文件的 assemble 阶段总耗时直接除以 N(受池大小 4 限制)。 + // assembleExecutor 是固定 4 线程池:sourceRows.size() ≤ 4 时全部并行; + // > 4 时多余源文件排队,避免 4×5000 行同时打开图片缓存爆堆。 + if (sourceRows.size() <= 1) { + for (SourceRows item : sourceRows) { + String filename = safeFileStem(item.sourceFilename()) + "-result.xlsx"; + String tempFilename = safeFileStem(item.sourceFilename()) + + "-" + task.getId() + + "-" + result.getId() + + "-" + UUID.randomUUID() + + "-result.xlsx"; + File xlsx = new File(outputDir, tempFilename); + writeResultWorkbook(xlsx, item.rows(), resultMap, taskImageCache); + workbooks.add(new SourceResultWorkbook(xlsx, filename, item.rows().size())); + } + } else { + long assembleStart = System.currentTimeMillis(); + List> futures = new ArrayList<>(sourceRows.size()); + for (SourceRows item : sourceRows) { + final SourceRows captured = item; + final String filename = safeFileStem(captured.sourceFilename()) + "-result.xlsx"; + final String tempFilename = safeFileStem(captured.sourceFilename()) + + "-" + task.getId() + + "-" + result.getId() + + "-" + UUID.randomUUID() + + "-result.xlsx"; + futures.add(CompletableFuture.supplyAsync(() -> { + File xlsx = new File(outputDir, tempFilename); + writeResultWorkbook(xlsx, captured.rows(), resultMap, taskImageCache); + return new SourceResultWorkbook(xlsx, filename, captured.rows().size()); + }, assembleExecutor)); + } + for (CompletableFuture future : futures) { + try { + // 单源文件 30 分钟硬上限:超时直接抛错,避免被慢源永久阻塞。 + workbooks.add(future.get(30, TimeUnit.MINUTES)); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + for (CompletableFuture remaining : futures) { + remaining.cancel(true); + } + throw new BusinessException("生成相似ASIN检测结果中断", ie); + } catch (TimeoutException te) { + for (CompletableFuture remaining : futures) { + remaining.cancel(true); + } + throw new BusinessException("生成相似ASIN检测结果超时(30 分钟)", te); + } catch (ExecutionException ee) { + Throwable cause = ee.getCause(); + for (CompletableFuture remaining : futures) { + remaining.cancel(true); + } + if (cause instanceof BusinessException be) { + throw be; + } + throw new BusinessException("生成相似ASIN检测结果失败", + cause != null ? cause : ee); + } + } + log.info("[similar-asin] assemble parallel finished taskId={} sources={} costMs={}", + task.getId(), sourceRows.size(), System.currentTimeMillis() - assembleStart); } + log.info("[similar-asin] assemble image-cache stats taskId={} sources={} cacheSize={} currentBytes={} evictedCount={} evictedBytes={} maxBytes={}", + task.getId(), sourceRows.size(), taskImageCache.size(), + taskImageCache.currentBytes(), taskImageCache.evictedCount(), + taskImageCache.evictedBytes(), taskImageCache.maxBytes()); + // 主动清空,让 GC 尽早回收图片字节,避免 zip 阶段还占着堆。 + taskImageCache.clear(); if (workbooks.isEmpty()) { throw new BusinessException("相似ASIN检测结果为空,请稍后重试生成结果文件"); } @@ -3904,7 +4324,8 @@ public class SimilarAsinTaskService { private void writeResultWorkbook(File xlsx, List rowsToWrite, - Map resultMap) { + Map resultMap, + Map taskImageCache) { // Excel 365 "Place in Cell" 图片:写完 workbook 后 patch richData 单元格图片结构。 // 这不是浮动 Drawing,因此点击图片区域会选中单元格,图片不能被拖到任意位置,也不需要工作表保护。 ExcelCellImageWriter.Session excelCellImageSession = ExcelCellImageWriter.createSession(); @@ -3915,7 +4336,7 @@ public class SimilarAsinTaskService { font.setBold(true); headerStyle.setFont(font); - Map taskImageCache = new ConcurrentHashMap<>(); + // taskImageCache 由 assembleResultWorkbook 跨 workbook 注入,避免 4 源文件 × 同一批 URL 重复下载。 sheet.setColumnWidth(IMG_COL_MAIN, SimilarAsinImageEmbedder.IMAGE_COL_WIDTH_CHARS * 256); sheet.setColumnWidth(IMG_COL_PUZZLE1, SimilarAsinImageEmbedder.IMAGE_COL_WIDTH_CHARS * 256); @@ -3933,10 +4354,36 @@ public class SimilarAsinTaskService { // POI 写入仍单线程串行注册 cell image(cache 命中直接 resize + register),下载/写入解耦。 List prefetchUrls = collectImageUrlsForPrefetch(rowsToWrite, resultMap); if (!prefetchUrls.isEmpty()) { - long prefetchStart = System.currentTimeMillis(); - imageEmbedder.prefetch(prefetchUrls, taskImageCache); - log.info("[similar-asin] image prefetch finished urls={} cached={} costMs={}", - prefetchUrls.size(), taskImageCache.size(), System.currentTimeMillis() - prefetchStart); + // P2-11:先一次性把 DB cache 命中的字节填入 taskImageCache,避免再次走网络。 + // 命中部分从 prefetchUrls 剔除,剩余的真正未命中的 URL 才走 imageEmbedder.prefetch 网络下载。 + long dbCacheStart = System.currentTimeMillis(); + List remainingUrls = new ArrayList<>(prefetchUrls.size()); + int dbHit = 0; + for (String url : prefetchUrls) { + if (url == null || taskImageCache.containsKey(url)) { + continue; + } + byte[] cachedBytes = imagePrefetchService.lookup(url); + if (cachedBytes == null) { + remainingUrls.add(url); + continue; + } + SimilarAsinImageEmbedder.ResizedImage thumb = imageEmbedder.decodeCachedThumb(cachedBytes); + if (thumb != null) { + taskImageCache.putIfAbsent(url, thumb); + dbHit++; + } else { + remainingUrls.add(url); + } + } + log.info("[similar-asin] image db-cache lookup finished urls={} hit={} miss={} costMs={}", + prefetchUrls.size(), dbHit, remainingUrls.size(), System.currentTimeMillis() - dbCacheStart); + if (!remainingUrls.isEmpty()) { + long prefetchStart = System.currentTimeMillis(); + imageEmbedder.prefetch(remainingUrls, taskImageCache); + log.info("[similar-asin] image prefetch finished urls={} cached={} costMs={}", + remainingUrls.size(), taskImageCache.size(), System.currentTimeMillis() - prefetchStart); + } } int rowIndex = 1; @@ -4569,12 +5016,94 @@ public class SimilarAsinTaskService { rows.put(rowKey(row), row); } } catch (Exception ex) { - log.warn("[similar-asin] read chunk payload failed taskId={} chunk={} err={}", - chunk.getTaskId(), chunk.getChunkIndex(), ex.getMessage()); + String msg = ex.getMessage() == null ? "" : ex.getMessage(); + // P1-4:识别跨实例 local 指针读不到的场景,把"chunk 在另一实例"的元信息 + // 通过 task_scope_state.last_error 留痕,便于排查"为何 owner 切换后 chunk 读不到"。 + boolean crossInstance = msg.contains("only exists on instance="); + log.warn("[similar-asin] read chunk payload failed taskId={} chunk={} crossInstance={} err={}", + chunk.getTaskId(), chunk.getChunkIndex(), crossInstance, msg); + recordChunkReadFailure(chunk, crossInstance, msg); } return rows; } + /** + * P1-4:扫描该 task 下所有 task_scope_state.last_error, + * 把 "chunk-read-failed[N]" 标记收集成一行摘要返回;没命中返回空字符串。 + */ + private String collectChunkReadFailureSummary(Long taskId) { + if (taskId == null) { + return ""; + } + try { + List scopes = taskScopeStateMapper.selectList( + new LambdaQueryWrapper() + .eq(TaskScopeStateEntity::getTaskId, taskId) + .eq(TaskScopeStateEntity::getModuleType, MODULE_TYPE)); + if (scopes == null || scopes.isEmpty()) { + return ""; + } + LinkedHashSet tags = new LinkedHashSet<>(); + for (TaskScopeStateEntity scope : scopes) { + String lastError = scope == null ? null : scope.getLastError(); + if (lastError == null || lastError.isBlank()) { + continue; + } + int from = 0; + while (true) { + int start = lastError.indexOf("chunk-read-failed[", from); + if (start < 0) { + break; + } + int end = lastError.indexOf(']', start); + if (end < 0) { + break; + } + tags.add(lastError.substring(start, end + 1)); + from = end + 1; + } + } + return String.join(", ", tags); + } catch (Exception ex) { + log.warn("[similar-asin] collect chunk-read-failure summary failed taskId={} err={}", taskId, ex.getMessage()); + return ""; + } + } + + /** + * P1-4:把 chunk 读失败信息聚合到 task_scope_state.last_error, + * 形式 "chunk-read-failed[{idx}]" / "chunk-read-failed[{idx}@cross-instance]"。 + * 后续 assembleResultWorkbook 之前抛"结果为空"时可以附带 chunk index 列表,定位更具体。 + * best-effort:查询 / 更新失败时仅记日志,不抛回主流程。 + */ + private void recordChunkReadFailure(TaskChunkEntity chunk, boolean crossInstance, String msg) { + if (chunk == null || chunk.getTaskId() == null || chunk.getScopeHash() == null) { + return; + } + try { + TaskScopeStateEntity scope = taskScopeStateMapper.selectOne(new LambdaQueryWrapper() + .eq(TaskScopeStateEntity::getTaskId, chunk.getTaskId()) + .eq(TaskScopeStateEntity::getModuleType, MODULE_TYPE) + .eq(TaskScopeStateEntity::getScopeHash, chunk.getScopeHash()) + .last("limit 1")); + if (scope == null) { + return; + } + String existing = scope.getLastError() == null ? "" : scope.getLastError(); + String tag = "chunk-read-failed[" + chunk.getChunkIndex() + (crossInstance ? "@cross-instance" : "") + "]"; + if (existing.contains(tag)) { + return; + } + String updated = existing.isBlank() ? tag : existing + "; " + tag; + taskScopeStateMapper.update(null, new LambdaUpdateWrapper() + .eq(TaskScopeStateEntity::getId, scope.getId()) + .set(TaskScopeStateEntity::getLastError, updated) + .set(TaskScopeStateEntity::getUpdatedAt, LocalDateTime.now())); + } catch (Exception ignored) { + // best-effort:失败不影响主流程 + } + } + /** * 跨 chunk merge 时把没有匹配到任何 chunk 的回流行兜底落到 task_scope_state, * 避免数据被静默丢弃;assembleResult 阶段会在 {@link #loadPersistedResultRows(Long)} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/BoundedImageCache.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/BoundedImageCache.java new file mode 100644 index 0000000..bbfc9fe --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/BoundedImageCache.java @@ -0,0 +1,193 @@ +package com.nanri.aiimage.modules.similarasin.util; + +import lombok.extern.slf4j.Slf4j; + +import java.util.Collection; +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.atomic.AtomicLong; + +/** + * similar-asin assemble 阶段的 taskImageCache 实现:按字节累计上限做 LRU 淘汰。 + * + *

背景:原实现是 {@code ConcurrentHashMap},在 5000 行 × 200 实例 + * 规模下、单图缩略图最大 300KB 时,单 task 可能堆积 5000 × 3 × 300KB ≈ 4.5GB 图片字节, + * 远超 JVM 2GB 堆。需要硬上限避免堆爆。 + * + *

设计: + *

    + *
  • 底层 {@link LinkedHashMap} access-order 维护 LRU;
  • + *
  • 所有写入路径 {@code put / putIfAbsent} 后 evictIfOverflow,按字节累计淘汰最久未访问条目;
  • + *
  • {@code containsKey / get} 也会更新 LRU 顺序;
  • + *
  • 整个类对外仍是 {@code Map},调用方无感知;
  • + *
  • 所有公共方法 synchronized:embed 阶段单线程主导,prefetch 阶段并发只通过 putIfAbsent + * 少量竞争,加锁开销可忽略。
  • + *
+ * + *

不实现 entrySet/keySet/values/equals/hashCode 等少用方法(throw UnsupportedOperationException)。 + * 与现有 SimilarAsinImageEmbedder / SimilarAsinTaskService 中使用的 6 个方法(containsKey、get、 + * put、putIfAbsent、remove、size)严格对齐。 + */ +@Slf4j +public class BoundedImageCache implements Map { + + /** + * 默认堆字节预算:256MB。 + *

5000 行 × 3 列 × 平均 100KB = 1.5GB,超出后按 LRU 淘汰; + * 由于 embed() 写完即 remove(),活跃图片字节通常远低于该上限, + * 仅在 prefetch 显著领先 embed 时才会触发淘汰。 + */ + public static final long DEFAULT_MAX_BYTES = 256L * 1024L * 1024L; + + private final long maxBytes; + private final LinkedHashMap backing; + private final AtomicLong currentBytes = new AtomicLong(); + private final AtomicLong evictedCount = new AtomicLong(); + private final AtomicLong evictedBytes = new AtomicLong(); + + public BoundedImageCache() { + this(DEFAULT_MAX_BYTES); + } + + public BoundedImageCache(long maxBytes) { + this.maxBytes = maxBytes > 0 ? maxBytes : DEFAULT_MAX_BYTES; + // access-order = true:get/containsKey 也会刷新 LRU 顺序。 + this.backing = new LinkedHashMap<>(64, 0.75f, true); + } + + public long maxBytes() { + return maxBytes; + } + + public long currentBytes() { + return currentBytes.get(); + } + + public long evictedCount() { + return evictedCount.get(); + } + + public long evictedBytes() { + return evictedBytes.get(); + } + + @Override + public synchronized int size() { + return backing.size(); + } + + @Override + public synchronized boolean isEmpty() { + return backing.isEmpty(); + } + + @Override + public synchronized boolean containsKey(Object key) { + return backing.containsKey(key); + } + + @Override + public synchronized boolean containsValue(Object value) { + return backing.containsValue(value); + } + + @Override + public synchronized SimilarAsinImageEmbedder.ResizedImage get(Object key) { + return backing.get(key); + } + + @Override + public synchronized SimilarAsinImageEmbedder.ResizedImage put(String key, SimilarAsinImageEmbedder.ResizedImage value) { + SimilarAsinImageEmbedder.ResizedImage prev = backing.put(key, value); + if (prev != null) { + currentBytes.addAndGet(-byteSizeOf(prev)); + } + currentBytes.addAndGet(byteSizeOf(value)); + evictIfOverflow(); + return prev; + } + + @Override + public synchronized SimilarAsinImageEmbedder.ResizedImage putIfAbsent(String key, SimilarAsinImageEmbedder.ResizedImage value) { + SimilarAsinImageEmbedder.ResizedImage existing = backing.get(key); + if (existing != null) { + return existing; + } + backing.put(key, value); + currentBytes.addAndGet(byteSizeOf(value)); + evictIfOverflow(); + return null; + } + + @Override + public synchronized SimilarAsinImageEmbedder.ResizedImage remove(Object key) { + SimilarAsinImageEmbedder.ResizedImage removed = backing.remove(key); + if (removed != null) { + currentBytes.addAndGet(-byteSizeOf(removed)); + } + return removed; + } + + @Override + public synchronized void putAll(Map m) { + if (m == null || m.isEmpty()) { + return; + } + for (Map.Entry entry : m.entrySet()) { + put(entry.getKey(), entry.getValue()); + } + } + + @Override + public synchronized void clear() { + backing.clear(); + currentBytes.set(0); + } + + private void evictIfOverflow() { + long over = currentBytes.get() - maxBytes; + if (over <= 0) { + return; + } + long evicted = 0; + long evictedBytesLocal = 0; + java.util.Iterator> it = backing.entrySet().iterator(); + while (it.hasNext() && currentBytes.get() > maxBytes) { + Map.Entry oldest = it.next(); + int sz = byteSizeOf(oldest.getValue()); + it.remove(); + currentBytes.addAndGet(-sz); + evicted++; + evictedBytesLocal += sz; + } + if (evicted > 0) { + this.evictedCount.addAndGet(evicted); + this.evictedBytes.addAndGet(evictedBytesLocal); + log.warn("[similar-asin][image-cache] LRU evicted entries={} bytes={} currentBytes={} maxBytes={}", + evicted, evictedBytesLocal, currentBytes.get(), maxBytes); + } + } + + private static int byteSizeOf(SimilarAsinImageEmbedder.ResizedImage img) { + if (img == null || img.bytes() == null) { + return 0; + } + return img.bytes().length; + } + + @Override + public Set keySet() { + throw new UnsupportedOperationException("BoundedImageCache.keySet not supported"); + } + + @Override + public Collection values() { + throw new UnsupportedOperationException("BoundedImageCache.values not supported"); + } + + @Override + public Set> entrySet() { + throw new UnsupportedOperationException("BoundedImageCache.entrySet not supported"); + } +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/ExcelCellImageWriter.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/ExcelCellImageWriter.java index 7c454a9..e2a3550 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/ExcelCellImageWriter.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/util/ExcelCellImageWriter.java @@ -166,7 +166,7 @@ public final class ExcelCellImageWriter { private static PatchSheetResult patchSheet(byte[] original, Session session) { String content = new String(original, StandardCharsets.UTF_8); - int idx = 1; + int idx = 0; int patchedCount = 0; for (RegisteredCellImage image : session.images) { String replacement = "#VALUE!"; 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 4a4e730..5f986ef 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 @@ -39,6 +39,8 @@ import java.util.List; import java.util.Locale; import java.util.Map; import java.util.Set; +import java.util.concurrent.CompletionService; +import java.util.concurrent.ExecutorCompletionService; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; @@ -55,17 +57,36 @@ import java.util.concurrent.atomic.AtomicInteger; @Slf4j public class SimilarAsinImageEmbedder { - /** P2-8 默认值:原硬编码 5s,放宽到 8s 配合 1 次重试,整体更稳。可通过 properties 覆盖。 */ - static final int DEFAULT_DOWNLOAD_TIMEOUT_SECONDS = 8; - static final int DOWNLOAD_MAX_RETRY = 1; - /** P2-8 默认值:原硬编码 8,几千行 ×3 列图片场景下提升到 16 显著缩短 xlsx 组装阶段。 */ - static final int DEFAULT_DOWNLOAD_POOL_SIZE = 16; + /** + * P2-10 默认值:单图下载超时 5s。 + * 历史 P2-8 调到 8s + retry=1,但慢源(cbu01.alicdn)会一直挂 8s 才重试, + * 整体下载时间被尾延迟放大;改回 5s + retry=2 让慢源更快进入下次尝试。 + */ + static final int DEFAULT_DOWNLOAD_TIMEOUT_SECONDS = 5; + /** P2-10:retry 由 1 升到 2,配合 5s timeout 单图最坏耗时 ≈ 15s。 */ + static final int DOWNLOAD_MAX_RETRY = 2; + /** P2-10:原 8 → P2-8 16 → P2-10 32,1000+ 行 ×3 列场景下显著拉低 assemble 阶段尾延迟。 */ + static final int DEFAULT_DOWNLOAD_POOL_SIZE = 32; // 单元格固定尺寸,图片在其中等比缩放(不拉伸);resize 仍按长边 1280 px 控制堆体积。 public static final float IMAGE_ROW_HEIGHT_POINTS = 409f; public static final int IMAGE_COL_WIDTH_CHARS = 80; static final int TARGET_LONG_EDGE_PX = 1280; static final float JPEG_QUALITY = 0.75f; - static final int MAX_THUMB_SIZE_BYTES = 300 * 1024; + /** + * B 方案缩略图字节硬上限:300KB → 150KB(5000 行 × 3 列规模下,单 task 活跃图片字节 + * 上限直接砍半:300KB×N → 150KB×N,配合 BoundedImageCache 的 256MB LRU, + * 200 实例并发不再爆 2GB 堆)。 + * 出口超限时 resizeImage 会按 MAX_THUMB_SIZE → 长边 → 质量的顺序迭代降级, + * 仍然超限才抛 ResizeOversizeException。 + */ + static final int MAX_THUMB_SIZE_BYTES = 150 * 1024; + /** + * 迭代降级时的备选长边像素,按 1280 → 960 → 720 降; + * 不再缩到更小,因为 Excel 单元格列宽 80 字符(≈ 600 px)已是显示下限。 + */ + private static final int[] FALLBACK_LONG_EDGES = new int[]{1280, 960, 720}; + /** 迭代降级时的备选 JPEG 质量;末位 0.55 是肉眼可接受下限。 */ + private static final float[] FALLBACK_QUALITIES = new float[]{0.75f, 0.65f, 0.55f}; static final int MAX_DOWNLOAD_BYTES = 5 * 1024 * 1024; static final int MAX_DECODE_PIXELS = 6000 * 6000; @@ -114,6 +135,12 @@ public class SimilarAsinImageEmbedder { * - POI 写入是单线程,下载并行化不影响写入顺序; * - 失败 URL 不写 cache,由 embed() 沿用既有异常分类做文本兜底; * - 调用方需保证传入同一个 taskImageCache 给 embed()。 + * + *

P2-10 改造:原实现按"每个 future 等满 perTaskWaitMs"串行 join, + * 1000+ url 的尾部慢源会把整体时长堆到几百秒(实测 244s / 918s)。 + * 改为 {@link ExecutorCompletionService} + 全局 deadline: + * 已完成的 future 立即被收割,慢源在 deadline 后整体取消,避免被尾延迟拖死。 + * 全局 deadline = clamp(urlCount * 200ms, 15s, 120s),与 P2-10 retry/timeout 调整匹配。 */ public void prefetch(Collection urls, Map taskImageCache) { if (urls == null || urls.isEmpty() || taskImageCache == null) { @@ -133,9 +160,11 @@ public class SimilarAsinImageEmbedder { if (distinctUrls.isEmpty()) { return; } + CompletionService completion = new ExecutorCompletionService<>(downloadPool); List> futures = new ArrayList<>(distinctUrls.size()); for (String url : distinctUrls) { - futures.add(downloadPool.submit(() -> { + // 显式作为 Runnable 提交(带 null result),避免与 Callable 重载产生歧义。 + Runnable task = () -> { try { if (taskImageCache.containsKey(url)) { return; @@ -147,24 +176,93 @@ public class SimilarAsinImageEmbedder { // 预下载失败不抛出:embed() 时同 url 会再次尝试并走原有兜底链路。 log.debug("[similar-asin][image] prefetch-fail url={} err={}", url, ex.getMessage()); } - })); + }; + futures.add(completion.submit(task, null)); } - // 等待全部预下载结束(含失败)。下载池容量 = downloadPoolSize, - // 即使部分任务超时,单图也最多被 downloadTimeoutSeconds * 2 锁定。 - long perTaskWaitMs = downloadTimeoutSeconds * 2L * 1000L + 5000L; - for (Future f : futures) { + // 全局 deadline:每 url 给 200ms 的预算,clamp 到 [15s, 120s]。 + long globalDeadlineMs = Math.min(120_000L, Math.max(15_000L, distinctUrls.size() * 200L)); + long deadline = System.currentTimeMillis() + globalDeadlineMs; + int total = distinctUrls.size(); + int done = 0; + while (done < total) { + long left = deadline - System.currentTimeMillis(); + if (left <= 0) { + break; + } try { - f.get(perTaskWaitMs, TimeUnit.MILLISECONDS); - } catch (java.util.concurrent.TimeoutException te) { - f.cancel(true); + Future f = completion.poll(left, TimeUnit.MILLISECONDS); + if (f == null) { + break; + } + done++; } catch (InterruptedException ie) { Thread.currentThread().interrupt(); - f.cancel(true); - return; - } catch (Exception ignored) { - // 单任务失败不影响其它预下载——已在内部 catch 了。 + break; } } + // 超时未完成的 future 主动 cancel,避免 idle 持有 OkHttp 连接。 + if (done < total) { + for (Future f : futures) { + if (!f.isDone()) { + f.cancel(true); + } + } + log.info("[similar-asin][image] prefetch deadline reached total={} done={} cancelled={}", + total, done, total - done); + } + } + + /** + * P2-11:给 {@code SimilarAsinImagePrefetchService} 用的对外入口。 + * 仅做下载 + resize,不写 taskImageCache(DB cache 由 service 层处理)。 + * 失败统一返回 null,调用方决定是否落表/重试。 + */ + public ResizedImage fetchAndResizeForCache(String url) { + if (url == null) { + return null; + } + String trimmed = url.trim(); + if (trimmed.isEmpty()) { + return null; + } + try { + byte[] raw = downloadWithRetry(trimmed); + return resizeImage(trimmed, raw); + } catch (Exception ex) { + log.debug("[similar-asin][image] prefetch-cache-fail url={} err={}", trimmed, ex.getMessage()); + return null; + } + } + + /** + * P2-11:从 DB cache 中读到的字节恢复 ResizedImage。仅做尺寸读取,不做二次压缩; + * 字节保持与原写入一致(避免 JPEG 反复编码导致的画质退化与字节膨胀)。 + */ + public ResizedImage decodeCachedThumb(byte[] cachedBytes) { + if (cachedBytes == null || cachedBytes.length == 0) { + return null; + } + try (ImageInputStream iis = ImageIO.createImageInputStream(new ByteArrayInputStream(cachedBytes))) { + if (iis == null) { + return null; + } + Iterator readers = ImageIO.getImageReaders(iis); + if (!readers.hasNext()) { + return null; + } + ImageReader reader = readers.next(); + try { + reader.setInput(iis, true, true); + int width = reader.getWidth(0); + int height = reader.getHeight(0); + return new ResizedImage(cachedBytes, width, height); + } finally { + reader.dispose(); + } + } catch (Exception ex) { + log.debug("[similar-asin][image] decode-cached-thumb-fail bytes={} err={}", cachedBytes.length, ex.getMessage()); + return null; + } } /** @@ -345,7 +443,8 @@ public class SimilarAsinImageEmbedder { /** * 等比缩放到长边 TARGET_LONG_EDGE_PX,JPEG q=0.75 输出;返回字节 + 实际像素,供调用方按图自适应单元格尺寸。 - * 出口校验字节数 ≤ MAX_THUMB_SIZE_BYTES,超限抛 ResizeOversizeException 触发文本兜底。 + * 出口校验字节数 ≤ MAX_THUMB_SIZE_BYTES,超限按 MAX_THUMB_SIZE → 长边 → 质量的顺序迭代降级, + * 仍然超限才抛 ResizeOversizeException 触发文本兜底。 */ ResizedImage resizeImage(String sourceUrl, byte[] raw) throws IOException { guardImageDimensions(sourceUrl, raw); @@ -355,7 +454,36 @@ public class SimilarAsinImageEmbedder { } int srcW = src.getWidth(); int srcH = src.getHeight(); - double ratio = (double) Math.max(srcW, srcH) / TARGET_LONG_EDGE_PX; + // B 方案降级顺序:固定 MAX_THUMB_SIZE 上限 → 优先调整长边像素 → 再调质量。 + // 同一 src BufferedImage 解码一次,下面 9 种组合复用,避免重复 ImageIO.read。 + ResizedImage candidate = null; + ResizedImage smallest = null; + for (int longEdge : FALLBACK_LONG_EDGES) { + for (float quality : FALLBACK_QUALITIES) { + ResizedImage tried = encodeAt(src, srcW, srcH, longEdge, quality); + if (smallest == null || tried.bytes().length < smallest.bytes().length) { + smallest = tried; + } + if (tried.bytes().length <= MAX_THUMB_SIZE_BYTES) { + candidate = tried; + break; + } + } + if (candidate != null) { + break; + } + } + if (candidate != null) { + return candidate; + } + // 9 组合都没压到上限:抛 ResizeOversizeException 走文本兜底, + // 同时报告 smallest 字节让运维直观知道当前压缩极限。 + int reportedSize = smallest != null ? smallest.bytes().length : -1; + throw new ResizeOversizeException(sourceUrl, reportedSize); + } + + private ResizedImage encodeAt(BufferedImage src, int srcW, int srcH, int longEdgePx, float quality) throws IOException { + double ratio = (double) Math.max(srcW, srcH) / longEdgePx; int dstW = ratio > 1 ? Math.max(1, (int) Math.round(srcW / ratio)) : srcW; int dstH = ratio > 1 ? Math.max(1, (int) Math.round(srcH / ratio)) : srcH; BufferedImage dst = new BufferedImage(dstW, dstH, BufferedImage.TYPE_INT_RGB); @@ -376,7 +504,7 @@ public class SimilarAsinImageEmbedder { try { ImageWriteParam param = writer.getDefaultWriteParam(); param.setCompressionMode(ImageWriteParam.MODE_EXPLICIT); - param.setCompressionQuality(JPEG_QUALITY); + param.setCompressionQuality(quality); ImageOutputStream ios = ImageIO.createImageOutputStream(baos); try { writer.setOutput(ios); @@ -387,9 +515,6 @@ public class SimilarAsinImageEmbedder { } finally { writer.dispose(); } - if (baos.size() > MAX_THUMB_SIZE_BYTES) { - throw new ResizeOversizeException(sourceUrl, baos.size()); - } return new ResizedImage(baos.toByteArray(), dstW, dstH); } @@ -523,10 +648,13 @@ public class SimilarAsinImageEmbedder { } /** - * resize 出口字节硬上限保护:B 方案下单图最坏 ~300KB;POI picture pool 在 SXSSFWorkbook.dispose() 前不会 spill 到临时文件, - * 因此 N 行 × 3 列 × 300KB 全量驻留堆。配合 -Xmx2048M:≈1500 行 × 3 列 ≈ 1.32GB picture pool,已逼近安全水位; - * embed() 成功后会立刻 taskImageCache.remove() 释放 A 副本,B 副本(picture pool)仍随 workbook 生命周期驻留。 - * 超过 ~1500 行需要降低 MAX_THUMB_SIZE_BYTES 或调小 TARGET_LONG_EDGE_PX,必要时再考虑拆任务。 + * resize 出口字节硬上限保护:B 方案降到 150KB(原 300KB)。POI picture pool 在 SXSSFWorkbook.dispose() + * 前不会 spill 到临时文件,N 行 × 3 列 × 150KB 全量驻留堆。配合 -Xmx2048M:5000 行 × 3 列 ≈ 2.25GB + * picture pool 仍超出安全水位,因此 5000 行任务必须依赖: + * (a) BoundedImageCache 的 LRU evict(assemble 阶段共享 256MB 上限); + * (b) embed() 成功后立刻 taskImageCache.remove() 释放 A 副本,但 B 副本随 workbook 生命周期驻留; + * (c) 单 task 触发的 N 个 source workbook 串行/最多 4 并发,避免 4×5000×3×150KB 同时驻留。 + * 仍出现堆压力时再调小 TARGET_LONG_EDGE_PX 或拆任务。 */ public static class ResizeOversizeException extends ResizeException { private final String url; diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/mapper/TaskImageCacheMapper.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/mapper/TaskImageCacheMapper.java new file mode 100644 index 0000000..7f13150 --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/mapper/TaskImageCacheMapper.java @@ -0,0 +1,29 @@ +package com.nanri.aiimage.modules.task.mapper; + +import com.baomidou.mybatisplus.core.mapper.BaseMapper; +import com.nanri.aiimage.modules.task.model.entity.TaskImageCacheEntity; +import org.apache.ibatis.annotations.Mapper; +import org.apache.ibatis.annotations.Param; +import org.apache.ibatis.annotations.Select; +import org.apache.ibatis.annotations.Update; + +/** + * P2-12:图片缩略图缓存 mapper。提供 LRU 命中刷新、按 url_hash 直读字节两个轻量入口, + * 其余 CRUD 走 {@link BaseMapper} 默认实现。 + */ +@Mapper +public interface TaskImageCacheMapper extends BaseMapper { + + /** + * 命中时同步 bumping last_used_at,便于后续按 LRU 清理 7 天前未用记录。 + */ + @Update("UPDATE biz_task_image_cache SET last_used_at = NOW(3) WHERE url_hash = #{urlHash}") + int touchLastUsed(@Param("urlHash") String urlHash); + + /** + * 按 url_hash 直读缩略图字节,命中时返回 BLOB;未命中返回 null。 + * 选择只 select image_bytes 一列,避免把整行 entity(含 url 字符串)拉回 JVM。 + */ + @Select("SELECT image_bytes FROM biz_task_image_cache WHERE url_hash = #{urlHash} LIMIT 1") + byte[] selectBytesByUrlHash(@Param("urlHash") String urlHash); +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/entity/TaskImageCacheEntity.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/entity/TaskImageCacheEntity.java new file mode 100644 index 0000000..bffff0f --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/entity/TaskImageCacheEntity.java @@ -0,0 +1,36 @@ +package com.nanri.aiimage.modules.task.model.entity; + +import com.baomidou.mybatisplus.annotation.IdType; +import com.baomidou.mybatisplus.annotation.TableId; +import com.baomidou.mybatisplus.annotation.TableName; +import lombok.Data; + +import java.time.LocalDateTime; + +/** + * P2-12:图片缩略图缓存。配合 {@code SimilarAsinImagePrefetchService} 跨任务复用 Coze + * 回包中的 main_url / puzzle_img 图片,避免 assemble 阶段每次都重新下载。 + * + *

对应表 {@code biz_task_image_cache}(migration V53)。 + *

    + *
  • {@link #urlHash}:sha256 lowercase hex,作为唯一键避免 1024 长 URL 命中索引长度限制;
  • + *
  • {@link #imageBytes}:已 resize 后的 JPEG 缩略图字节,最大约 300KB(保持与 + * {@code SimilarAsinImageEmbedder.MAX_THUMB_SIZE_BYTES} 对齐);
  • + *
  • {@link #lastUsedAt}:用于后续 LRU 清理(>7 天未用)。
  • + *
+ */ +@Data +@TableName("biz_task_image_cache") +public class TaskImageCacheEntity { + + @TableId(type = IdType.AUTO) + private Long id; + private String urlHash; + private String url; + private byte[] imageBytes; + private Integer byteSize; + private Integer width; + private Integer height; + private LocalDateTime createdAt; + private LocalDateTime lastUsedAt; +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TransientPayloadStorageService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TransientPayloadStorageService.java index b665cdf..2ddcb60 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TransientPayloadStorageService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TransientPayloadStorageService.java @@ -1,11 +1,16 @@ package com.nanri.aiimage.modules.task.service; +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.fasterxml.jackson.databind.ObjectMapper; import com.nanri.aiimage.config.InstanceMetadata; import com.nanri.aiimage.config.StorageProperties; import com.nanri.aiimage.config.TransientStorageProperties; import com.nanri.aiimage.modules.file.service.object.RustfsObjectStorageService; import com.nanri.aiimage.modules.file.service.oss.OssStorageService; +import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; +import com.nanri.aiimage.modules.task.mapper.TaskScopeStateMapper; +import com.nanri.aiimage.modules.task.model.entity.TaskChunkEntity; +import com.nanri.aiimage.modules.task.model.entity.TaskScopeStateEntity; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -16,9 +21,13 @@ import java.nio.file.Path; import java.nio.file.StandardOpenOption; import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; +import java.util.ArrayList; import java.util.Base64; +import java.util.LinkedHashSet; +import java.util.List; import java.util.Locale; import java.util.Objects; +import java.util.Set; import java.util.UUID; import java.util.zip.GZIPInputStream; import java.util.zip.GZIPOutputStream; @@ -34,12 +43,43 @@ public class TransientPayloadStorageService { private static final String OSS_POINTER_PREFIX = "oss:"; private static final String LOCAL_PAYLOAD_DIR = "transient-payload"; + /** + * P0-3:记录"最近一次 store 是否因 RustFS 失败回落到 local"。 + * 调用方(如 SimilarAsinTaskService.submitResult)拿到 storedPayload 后可以查询 + * 此 ThreadLocal 决定是否把 task 绑定到 owner instance、是否更新 state_json + * 的 localFallback 标记,从而保证后续 assemble 走对实例。 + */ + private static final ThreadLocal LAST_STORE_LOCAL_FALLBACK = ThreadLocal.withInitial(() -> Boolean.FALSE); + private final TransientStorageProperties properties; private final StorageProperties storageProperties; private final RustfsObjectStorageService rustfsObjectStorageService; private final OssStorageService ossStorageService; private final ObjectMapper objectMapper; private final InstanceMetadata instanceMetadata; + /** + * P2-9:全局引用计数依赖的两个 mapper。 + * 删除 transient payload 之前要反查 biz_task_chunk / biz_task_scope_state, + * 任何一张表还有其它行引用同一个 pointer/value 就跳过删除,避免 deterministic key + * 残留场景下"赢家写、输家删"造成的下游悬空。 + */ + private final TaskChunkMapper taskChunkMapper; + private final TaskScopeStateMapper taskScopeStateMapper; + + /** + * P0-3:返回上一次 store 调用是否走了 RustFS → local 兜底。 + * 注意:跨线程不传递;调用方在同一线程内 store 完后立即读取。 + */ + public boolean wasLastStoreLocalFallback() { + return Boolean.TRUE.equals(LAST_STORE_LOCAL_FALLBACK.get()); + } + + /** + * P0-3:返回当前实例 ID,供调用方写入 ownerInstance 锁定。 + */ + public String currentInstanceId() { + return instanceMetadata.getInstanceId(); + } public boolean isWriteEnabled() { return properties.isEnabled() @@ -124,6 +164,13 @@ public class TransientPayloadStorageService { if (pointer == null) { return; } + // P2-9:全局引用计数兜底——deterministic key 残留场景下, + // 同一个 pointer 可能被多个 biz_task_chunk / biz_task_scope_state 行共享。 + // 删之前先反查一次,确认没有别的行还在引用,再真正落地物理删除。 + if (isStillReferenced(pointer, value)) { + log.info("[transient-payload] skip delete, still referenced pointer={}", pointer); + return; + } if (pointer.startsWith(LOCAL_POINTER_PREFIX)) { String localKey = pointer.substring(LOCAL_POINTER_PREFIX.length()); String ownerInstance = extractLocalInstanceId(localKey); @@ -145,6 +192,67 @@ public class TransientPayloadStorageService { } } + /** + * P2-9:判断给定 pointer 是否仍被 biz_task_chunk / biz_task_scope_state 中的其它行引用。 + * + *

背景:历史 deterministic key({@code storeChunkPayload} 等)在并发回传时多个调用方会 + * 拿到同一个 pointer,DB 里也会出现多条不同的 task_chunk 行共享同一个 payload 字段值。 + * 老路径 caller 把自己那行覆盖/删除后立即调用 {@link #deletePayloadIfPresent}, + * 此时若不做反查就会把"赢家"对象一并物理删掉,下游读 chunk 直接变成 + * {@code read chunk payload failed}。 + * + *

判断口径: + *

    + *
  • {@code biz_task_chunk.payload_json} 命中 > 1 行(> 1 表示除了 caller 视角下 + * 自己即将释放的那一行之外,至少还有别的 chunk 行也指向同一对象)→ 视为仍被引用。
  • + *
  • {@code biz_task_scope_state.parsed_payload_json}/{@code state_json} 命中 > 0 行 → + * 视为仍被引用(这两个字段不是 caller 自身行的常见持有者,命中即非自我引用)。
  • + *
+ * + *

查询时同时把 caller 传入的 raw value 与 {@code pointer}、{@code "pointer"}(JSON 编码后 + * 的字符串形式)都纳入候选,覆盖以下两类常见 DB 字段写入形式: + *

    + *
  • {@code objectMapper.writeValueAsString(pointer)} 写入的 JSON 字符串包裹形式;
  • + *
  • 少数老路径直接写入裸 pointer 的形式(如 {@link #deleteReplacedPayloadIfNeeded})。
  • + *
+ * + *

查询出现异常时保守返回 {@code true}(不删),由后续清理任务兜底。 + */ + private boolean isStillReferenced(String pointer, String originalValue) { + try { + Set candidates = new LinkedHashSet<>(); + if (originalValue != null && !originalValue.isBlank()) { + candidates.add(originalValue); + } + if (pointer != null && !pointer.isBlank()) { + candidates.add(pointer); + try { + candidates.add(objectMapper.writeValueAsString(pointer)); + } catch (Exception ignored) { + // 编码失败时仅依赖其它候选 + } + } + if (candidates.isEmpty()) { + return false; + } + List values = new ArrayList<>(candidates); + Long chunkCount = taskChunkMapper.selectCount(new LambdaQueryWrapper() + .in(TaskChunkEntity::getPayloadJson, values)); + if (chunkCount != null && chunkCount > 1L) { + return true; + } + Long scopeStateCount = taskScopeStateMapper.selectCount(new LambdaQueryWrapper() + .and(w -> w.in(TaskScopeStateEntity::getParsedPayloadJson, values) + .or() + .in(TaskScopeStateEntity::getStateJson, values))); + return scopeStateCount != null && scopeStateCount > 0L; + } catch (Exception ex) { + log.warn("[transient-payload] reference check failed pointer={} err={}", pointer, ex.getMessage()); + // 查询出错时保守:不删,等周期性清理兜底 + return true; + } + } + public void deleteReplacedPayloadIfNeeded(String oldValue, String newValue) { String oldPointer = extractPointer(oldValue); String newPointer = extractPointer(newValue); @@ -206,18 +314,22 @@ public class TransientPayloadStorageService { String objectKey = buildObjectKey(category, moduleType, taskId, scopeHash, entryKey); String storedContent = encodeStoredPayload(content); String pointer = null; + boolean rustfsFallbackToLocal = false; if (rustfsObjectStorageService.isConfigured()) { try { pointer = RUSTFS_POINTER_PREFIX + rustfsObjectStorageService.uploadText(objectKey, storedContent, verifyAfterUpload); } catch (Exception ex) { + rustfsFallbackToLocal = true; // 升级为 ERROR:rustfs 失败后只能落到本地,多实例下其他节点读不到,必须能告警。 - log.error("[transient-payload] rustfs upload failed, fallback to local store instanceId={} objectKey={} err={}", - instanceMetadata.getInstanceId(), objectKey, ex.getMessage()); + log.error("[transient-payload] rustfs upload failed, fallback to local store instanceId={} category={} taskId={} objectKey={} err={}", + instanceMetadata.getInstanceId(), category, taskId, objectKey, ex.getMessage()); } } if (pointer == null) { pointer = storeLocal(objectKey, storedContent); } + // P0-3:记录本次 store 是否走 local 兜底,供调用方在拿到 pointer 后立即查询。 + LAST_STORE_LOCAL_FALLBACK.set(rustfsFallbackToLocal && pointer != null && pointer.startsWith(LOCAL_POINTER_PREFIX)); if (pointer == null) { return content; } diff --git a/backend-java/src/main/resources/application-local.example.yml b/backend-java/src/main/resources/application-local.example.yml index 50e2d9e..a0f5eb6 100644 --- a/backend-java/src/main/resources/application-local.example.yml +++ b/backend-java/src/main/resources/application-local.example.yml @@ -43,7 +43,7 @@ AIIMAGE_SIMILAR_ASIN_COZE_BASE_URL=https://api.coze.cn AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_PATH=/v1/workflow/run AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_ID=7635328462404583478 AIIMAGE_SIMILAR_ASIN_COZE_TOKEN= -AIIMAGE_SIMILAR_ASIN_COZE_BATCH_SIZE=50 +AIIMAGE_SIMILAR_ASIN_COZE_BATCH_SIZE=3 AIIMAGE_SIMILAR_ASIN_COZE_READ_TIMEOUT_MILLIS=60000 AIIMAGE_SIMILAR_ASIN_STALE_TIMEOUT_MINUTES=30 diff --git a/backend-java/src/main/resources/application.yml b/backend-java/src/main/resources/application.yml index 0d02cce..4e619bb 100644 --- a/backend-java/src/main/resources/application.yml +++ b/backend-java/src/main/resources/application.yml @@ -166,17 +166,18 @@ aiimage: coze-workflow-path: ${AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_PATH:/v1/workflow/run} coze-workflow-id: ${AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_ID:7639708860686024756} coze-token: ${AIIMAGE_SIMILAR_ASIN_COZE_TOKEN:} - coze-batch-size: ${AIIMAGE_SIMILAR_ASIN_COZE_BATCH_SIZE:10} + coze-batch-size: ${AIIMAGE_SIMILAR_ASIN_COZE_BATCH_SIZE:3} coze-credential-stripe-size: ${AIIMAGE_SIMILAR_ASIN_COZE_CREDENTIAL_STRIPE_SIZE:0} coze-connect-timeout-millis: ${AIIMAGE_SIMILAR_ASIN_COZE_CONNECT_TIMEOUT_MILLIS:10000} coze-read-timeout-millis: ${AIIMAGE_SIMILAR_ASIN_COZE_READ_TIMEOUT_MILLIS:60000} coze-poll-interval-millis: ${AIIMAGE_SIMILAR_ASIN_COZE_POLL_INTERVAL_MILLIS:30000} coze-poll-timeout-millis: ${AIIMAGE_SIMILAR_ASIN_COZE_POLL_TIMEOUT_MILLIS:1800000} coze-submit-min-interval-millis: ${AIIMAGE_SIMILAR_ASIN_COZE_SUBMIT_MIN_INTERVAL_MILLIS:5000} - coze-flush-pending-minutes: ${AIIMAGE_SIMILAR_ASIN_COZE_FLUSH_PENDING_MINUTES:15} + coze-flush-pending-minutes: ${AIIMAGE_SIMILAR_ASIN_COZE_FLUSH_PENDING_MINUTES:1} coze-submit-max-retry-count: ${AIIMAGE_SIMILAR_ASIN_COZE_SUBMIT_MAX_RETRY_COUNT:5} - image-download-pool-size: ${AIIMAGE_SIMILAR_ASIN_IMAGE_DOWNLOAD_POOL_SIZE:16} - image-download-timeout-seconds: ${AIIMAGE_SIMILAR_ASIN_IMAGE_DOWNLOAD_TIMEOUT_SECONDS:8} + image-download-pool-size: ${AIIMAGE_SIMILAR_ASIN_IMAGE_DOWNLOAD_POOL_SIZE:32} + image-download-timeout-seconds: ${AIIMAGE_SIMILAR_ASIN_IMAGE_DOWNLOAD_TIMEOUT_SECONDS:5} + image-cache-max-bytes: ${AIIMAGE_SIMILAR_ASIN_IMAGE_CACHE_MAX_BYTES:268435456} stale-timeout-minutes: ${AIIMAGE_SIMILAR_ASIN_STALE_TIMEOUT_MINUTES:30} stale-finalize-cron: ${AIIMAGE_SIMILAR_ASIN_STALE_FINALIZE_CRON:0 */2 * * * *} coze-include-legacy-api-key: ${AIIMAGE_SIMILAR_ASIN_COZE_INCLUDE_LEGACY_API_KEY:true} diff --git a/backend-java/src/main/resources/db/V53__task_image_cache.sql b/backend-java/src/main/resources/db/V53__task_image_cache.sql new file mode 100644 index 0000000..d3e8e4a --- /dev/null +++ b/backend-java/src/main/resources/db/V53__task_image_cache.sql @@ -0,0 +1,18 @@ +-- P2-12:相似ASIN/通用 图片缩略图缓存表。 +-- 用于跨任务复用 Coze 回包中的 main_url / puzzle_img1 / puzzle_img2 等图片, +-- 与 SimilarAsinImagePrefetchService 协作,避免 assemble 阶段每次都重新下载远程图片。 +-- url_hash 走 sha256(lowercase hex),避免 1024 长 URL 作为唯一键命中索引长度限制。 +CREATE TABLE IF NOT EXISTS biz_task_image_cache ( + id BIGINT NOT NULL AUTO_INCREMENT, + url_hash CHAR(64) NOT NULL COMMENT 'sha256 lowercase hex of url', + url VARCHAR(1024) NOT NULL, + image_bytes MEDIUMBLOB NOT NULL COMMENT '已 resize 缩略图字节,最大约 300KB', + byte_size INT NOT NULL DEFAULT 0, + width INT NOT NULL DEFAULT 0, + height INT NOT NULL DEFAULT 0, + created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3), + last_used_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3), + PRIMARY KEY (id), + UNIQUE KEY uk_task_image_cache_url_hash (url_hash), + KEY idx_task_image_cache_last_used_at (last_used_at) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='相似ASIN/通用 图片缩略图缓存,跨任务复用'; diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/util/ExcelCellImageWriterTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/util/ExcelCellImageWriterTest.java new file mode 100644 index 0000000..9494415 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/util/ExcelCellImageWriterTest.java @@ -0,0 +1,105 @@ +package com.nanri.aiimage.modules.similarasin.util; + +import org.junit.jupiter.api.Test; + +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.regex.Matcher; +import java.util.regex.Pattern; +import java.util.zip.ZipEntry; +import java.util.zip.ZipFile; +import java.util.zip.ZipOutputStream; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class ExcelCellImageWriterTest { + + @Test + void patchedCellImageValueMetadataUsesZeroBasedVmIndexes() throws Exception { + Path dir = Files.createTempDirectory("excel-cell-image-test-"); + Path xlsx = dir.resolve("result.xlsx"); + writeMinimalWorkbook(xlsx); + + ExcelCellImageWriter.Session session = ExcelCellImageWriter.createSession(); + session.registerImage(1, 9, new byte[]{1, 2, 3}); + session.registerImage(1, 10, new byte[]{4, 5, 6}); + + ExcelCellImageWriter.patchXlsxFile(xlsx.toFile(), session); + + try (ZipFile zip = new ZipFile(xlsx.toFile())) { + String sheet = read(zip, "xl/worksheets/sheet1.xml"); + assertTrue(sheet.contains("#VALUE!")); + assertTrue(sheet.contains("#VALUE!")); + + String metadata = read(zip, "xl/metadata.xml"); + assertTrue(metadata.contains("")); + assertTrue(metadata.contains("")); + assertEquals(2, maxVm(sheet) + 1); + } + } + + private static int maxVm(String sheet) { + Matcher matcher = Pattern.compile("vm=\"(\\d+)\"").matcher(sheet); + int max = -1; + while (matcher.find()) { + max = Math.max(max, Integer.parseInt(matcher.group(1))); + } + return max; + } + + private static String read(ZipFile zip, String name) throws Exception { + return new String(zip.getInputStream(zip.getEntry(name)).readAllBytes(), StandardCharsets.UTF_8); + } + + private static void writeMinimalWorkbook(Path xlsx) throws Exception { + try (ZipOutputStream out = new ZipOutputStream(Files.newOutputStream(xlsx))) { + write(out, "[Content_Types].xml", """ + + + + + + + + """); + write(out, "_rels/.rels", """ + + + + + """); + write(out, "xl/workbook.xml", """ + + + + + """); + write(out, "xl/_rels/workbook.xml.rels", """ + + + + + """); + write(out, "xl/worksheets/sheet1.xml", """ + + + + + main + puzzle + + + + """); + } + } + + private static void write(ZipOutputStream out, String name, String content) throws Exception { + out.putNextEntry(new ZipEntry(name)); + out.write(content.getBytes(StandardCharsets.UTF_8)); + out.closeEntry(); + } +}