提交货源流程优化
This commit is contained in:
4
app/.env
4
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
|
||||
|
||||
@@ -16,7 +16,14 @@ public class SimilarAsinProperties {
|
||||
private String cozeToken = "";
|
||||
private List<CozeCredential> 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;
|
||||
|
||||
@@ -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 图片异步预热服务。
|
||||
*
|
||||
* <p>背景:assemble 阶段({@code assembleResultWorkbook})需要把 Coze 回包中的 main_url /
|
||||
* puzzle_img1 / puzzle_img2 下载并 resize 后嵌入 xlsx。当任务行数到 1000+ 时,串行 +
|
||||
* 短池下载会把整个 assemble 拖到 244s / 918s。改造点:
|
||||
*
|
||||
* <ul>
|
||||
* <li>每次 {@code mergeCozeRowsIntoChunk} 拿到新 cozeRows 时,调用 {@link #enqueue}
|
||||
* 立即丢入预热队列;同 task 串行排队({@link #inflight}),避免多个 batch 同时打爆图片源站;</li>
|
||||
* <li>预热成功的缩略图字节落表 {@code biz_task_image_cache}(由 P2-12 提供),跨任务复用;</li>
|
||||
* <li>所有路径 best-effort:预热失败、DB 写入失败都吞掉,assemble 阶段会回退到原下载链路兜底。</li>
|
||||
* </ul>
|
||||
*/
|
||||
@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<Long, Future<?>> 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<String> urls) {
|
||||
if (taskId == null || urls == null || urls.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
Set<String> 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<String> 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<String> 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<TaskImageCacheEntity>()
|
||||
.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;
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -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<String, Long> 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<Long, AtomicReferenceArray<Boolean>> recentPoisonByTask = new ConcurrentHashMap<>();
|
||||
private final ConcurrentHashMap<Long, AtomicInteger> 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<SimilarAsinFilterConditionVo> listFilterConditions(Long userId) {
|
||||
validateUserId(userId);
|
||||
List<SimilarAsinFilterConditionEntity> 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<CozeCandidate> candidates = collectPendingCozeCandidates(chunks, allRowsByBaseId);
|
||||
if (candidates.isEmpty()) {
|
||||
return countPendingCozeStates(task.getId()) > 0;
|
||||
}
|
||||
List<CozeCandidate> 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<SimilarAsinResultRowDto> isMissingRequired = row ->
|
||||
row == null
|
||||
|| !row.hasImageUrl()
|
||||
|| normalize(row.getAsin()).isBlank()
|
||||
|| normalize(row.getTitle()).isBlank();
|
||||
List<CozeCandidate> 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<CozeCandidate> 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<Boolean> 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<Boolean> 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<TaskScopeStateEntity>()
|
||||
.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<SimilarAsinResultRowDto> 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<String> 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<String, Map<String, SimilarAsinResultRowDto>> rowsByChunk = new LinkedHashMap<>();
|
||||
Map<String, TaskChunkEntity> 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<TaskScopeStateEntity>()
|
||||
.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> sourceRows = splitRowsBySourceFile(parsed, parsed.getAllItems(), result.getSourceFilename());
|
||||
List<SourceResultWorkbook> 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<CompletableFuture<SourceResultWorkbook>> 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<SourceResultWorkbook> future : futures) {
|
||||
try {
|
||||
// 单源文件 30 分钟硬上限:超时直接抛错,避免被慢源永久阻塞。
|
||||
workbooks.add(future.get(30, TimeUnit.MINUTES));
|
||||
} catch (InterruptedException ie) {
|
||||
Thread.currentThread().interrupt();
|
||||
for (CompletableFuture<SourceResultWorkbook> remaining : futures) {
|
||||
remaining.cancel(true);
|
||||
}
|
||||
throw new BusinessException("生成相似ASIN检测结果中断", ie);
|
||||
} catch (TimeoutException te) {
|
||||
for (CompletableFuture<SourceResultWorkbook> remaining : futures) {
|
||||
remaining.cancel(true);
|
||||
}
|
||||
throw new BusinessException("生成相似ASIN检测结果超时(30 分钟)", te);
|
||||
} catch (ExecutionException ee) {
|
||||
Throwable cause = ee.getCause();
|
||||
for (CompletableFuture<SourceResultWorkbook> 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<SimilarAsinParsedRowVo> rowsToWrite,
|
||||
Map<String, SimilarAsinResultRowDto> resultMap) {
|
||||
Map<String, SimilarAsinResultRowDto> resultMap,
|
||||
Map<String, SimilarAsinImageEmbedder.ResizedImage> 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<String, SimilarAsinImageEmbedder.ResizedImage> 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<String> 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<String> 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<TaskScopeStateEntity> scopes = taskScopeStateMapper.selectList(
|
||||
new LambdaQueryWrapper<TaskScopeStateEntity>()
|
||||
.eq(TaskScopeStateEntity::getTaskId, taskId)
|
||||
.eq(TaskScopeStateEntity::getModuleType, MODULE_TYPE));
|
||||
if (scopes == null || scopes.isEmpty()) {
|
||||
return "";
|
||||
}
|
||||
LinkedHashSet<String> 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<TaskScopeStateEntity>()
|
||||
.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<TaskScopeStateEntity>()
|
||||
.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)}
|
||||
|
||||
@@ -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 淘汰。
|
||||
*
|
||||
* <p>背景:原实现是 {@code ConcurrentHashMap<String, ResizedImage>},在 5000 行 × 200 实例
|
||||
* 规模下、单图缩略图最大 300KB 时,单 task 可能堆积 5000 × 3 × 300KB ≈ 4.5GB 图片字节,
|
||||
* 远超 JVM 2GB 堆。需要硬上限避免堆爆。
|
||||
*
|
||||
* <p>设计:
|
||||
* <ul>
|
||||
* <li>底层 {@link LinkedHashMap} access-order 维护 LRU;</li>
|
||||
* <li>所有写入路径 {@code put / putIfAbsent} 后 evictIfOverflow,按字节累计淘汰最久未访问条目;</li>
|
||||
* <li>{@code containsKey / get} 也会更新 LRU 顺序;</li>
|
||||
* <li>整个类对外仍是 {@code Map<String, ResizedImage>},调用方无感知;</li>
|
||||
* <li>所有公共方法 synchronized:embed 阶段单线程主导,prefetch 阶段并发只通过 putIfAbsent
|
||||
* 少量竞争,加锁开销可忽略。</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>不实现 entrySet/keySet/values/equals/hashCode 等少用方法(throw UnsupportedOperationException)。
|
||||
* 与现有 SimilarAsinImageEmbedder / SimilarAsinTaskService 中使用的 6 个方法(containsKey、get、
|
||||
* put、putIfAbsent、remove、size)严格对齐。
|
||||
*/
|
||||
@Slf4j
|
||||
public class BoundedImageCache implements Map<String, SimilarAsinImageEmbedder.ResizedImage> {
|
||||
|
||||
/**
|
||||
* 默认堆字节预算:256MB。
|
||||
* <p>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<String, SimilarAsinImageEmbedder.ResizedImage> 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<? extends String, ? extends SimilarAsinImageEmbedder.ResizedImage> m) {
|
||||
if (m == null || m.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
for (Map.Entry<? extends String, ? extends SimilarAsinImageEmbedder.ResizedImage> 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<Map.Entry<String, SimilarAsinImageEmbedder.ResizedImage>> it = backing.entrySet().iterator();
|
||||
while (it.hasNext() && currentBytes.get() > maxBytes) {
|
||||
Map.Entry<String, SimilarAsinImageEmbedder.ResizedImage> 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<String> keySet() {
|
||||
throw new UnsupportedOperationException("BoundedImageCache.keySet not supported");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Collection<SimilarAsinImageEmbedder.ResizedImage> values() {
|
||||
throw new UnsupportedOperationException("BoundedImageCache.values not supported");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Entry<String, SimilarAsinImageEmbedder.ResizedImage>> entrySet() {
|
||||
throw new UnsupportedOperationException("BoundedImageCache.entrySet not supported");
|
||||
}
|
||||
}
|
||||
@@ -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 = "<c r=\"" + image.cellRef + "\" t=\"e\" vm=\"" + idx + "\"><v>#VALUE!</v></c>";
|
||||
|
||||
@@ -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()。
|
||||
*
|
||||
* <p>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<String> urls, Map<String, ResizedImage> taskImageCache) {
|
||||
if (urls == null || urls.isEmpty() || taskImageCache == null) {
|
||||
@@ -133,9 +160,11 @@ public class SimilarAsinImageEmbedder {
|
||||
if (distinctUrls.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
CompletionService<Object> completion = new ExecutorCompletionService<>(downloadPool);
|
||||
List<Future<?>> futures = new ArrayList<>(distinctUrls.size());
|
||||
for (String url : distinctUrls) {
|
||||
futures.add(downloadPool.submit(() -> {
|
||||
// 显式作为 Runnable 提交(带 null result),避免与 Callable<Object> 重载产生歧义。
|
||||
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<Object> 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<ImageReader> 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;
|
||||
|
||||
@@ -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<TaskImageCacheEntity> {
|
||||
|
||||
/**
|
||||
* 命中时同步 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);
|
||||
}
|
||||
@@ -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 阶段每次都重新下载。
|
||||
*
|
||||
* <p>对应表 {@code biz_task_image_cache}(migration V53)。
|
||||
* <ul>
|
||||
* <li>{@link #urlHash}:sha256 lowercase hex,作为唯一键避免 1024 长 URL 命中索引长度限制;</li>
|
||||
* <li>{@link #imageBytes}:已 resize 后的 JPEG 缩略图字节,最大约 300KB(保持与
|
||||
* {@code SimilarAsinImageEmbedder.MAX_THUMB_SIZE_BYTES} 对齐);</li>
|
||||
* <li>{@link #lastUsedAt}:用于后续 LRU 清理(>7 天未用)。</li>
|
||||
* </ul>
|
||||
*/
|
||||
@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;
|
||||
}
|
||||
@@ -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<Boolean> 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 中的其它行引用。
|
||||
*
|
||||
* <p>背景:历史 deterministic key({@code storeChunkPayload} 等)在并发回传时多个调用方会
|
||||
* 拿到同一个 pointer,DB 里也会出现多条不同的 task_chunk 行共享同一个 payload 字段值。
|
||||
* 老路径 caller 把自己那行覆盖/删除后立即调用 {@link #deletePayloadIfPresent},
|
||||
* 此时若不做反查就会把"赢家"对象一并物理删掉,下游读 chunk 直接变成
|
||||
* {@code read chunk payload failed}。
|
||||
*
|
||||
* <p>判断口径:
|
||||
* <ul>
|
||||
* <li>{@code biz_task_chunk.payload_json} 命中 > 1 行(> 1 表示除了 caller 视角下
|
||||
* 自己即将释放的那一行之外,至少还有别的 chunk 行也指向同一对象)→ 视为仍被引用。</li>
|
||||
* <li>{@code biz_task_scope_state.parsed_payload_json}/{@code state_json} 命中 > 0 行 →
|
||||
* 视为仍被引用(这两个字段不是 caller 自身行的常见持有者,命中即非自我引用)。</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>查询时同时把 caller 传入的 raw value 与 {@code pointer}、{@code "pointer"}(JSON 编码后
|
||||
* 的字符串形式)都纳入候选,覆盖以下两类常见 DB 字段写入形式:
|
||||
* <ul>
|
||||
* <li>{@code objectMapper.writeValueAsString(pointer)} 写入的 JSON 字符串包裹形式;</li>
|
||||
* <li>少数老路径直接写入裸 pointer 的形式(如 {@link #deleteReplacedPayloadIfNeeded})。</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>查询出现异常时保守返回 {@code true}(不删),由后续清理任务兜底。
|
||||
*/
|
||||
private boolean isStillReferenced(String pointer, String originalValue) {
|
||||
try {
|
||||
Set<String> 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<String> values = new ArrayList<>(candidates);
|
||||
Long chunkCount = taskChunkMapper.selectCount(new LambdaQueryWrapper<TaskChunkEntity>()
|
||||
.in(TaskChunkEntity::getPayloadJson, values));
|
||||
if (chunkCount != null && chunkCount > 1L) {
|
||||
return true;
|
||||
}
|
||||
Long scopeStateCount = taskScopeStateMapper.selectCount(new LambdaQueryWrapper<TaskScopeStateEntity>()
|
||||
.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;
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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}
|
||||
|
||||
18
backend-java/src/main/resources/db/V53__task_image_cache.sql
Normal file
18
backend-java/src/main/resources/db/V53__task_image_cache.sql
Normal file
@@ -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/通用 图片缩略图缓存,跨任务复用';
|
||||
@@ -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("<c r=\"J2\" t=\"e\" vm=\"0\"><v>#VALUE!</v></c>"));
|
||||
assertTrue(sheet.contains("<c r=\"K2\" t=\"e\" vm=\"1\"><v>#VALUE!</v></c>"));
|
||||
|
||||
String metadata = read(zip, "xl/metadata.xml");
|
||||
assertTrue(metadata.contains("<xlrd:rvb i=\"0\"/>"));
|
||||
assertTrue(metadata.contains("<xlrd:rvb i=\"1\"/>"));
|
||||
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", """
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<Types xmlns="http://schemas.openxmlformats.org/package/2006/content-types">
|
||||
<Default Extension="rels" ContentType="application/vnd.openxmlformats-package.relationships+xml"/>
|
||||
<Default Extension="xml" ContentType="application/xml"/>
|
||||
<Override PartName="/xl/workbook.xml" ContentType="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet.main+xml"/>
|
||||
<Override PartName="/xl/worksheets/sheet1.xml" ContentType="application/vnd.openxmlformats-officedocument.spreadsheetml.worksheet+xml"/>
|
||||
</Types>
|
||||
""");
|
||||
write(out, "_rels/.rels", """
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<Relationships xmlns="http://schemas.openxmlformats.org/package/2006/relationships">
|
||||
<Relationship Id="rId1" Type="http://schemas.openxmlformats.org/officeDocument/2006/relationships/officeDocument" Target="xl/workbook.xml"/>
|
||||
</Relationships>
|
||||
""");
|
||||
write(out, "xl/workbook.xml", """
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<workbook xmlns="http://schemas.openxmlformats.org/spreadsheetml/2006/main"
|
||||
xmlns:r="http://schemas.openxmlformats.org/officeDocument/2006/relationships">
|
||||
<sheets><sheet name="Sheet1" r:id="rId1" sheetId="1"/></sheets>
|
||||
</workbook>
|
||||
""");
|
||||
write(out, "xl/_rels/workbook.xml.rels", """
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<Relationships xmlns="http://schemas.openxmlformats.org/package/2006/relationships">
|
||||
<Relationship Id="rId1" Type="http://schemas.openxmlformats.org/officeDocument/2006/relationships/worksheet" Target="worksheets/sheet1.xml"/>
|
||||
</Relationships>
|
||||
""");
|
||||
write(out, "xl/worksheets/sheet1.xml", """
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<worksheet xmlns="http://schemas.openxmlformats.org/spreadsheetml/2006/main">
|
||||
<sheetData>
|
||||
<row r="2">
|
||||
<c r="J2" t="inlineStr"><is><t>main</t></is></c>
|
||||
<c r="K2" t="inlineStr"><is><t>puzzle</t></is></c>
|
||||
</row>
|
||||
</sheetData>
|
||||
</worksheet>
|
||||
""");
|
||||
}
|
||||
}
|
||||
|
||||
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();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user