From 3634ea1d62304c0102380bbd80d38990508667fe Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=BB=84=E8=87=AA=E8=BE=BE?= <980324341@qq.com> Date: Thu, 17 Sep 2026 23:08:15 +0800 Subject: [PATCH] =?UTF-8?q?feat(=E6=92=9E=E6=AC=BE):=20=E9=87=87=E9=9B=86?= =?UTF-8?q?=E6=98=8E=E7=BB=86=E8=90=BD=E5=BA=93=E5=90=8E=E8=87=AA=E5=8A=A8?= =?UTF-8?q?=E8=A7=A6=E5=8F=91=E9=87=8D=E6=89=AB=EF=BC=8C=E9=87=8D=E5=A4=8D?= =?UTF-8?q?=E6=A3=80=E6=9F=A5=E6=97=A0=E9=9C=80=E7=AD=89=E6=AC=A1=E6=97=A5?= =?UTF-8?q?=2000:00?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - shopdatacrawl 新增 DuplicateCheckRefreshPort 端口,TaskService 在明细落库成功后请求重扫 - shopduplicatecheck 新增 DuplicateCheckRefreshScheduler:10s 合并窗口去抖、单飞、 扫描锁忙自动重试(最多 6 次),失败不影响采集归档 - 复用现有 SCAN_LOCK,与定时扫描、手动「重新分析」互斥;新增 5 个调度器单测 --- .../service/ShopDataCrawlTaskService.java | 9 ++ .../spi/DuplicateCheckRefreshPort.java | 17 +++ .../ShopDataDuplicateCheckScanService.java | 34 ++++- .../DuplicateCheckRefreshScheduler.java | 138 ++++++++++++++++++ .../service/ShopDataCrawlChunkUpsertTest.java | 2 + .../service/ShopDataCrawlCleanupTest.java | 2 + ...ShopDataCrawlDailyFileIncrementalTest.java | 2 + .../ShopDataCrawlDailyFileJobSplitTest.java | 2 + .../ShopDataCrawlDailyFileLockTest.java | 2 + .../ShopDataCrawlLightweightProgressTest.java | 2 + .../service/ShopDataCrawlOwnerColumnTest.java | 2 + .../ShopDataCrawlProgressQueryTest.java | 2 + .../service/ShopDataCrawlRowDedupKeyTest.java | 2 + .../ShopDataCrawlScopeCounterTest.java | 2 + .../service/ShopDataCrawlScopeMergeTest.java | 2 + .../ShopDataCrawlTaskServiceChunkTest.java | 2 + .../DuplicateCheckRefreshSchedulerTest.java | 123 ++++++++++++++++ 17 files changed, 343 insertions(+), 2 deletions(-) create mode 100644 backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/spi/DuplicateCheckRefreshPort.java create mode 100644 backend-java/src/main/java/com/nanri/aiimage/modules/shopduplicatecheck/service/support/DuplicateCheckRefreshScheduler.java create mode 100644 backend-java/src/test/java/com/nanri/aiimage/modules/shopduplicatecheck/service/support/DuplicateCheckRefreshSchedulerTest.java diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskService.java index 42b47c62..e59bf711 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskService.java @@ -27,6 +27,7 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlTaskBatchVo import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlCreateTaskVo; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlHistoryVo; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlResultItemVo; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.common.model.vo.ProductRiskDashboardVo; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; @@ -121,6 +122,7 @@ public class ShopDataCrawlTaskService { private final InstanceMetadata instanceMetadata; private final ShopDataCrawlDailyFileService dailyFileService; private final ShopDataCrawlItemStoreService shopDataCrawlItemStoreService; + private final DuplicateCheckRefreshPort duplicateCheckRefreshPort; private final PlatformTransactionManager transactionManager; private final TaskProgressLightAssembler taskProgressLightAssembler; @@ -2116,6 +2118,13 @@ public class ShopDataCrawlTaskService { shopDataCrawlItemStoreService.saveShopBatchFromSnapshot( snapshot.getShopName(), itemBatchDate, accumulatedItems, snapshot.getResultId(), task.getId(), baseDailyFile == null ? null : baseDailyFile.getId()); + // 明细已落库:请求撞款重扫(异步合并执行,不阻塞归档;端口契约保证不抛错) + try { + duplicateCheckRefreshPort.requestRefresh("shop-data-crawl:" + snapshot.getShopName()); + } catch (RuntimeException ex) { + log.warn("[shop-data-crawl] 请求撞款重扫失败(忽略,不影响归档) shop={} msg={}", + snapshot.getShopName(), ex.getMessage()); + } int rowCount = excelAssemblyService.writeWorkbook(outputXlsx, accumulatedItems); String objectKey = ossStorageService.uploadResultFile(outputXlsx, MODULE_TYPE); if (blank(objectKey)) { diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/spi/DuplicateCheckRefreshPort.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/spi/DuplicateCheckRefreshPort.java new file mode 100644 index 00000000..f4530224 --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopdatacrawl/spi/DuplicateCheckRefreshPort.java @@ -0,0 +1,17 @@ +package com.nanri.aiimage.modules.shopdatacrawl.spi; + +/** + * 采集明细就绪后的撞款重扫触发端口(2026-09:店铺数据采集落库后即时刷新重复检查)。 + * + *

实现方在 shopduplicatecheck 模块({@code ShopDataDuplicateCheckScanService})。 + * 契约:实现必须异步执行、去抖合并,不得阻塞调用方、不得向外抛出异常。 + */ +public interface DuplicateCheckRefreshPort { + + /** + * 请求一次撞款重扫(异步;合并窗口内的多次触发聚合为一次扫描)。 + * + * @param reason 触发来源,仅用于日志排查 + */ + void requestRefresh(String reason); +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopduplicatecheck/service/ShopDataDuplicateCheckScanService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopduplicatecheck/service/ShopDataDuplicateCheckScanService.java index 22bc5955..acee19f2 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/shopduplicatecheck/service/ShopDataDuplicateCheckScanService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopduplicatecheck/service/ShopDataDuplicateCheckScanService.java @@ -9,6 +9,7 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlCountryRes import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlRowDto; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlResultItemVo; import com.nanri.aiimage.modules.shopdatacrawl.service.ShopDataCrawlItemStoreService; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.shopduplicatecheck.mapper.ShopDataDuplicateScanMapper; import com.nanri.aiimage.modules.shopduplicatecheck.mapper.ShopDuplicateCheckItemMapper; import com.nanri.aiimage.modules.shopduplicatecheck.mapper.ShopDuplicateCheckSourceMapper; @@ -21,6 +22,7 @@ import com.nanri.aiimage.modules.shopduplicatecheck.model.entity.ShopDataDuplica import com.nanri.aiimage.modules.shopduplicatecheck.model.payload.DuplicateScanPayload; import com.nanri.aiimage.modules.shopduplicatecheck.model.payload.DuplicateScanSummary; import com.nanri.aiimage.modules.shopduplicatecheck.service.support.DuplicateCheckAggregator; +import com.nanri.aiimage.modules.shopduplicatecheck.service.support.DuplicateCheckRefreshScheduler; import com.nanri.aiimage.modules.shopduplicatecheck.service.support.DuplicateCheckWorkbookParser; import com.nanri.aiimage.modules.shopduplicatecheck.service.support.RawRow; import com.nanri.aiimage.modules.shopduplicatecheck.service.support.ShopParsed; @@ -43,14 +45,14 @@ import java.util.Set; import java.util.concurrent.atomic.AtomicLong; /** - * 撞款扫描:每日 00:00 定时 + force 同步全量重扫。 + * 撞款扫描:每日 00:00 定时 + force 同步全量重扫 + 采集落库触发的异步合并重扫。 * 数据源 = 采集明细表(biz_shop_data_crawl_item,采集先落库再更新文件), * 直接查库聚合 shops/items/summary 落库 shop_data_duplicate_scan(输出契约不变)。 * 双实例通过 Redis 分布式锁防重;扫描失败落 FAILED 行并上抛,force 场景由端点转 409/500。 */ @Service @Slf4j -public class ShopDataDuplicateCheckScanService { +public class ShopDataDuplicateCheckScanService implements DuplicateCheckRefreshPort { public static final String SCAN_LOCK = "shop-data-duplicate-check:scan"; static final DateTimeFormatter SCANNED_AT_FORMAT = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); @@ -71,6 +73,9 @@ public class ShopDataDuplicateCheckScanService { private final AtomicLong cachedRowId = new AtomicLong(-1L); private volatile CachedScan cachedScan; + /** 采集落库触发的异步合并重扫调度器(单飞 + 去抖 + 锁忙重试)。 */ + private final DuplicateCheckRefreshScheduler refreshScheduler; + @Autowired public ShopDataDuplicateCheckScanService(ShopDataDuplicateScanMapper scanMapper, ShopDuplicateCheckSourceMapper sourceMapper, @@ -86,6 +91,7 @@ public class ShopDataDuplicateCheckScanService { this.objectMapper = objectMapper; this.itemMapper = itemMapper; this.itemStoreService = itemStoreService; + this.refreshScheduler = new DuplicateCheckRefreshScheduler(this::runRefreshOnce); } /** 读侧视图:scanned_at 为最新 SUCCESS 行 created_at(yyyy-MM-dd HH:mm:ss)。 */ @@ -122,6 +128,30 @@ public class ShopDataDuplicateCheckScanService { } } + /** 采集明细落库后的重扫请求(端口实现):异步合并执行,不阻塞、不抛错。 */ + @Override + public void requestRefresh(String reason) { + refreshScheduler.request(reason); + } + + /** 调度器单次扫描动作:锁被占返回 LOCK_BUSY 供其重试;失败只记日志(FAILED 行已落库)。 */ + private DuplicateCheckRefreshScheduler.Outcome runRefreshOnce() { + try { + scanNow(); + return DuplicateCheckRefreshScheduler.Outcome.DONE; + } catch (BusinessException ex) { + if (ex.getCode() != null && ex.getCode() == 409) { + log.info("[shop-duplicate-check] 自动重扫未执行:其它扫描进行中 msg={}", ex.getMessage()); + return DuplicateCheckRefreshScheduler.Outcome.LOCK_BUSY; + } + log.warn("[shop-duplicate-check] 自动重扫失败 code={} msg={}", ex.getCode(), ex.getMessage()); + return DuplicateCheckRefreshScheduler.Outcome.FAILED; + } catch (Exception ex) { + log.error("[shop-duplicate-check] 自动重扫异常", ex); + return DuplicateCheckRefreshScheduler.Outcome.FAILED; + } + } + /** 最新 SUCCESS 扫描视图;无扫描结果返回 null。 */ public DuplicateScanView loadLatest() { ScanLightRowDto light = scanMapper.selectLatestLightRow(); diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/shopduplicatecheck/service/support/DuplicateCheckRefreshScheduler.java b/backend-java/src/main/java/com/nanri/aiimage/modules/shopduplicatecheck/service/support/DuplicateCheckRefreshScheduler.java new file mode 100644 index 00000000..43492621 --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/shopduplicatecheck/service/support/DuplicateCheckRefreshScheduler.java @@ -0,0 +1,138 @@ +package com.nanri.aiimage.modules.shopduplicatecheck.service.support; + +import com.nanri.aiimage.common.util.ThreadPools; +import lombok.extern.slf4j.Slf4j; + +import java.util.concurrent.ExecutorService; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Supplier; + +/** + * 采集落库触发的撞款重扫调度器:合并窗口去抖 + 单飞 + 锁忙重试。 + * + *

语义:{@link #request} 只置位并异步执行,永不阻塞调用方、永不向外抛错; + * 合并窗口内的多次触发聚合为一次扫描;扫描动作执行期间到达的触发在下一轮执行; + * 扫描因分布式锁被占用未执行({@link Outcome#LOCK_BUSY})时按固定间隔重试有限次。 + */ +@Slf4j +public class DuplicateCheckRefreshScheduler { + + /** 单次扫描动作的终态:完成 / 锁被占(可重试)/ 失败(不重试,等下次触发或定时扫描)。 */ + public enum Outcome { + DONE, LOCK_BUSY, FAILED + } + + private static final long DEFAULT_DEBOUNCE_MILLIS = 10_000L; + private static final long DEFAULT_LOCK_RETRY_MILLIS = 20_000L; + private static final int DEFAULT_MAX_LOCK_RETRIES = 6; + + private final Supplier scanAction; + private final long debounceMillis; + private final long lockRetryMillis; + private final int maxLockRetries; + private final ExecutorService executor; + private final AtomicBoolean pending = new AtomicBoolean(false); + private final AtomicBoolean running = new AtomicBoolean(false); + + public DuplicateCheckRefreshScheduler(Supplier scanAction) { + this(scanAction, DEFAULT_DEBOUNCE_MILLIS, DEFAULT_LOCK_RETRY_MILLIS, DEFAULT_MAX_LOCK_RETRIES); + } + + /** 测试用:注入更短的窗口与重试参数。 */ + DuplicateCheckRefreshScheduler(Supplier scanAction, long debounceMillis, + long lockRetryMillis, int maxLockRetries) { + this.scanAction = scanAction; + this.debounceMillis = Math.max(0L, debounceMillis); + this.lockRetryMillis = Math.max(0L, lockRetryMillis); + this.maxLockRetries = Math.max(0, maxLockRetries); + this.executor = ThreadPools.boundedFixed("shop-dup-refresh", 1, 8); + } + + /** 请求一次重扫(异步、去抖合并)。调用方不被阻塞,也不会收到异常。 */ + public void request(String reason) { + pending.set(true); + if (running.compareAndSet(false, true)) { + submit(reason); + } + } + + private void submit(String reason) { + try { + log.info("[shop-duplicate-check] 触发撞款重扫(异步合并执行,窗口={}ms) reason={}", debounceMillis, reason); + executor.execute(this::drain); + } catch (Exception ex) { + // 提交失败(如线程池拒绝)时复位单飞标记,避免后续触发被永久吞掉 + running.set(false); + log.warn("[shop-duplicate-check] 撞款重扫任务提交失败 reason={} msg={}", reason, ex.getMessage()); + } + } + + private void drain() { + try { + while (true) { + // 合并窗口:窗口内到达的多次触发聚合为同一轮扫描 + if (!sleepQuietly(debounceMillis)) { + return; + } + if (!pending.compareAndSet(true, false)) { + return; + } + int lockRetries = 0; + while (true) { + Outcome outcome = runOnceSafely(); + if (outcome != Outcome.LOCK_BUSY) { + break; + } + if (lockRetries >= maxLockRetries) { + log.warn("[shop-duplicate-check] 撞款重扫连续 {} 次未取得扫描锁,放弃本轮(等待下次触发或定时扫描)", + lockRetries + 1); + break; + } + lockRetries++; + log.info("[shop-duplicate-check] 撞款重扫未取得扫描锁,{}ms 后重试(第 {}/{} 次)", + lockRetryMillis, lockRetries, maxLockRetries); + if (!sleepQuietly(lockRetryMillis)) { + return; + } + } + } + } finally { + running.set(false); + // 竞态兜底:running 复位前到达的触发可能没能提交,补一次 + if (pending.get() && running.compareAndSet(false, true)) { + submit("race-guard"); + } + } + } + + /** 执行一次扫描动作;动作自身异常也被吸收(调度器对外零抛出)。 */ + private Outcome runOnceSafely() { + long startedAt = System.currentTimeMillis(); + try { + Outcome outcome = scanAction.get(); + long elapsed = System.currentTimeMillis() - startedAt; + if (outcome == Outcome.DONE) { + log.info("[shop-duplicate-check] 采集后自动重扫完成 耗时={}ms", elapsed); + } else if (outcome == Outcome.FAILED) { + log.warn("[shop-duplicate-check] 采集后自动重扫失败 耗时={}ms", elapsed); + } + return outcome == null ? Outcome.FAILED : outcome; + } catch (Exception ex) { + log.error("[shop-duplicate-check] 采集后自动重扫异常", ex); + return Outcome.FAILED; + } + } + + private static boolean sleepQuietly(long millis) { + if (millis <= 0) { + return true; + } + try { + Thread.sleep(millis); + return true; + } catch (InterruptedException ex) { + Thread.currentThread().interrupt(); + return false; + } + } +} diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlChunkUpsertTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlChunkUpsertTest.java index 06722872..bac6fab2 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlChunkUpsertTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlChunkUpsertTest.java @@ -11,6 +11,7 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlCountryRes import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlRowDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlShopPayloadDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlSubmitResultRequest; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -137,6 +138,7 @@ class ShopDataCrawlChunkUpsertTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlCleanupTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlCleanupTest.java index 081b4bb4..ce8c6105 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlCleanupTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlCleanupTest.java @@ -16,6 +16,7 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlSubmitResu import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyFileEntity; import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyMemberEntity; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlResultItemVo; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -174,6 +175,7 @@ class ShopDataCrawlCleanupTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); ReflectionTestUtils.setField(service, "staleTimeoutMinutes", 30L); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileIncrementalTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileIncrementalTest.java index 6fd08a98..a2e54dcf 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileIncrementalTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileIncrementalTest.java @@ -11,6 +11,7 @@ import com.nanri.aiimage.modules.file.service.oss.OssStorageService; import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyFileEntity; import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyMemberEntity; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlResultItemVo; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -146,6 +147,7 @@ class ShopDataCrawlDailyFileIncrementalTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileJobSplitTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileJobSplitTest.java index 65a8961d..dd5fc022 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileJobSplitTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileJobSplitTest.java @@ -11,6 +11,7 @@ import com.nanri.aiimage.modules.file.service.oss.OssStorageService; import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyFileEntity; import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyMemberEntity; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlResultItemVo; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -152,6 +153,7 @@ class ShopDataCrawlDailyFileJobSplitTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileLockTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileLockTest.java index c8232388..29c7c292 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileLockTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlDailyFileLockTest.java @@ -11,6 +11,7 @@ import com.nanri.aiimage.modules.file.service.oss.OssStorageService; import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyFileEntity; import com.nanri.aiimage.modules.shopdatacrawl.model.entity.ShopDataCrawlDailyMemberEntity; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlResultItemVo; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -152,6 +153,7 @@ class ShopDataCrawlDailyFileLockTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlLightweightProgressTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlLightweightProgressTest.java index a179ca4a..12c0e2c3 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlLightweightProgressTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlLightweightProgressTest.java @@ -12,6 +12,7 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlRowDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlShopPayloadDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlSubmitResultRequest; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlResultItemVo; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -134,6 +135,7 @@ class ShopDataCrawlLightweightProgressTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlOwnerColumnTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlOwnerColumnTest.java index 3900b757..89fefad3 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlOwnerColumnTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlOwnerColumnTest.java @@ -12,6 +12,7 @@ import com.nanri.aiimage.modules.file.service.oss.OssStorageService; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlCreateTaskRequest; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlTaskItemDto; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlCreateTaskVo; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -139,6 +140,7 @@ class ShopDataCrawlOwnerColumnTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlProgressQueryTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlProgressQueryTest.java index 59c6d97c..23d4b499 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlProgressQueryTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlProgressQueryTest.java @@ -11,6 +11,7 @@ import com.nanri.aiimage.modules.file.service.oss.OssStorageService; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlHistoryVo; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlResultItemVo; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlTaskBatchVo; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -146,6 +147,7 @@ class ShopDataCrawlProgressQueryTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlRowDedupKeyTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlRowDedupKeyTest.java index f13d21b3..0eb54192 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlRowDedupKeyTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlRowDedupKeyTest.java @@ -12,6 +12,7 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlRowDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlShopPayloadDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlSubmitResultRequest; import com.nanri.aiimage.modules.shopdatacrawl.model.vo.ShopDataCrawlResultItemVo; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -133,6 +134,7 @@ class ShopDataCrawlRowDedupKeyTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlScopeCounterTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlScopeCounterTest.java index 1b943e85..638184a4 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlScopeCounterTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlScopeCounterTest.java @@ -11,6 +11,7 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlCountryRes import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlRowDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlShopPayloadDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlSubmitResultRequest; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -134,6 +135,7 @@ class ShopDataCrawlScopeCounterTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlScopeMergeTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlScopeMergeTest.java index ebf67626..25387391 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlScopeMergeTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlScopeMergeTest.java @@ -11,6 +11,7 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlCountryRes import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlRowDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlShopPayloadDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlSubmitResultRequest; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -135,6 +136,7 @@ class ShopDataCrawlScopeMergeTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskServiceChunkTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskServiceChunkTest.java index b20391fa..b8ed1f7e 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskServiceChunkTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlTaskServiceChunkTest.java @@ -11,6 +11,7 @@ import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlCountryRes import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlRowDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlShopPayloadDto; import com.nanri.aiimage.modules.shopdatacrawl.model.dto.ShopDataCrawlSubmitResultRequest; +import com.nanri.aiimage.modules.shopdatacrawl.spi.DuplicateCheckRefreshPort; import com.nanri.aiimage.modules.task.mapper.FileResultMapper; import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; @@ -129,6 +130,7 @@ class ShopDataCrawlTaskServiceChunkTest { instanceMetadata, dailyFileService, mock(ShopDataCrawlItemStoreService.class), + mock(DuplicateCheckRefreshPort.class), null, mock(TaskProgressLightAssembler.class)); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopduplicatecheck/service/support/DuplicateCheckRefreshSchedulerTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopduplicatecheck/service/support/DuplicateCheckRefreshSchedulerTest.java new file mode 100644 index 00000000..0c11a547 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopduplicatecheck/service/support/DuplicateCheckRefreshSchedulerTest.java @@ -0,0 +1,123 @@ +package com.nanri.aiimage.modules.shopduplicatecheck.service.support; + +import org.junit.jupiter.api.Test; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.BooleanSupplier; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assertions.fail; + +/** + * 撞款重扫调度器:合并窗口去抖(突发触发聚为一次扫描)、单飞、 + * 锁忙有限重试、动作异常吸收(对外零抛出)。 + */ +class DuplicateCheckRefreshSchedulerTest { + + private static final long AWAIT_TIMEOUT_MILLIS = 5_000L; + + @Test + void request_mergesBurstIntoSingleScan() throws Exception { + AtomicInteger scans = new AtomicInteger(); + CountDownLatch firstScan = new CountDownLatch(1); + DuplicateCheckRefreshScheduler scheduler = new DuplicateCheckRefreshScheduler(() -> { + scans.incrementAndGet(); + firstScan.countDown(); + return DuplicateCheckRefreshScheduler.Outcome.DONE; + }, 150L, 30L, 3); + + scheduler.request("burst-1"); + Thread.sleep(10L); + scheduler.request("burst-2"); + Thread.sleep(10L); + scheduler.request("burst-3"); + + assertTrue(firstScan.await(AWAIT_TIMEOUT_MILLIS, TimeUnit.MILLISECONDS), "合并窗口内的触发应执行扫描"); + Thread.sleep(300L); + assertEquals(1, scans.get(), "合并窗口内多次触发只执行一次扫描"); + } + + @Test + void request_retriesWhileLockBusyThenSucceeds() throws Exception { + AtomicInteger attempts = new AtomicInteger(); + CountDownLatch succeeded = new CountDownLatch(1); + DuplicateCheckRefreshScheduler scheduler = new DuplicateCheckRefreshScheduler(() -> { + if (attempts.incrementAndGet() <= 2) { + return DuplicateCheckRefreshScheduler.Outcome.LOCK_BUSY; + } + succeeded.countDown(); + return DuplicateCheckRefreshScheduler.Outcome.DONE; + }, 10L, 30L, 5); + + scheduler.request("retry"); + + assertTrue(succeeded.await(AWAIT_TIMEOUT_MILLIS, TimeUnit.MILLISECONDS), "锁忙重试后应成功执行"); + Thread.sleep(100L); + assertEquals(3, attempts.get(), "2 次锁忙 + 1 次成功"); + } + + @Test + void request_givesUpAfterMaxLockRetries() throws Exception { + AtomicInteger attempts = new AtomicInteger(); + CountDownLatch firstAttempt = new CountDownLatch(1); + DuplicateCheckRefreshScheduler scheduler = new DuplicateCheckRefreshScheduler(() -> { + attempts.incrementAndGet(); + firstAttempt.countDown(); + return DuplicateCheckRefreshScheduler.Outcome.LOCK_BUSY; + }, 10L, 20L, 2); + + scheduler.request("always-busy"); + + assertTrue(firstAttempt.await(AWAIT_TIMEOUT_MILLIS, TimeUnit.MILLISECONDS)); + awaitUntil(() -> attempts.get() >= 3, "应完成初试 + 2 次重试"); + Thread.sleep(200L); + assertEquals(3, attempts.get(), "超过重试上限后放弃,不再执行"); + } + + @Test + void request_absorbsActionFailure() throws Exception { + AtomicInteger attempts = new AtomicInteger(); + DuplicateCheckRefreshScheduler scheduler = new DuplicateCheckRefreshScheduler(() -> { + attempts.incrementAndGet(); + throw new IllegalStateException("模拟扫描动作异常"); + }, 10L, 20L, 1); + + assertDoesNotThrow(() -> scheduler.request("boom"), "request 不得向调用方抛错"); + awaitUntil(() -> attempts.get() >= 1, "动作应被执行"); + Thread.sleep(100L); + assertEquals(1, attempts.get(), "动作异常视为失败,不做锁忙重试"); + } + + @Test + void request_afterPreviousCycleAllowsNewScan() throws Exception { + AtomicInteger scans = new AtomicInteger(); + CountDownLatch twoScans = new CountDownLatch(2); + DuplicateCheckRefreshScheduler scheduler = new DuplicateCheckRefreshScheduler(() -> { + scans.incrementAndGet(); + twoScans.countDown(); + return DuplicateCheckRefreshScheduler.Outcome.DONE; + }, 20L, 20L, 2); + + scheduler.request("first"); + awaitUntil(() -> scans.get() >= 1, "首轮扫描应执行"); + scheduler.request("second"); + + assertTrue(twoScans.await(AWAIT_TIMEOUT_MILLIS, TimeUnit.MILLISECONDS), "新触发应再执行一次扫描"); + assertEquals(2, scans.get()); + } + + private static void awaitUntil(BooleanSupplier condition, String message) throws InterruptedException { + long deadline = System.currentTimeMillis() + AWAIT_TIMEOUT_MILLIS; + while (System.currentTimeMillis() < deadline) { + if (condition.getAsBoolean()) { + return; + } + Thread.sleep(10L); + } + fail(message); + } +}