fix(任务恢复): 停止意图优先于自动续跑 + 未覆盖模块的中断可见

自查上一提交时发现两处遗漏:

- 用户点了「停止循环」而子任务恰好因客户端中断失败时,续派逻辑跑在停止检查之前,
  会把循环从停止意图里拉回 RUNNING 自己转起来。改为停止请求优先:直接 STOPPED。
- 模块不支持自动续跑的中断任务此前完全无痕迹(用户只看到失败)。
  这类模块(上架/改价/审批/商品管理采集/跟价非循环任务)的续跑载荷需要用户在页面上
  选的执行参数(ziniao_version 等),而这份选择只存在于派发那一刻的浏览器里、没落到
  request_json —— 自动重排队会用错参数,所以它们**不纳入**续跑(白名单的设计原则就是
  「执行参数全在服务端」)。现补一条计数查询把这类中断数量带进 stale-check summary
  的 resume(u=N),运维据此人工重跑;将来把参数回写落库后即可纳入白名单。

测试:PriceTrackLoopRunServiceTest 新增「已请求停止时中断失败不续派而是停止」,
TaskResumeServiceTest 新增「统计不支持续跑的中断任务数」。
This commit is contained in:
2026-09-18 15:48:33 +08:00
parent bd359411a9
commit 986df86e89
5 changed files with 84 additions and 1 deletions
@@ -126,7 +126,7 @@ public class DeleteBrandStaleTaskService {
// 周期每 2 分钟一轮,各模块 summary 合并为单行,避免定期刷屏
// 注意:每段占位符数量必须与实参一致——此前每段 5 个占位符只传 4 个参数,
// 导致 withdraw 之后的取值整体错位、末尾 elapsedMs/thread 打成字面量
log.info("[stale-check] summary product-risk(s={} f={} x={} p={}) price-track(s={} f={} x={} p={}) shop-match(s={} f={} x={} p={}) patrol-delete(s={} f={} x={} p={}) query-asin(s={} f={} x={} p={}) withdraw(s={} f={} x={} p={}) no-upload(c={} f={} x={}) resume(s={} r={} k={}) elapsedMs={} thread={}",
log.info("[stale-check] summary product-risk(s={} f={} x={} p={}) price-track(s={} f={} x={} p={}) shop-match(s={} f={} x={} p={}) patrol-delete(s={} f={} x={} p={}) query-asin(s={} f={} x={} p={}) withdraw(s={} f={} x={} p={}) no-upload(c={} f={} x={}) resume(s={} r={} k={} u={}) elapsedMs={} thread={}",
stats.scannedTaskCount, stats.finalizedTaskCount, stats.failedTaskCount, stats.skippedTaskCount,
priceTrackStats.scannedTaskCount, priceTrackStats.finalizedTaskCount, priceTrackStats.failedTaskCount, priceTrackStats.skippedTaskCount,
shopMatchStats.scannedTaskCount, shopMatchStats.finalizedTaskCount, shopMatchStats.failedTaskCount, shopMatchStats.skippedTaskCount,
@@ -135,6 +135,7 @@ public class DeleteBrandStaleTaskService {
withdrawStats.scannedTaskCount, withdrawStats.finalizedTaskCount, withdrawStats.failedTaskCount, withdrawStats.skippedTaskCount,
noUploadStats.scannedTaskCount, noUploadStats.failedTaskCount, noUploadStats.skippedTaskCount,
resumeStats.scannedTaskCount, resumeStats.resumedTaskCount, resumeStats.skippedTaskCount,
resumeStats.unsupportedTaskCount,
System.currentTimeMillis() - startedAt,
Thread.currentThread().getName());
}
@@ -267,6 +267,11 @@ public class PriceTrackLoopRunService {
entity.setActiveTaskId(null);
entity.setUpdatedAt(LocalDateTime.now());
if (STATUS_FAILED.equals(task.getStatus())) {
// 用户已请求停止时绝不续派:停止意图优先于自动恢复(否则点完停止循环还会自己转起来)
if (Boolean.TRUE.equals(entity.getStopRequested())) {
markStopped(entity, null);
return;
}
// 客户端重启中断(markInterrupted 写入的前缀)不终止循环:清空 active_task_id 后保持
// RUNNING,客户端下次 dispatchNext 会拿到**同一店铺、同一轮次**的 childTaskRequest
// 等于原地续跑——页内已处理的 ASIN 由服务端 skip_asins 去重,不会重复改价。
@@ -96,8 +96,10 @@ public class TaskResumeService {
.orderByAsc(FileTaskEntity::getId)
.last("limit " + safeLimit));
if (candidates.isEmpty()) {
stats.unsupportedTaskCount = countUnsupportedInterrupts(cutoff);
return stats;
}
stats.unsupportedTaskCount = countUnsupportedInterrupts(cutoff);
stats.scannedTaskCount = candidates.size();
for (FileTaskEntity original : candidates) {
if (original.getUserId() == null || original.getUserId() <= 0) {
@@ -126,6 +128,28 @@ public class TaskResumeService {
return stats;
}
/**
* 统计「因客户端中断而失败、但模块不支持自动续跑」的任务数。
*
* <p>这些模块(上架/改价/审批/商品管理采集/跟价的非循环任务等)的续跑载荷需要**用户在页面上选的
* 执行参数**(如 ziniao_version),而这份选择只存在于派发那一刻的浏览器里、没落到服务端
* request_json —— 自动重排队会用错参数。因此它们只做**可见**:数量进巡检 summary
* 运维据此人工重跑;将来把这类参数回写落库后即可纳入续跑白名单。
*/
private int countUnsupportedInterrupts(LocalDateTime cutoff) {
try {
Long count = fileTaskMapper.selectCount(new LambdaQueryWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getStatus, STATUS_FAILED)
.likeRight(FileTaskEntity::getErrorMessage, CLIENT_INTERRUPT_PREFIX)
.notIn(FileTaskEntity::getModuleType, resumeHandlers.keySet())
.ge(FileTaskEntity::getFinishedAt, cutoff));
return count == null ? 0 : count.intValue();
} catch (Exception ex) {
log.warn("[task-resume] 统计不支持续跑的中断任务失败(忽略): {}", ex.getMessage());
return 0;
}
}
/** 该原任务是否已经有续跑任务(反查 resume_of_task_id)。 */
private boolean hasResumeChild(Long originalTaskId) {
Long count = fileTaskMapper.selectCount(new LambdaQueryWrapper<FileTaskEntity>()
@@ -160,5 +184,7 @@ public class TaskResumeService {
public int scannedTaskCount;
public int resumedTaskCount;
public int skippedTaskCount;
/** 模块不支持自动续跑的中断任务数(仅计数,供运维人工重跑)。 */
public int unsupportedTaskCount;
}
}
@@ -277,6 +277,44 @@ class PriceTrackLoopRunServiceTest {
assertEquals(0, loop.getResumeAttempt(), "成功一轮后中断计数归零,避免历史中断占用额度");
}
@Test
void 已请求停止时中断失败不续派而是停止() throws Exception {
ObjectMapper objectMapper = new ObjectMapper();
PriceTrackLoopRunMapper loopRunMapper = mock(PriceTrackLoopRunMapper.class);
FileTaskMapper fileTaskMapper = mock(FileTaskMapper.class);
PriceTrackLoopRunService service = new PriceTrackLoopRunService(
loopRunMapper, fileTaskMapper, objectMapper, mock(ZiniaoShopSwitchService.class));
PriceTrackLoopRunEntity loop = new PriceTrackLoopRunEntity();
loop.setId(2699L);
loop.setUserId(977L);
loop.setStatus("RUNNING");
loop.setExecutionMode("INFINITE");
loop.setCurrentRound(1);
loop.setCurrentShopIndex(0);
loop.setActiveTaskId(28587L);
loop.setResumeAttempt(0);
loop.setStopRequested(true);
loop.setShopsJson(objectMapper.writeValueAsString(List.of(shop("张美莺"))));
FileTaskEntity task = new FileTaskEntity();
task.setId(28587L);
task.setUserId(977L);
task.setModuleType("PRICE_TRACK");
task.setStatus("FAILED");
task.setErrorMessage("客户端异常中断: 客户端重启恢复上报");
task.setRequestJson(objectMapper.writeValueAsString(Map.of(
"loopRunId", 2699, "roundIndex", 1, "shopIndex", 0)));
when(loopRunMapper.selectList(any())).thenReturn(List.of(loop));
when(fileTaskMapper.selectById(28587L)).thenReturn(task);
service.syncLoopRunAfterChildTerminal(28587L);
assertEquals("STOPPED", loop.getStatus(), "停止意图优先于自动续跑");
assertEquals(0, loop.getResumeAttempt());
}
private PriceTrackMatchShopsVo.PriceTrackShopQueueItem shop(String name) {
PriceTrackMatchShopsVo.PriceTrackShopQueueItem item = new PriceTrackMatchShopsVo.PriceTrackShopQueueItem();
item.setShopName(name);
@@ -139,6 +139,19 @@ class TaskResumeServiceTest {
assertEquals(0, stats.resumedTaskCount);
}
@Test
void 无候选时统计不支持续跑的中断任务数() {
FileTaskMapper mapper = mock(FileTaskMapper.class);
when(mapper.selectList(any())).thenReturn(List.of());
when(mapper.selectCount(any())).thenReturn(3L);
TaskResumeService service = service(mapper, List.of(spi("SIMILAR_ASIN")), true, 3);
TaskResumeService.ResumeStats stats = service.resumeInterruptedTasks();
assertEquals(0, stats.resumedTaskCount);
assertEquals(3, stats.unsupportedTaskCount, "不支持续跑的模块也要可见,供运维人工重跑");
}
@Test
void 没有归属用户的中断任务跳过() {
FileTaskMapper mapper = mock(FileTaskMapper.class);