feat(任务恢复): 客户端中断的任务自动续跑(保留失败记录 + 重新排队)

2026-09-18 任务 28587(跟价,uid 977 店铺「张美莺」)在客户端 09:51/09:57 被更新脚本
重启后中断,任务一直挂 RUNNING 到 stale 兜底;跟价循环还有第二重问题——子任务失败
即整个循环判 FAILED,一次客户端更新就让无限循环彻底停下。

契约:中断任务**保留失败记录**(用户看得到「因客户端重启中断」),同时由服务端自动
重新排队续跑,长任务不再因为一次更新整个白跑。

- V129:biz_file_task 加 resume_of_task_id / resume_attempt,
  biz_price_track_loop_run 加 resume_attempt(代数封顶用);
- 新增 TaskResumeService:扫描最近 30 分钟内因客户端中断而失败、未续跑过、代数未超限的
  任务,复制请求参数重新排队成 PENDING,交给已有兜底拉取通道(客户端每分钟 pull-pending)
  领走执行——不新造第二套派发机制。只对注册了 ClientTaskPullSpi 的模块生效
  (相似ASIN/采集/外观专利),上架/改价等写操作模块刻意排除,避免盲目重跑;
- 幂等:按 resume_of_task_id 反查,同一原任务不会重复排队;
- 挂在 stale-check 巡检线(每 2 分钟、有分布式锁),并用隔离壳包住异常:
  续跑失败绝不拖垮判死主流程(判死优先级更高);
- PriceTrackLoopRunService:子任务因「客户端异常中断」失败时不终止循环,
  清 active_task_id 后保持 RUNNING,客户端下次 dispatchNext 拿到同一店铺/同一轮,
  即原地续跑(页内已处理 ASIN 由服务端 skip_asins 去重,不会重复改价);
  连续重派超过 3 次才真正判失败,避免会话持续不可用时无限重开浏览器
  (28587 就是这样白烧了 5 小时)。子任务成功一轮后计数归零;
- application.yml:aiimage.task-resume.enabled 默认跟随 client-task-pull 开关
  (兜底拉取关着时续跑任务无人领取,只会积压被告死)。

测试:TaskResumeServiceTest 6 例(排队/幂等/开关/无实现/插入失败/无归属用户)、
PriceTrackLoopRunServiceTest 新增 4 例(中断续跑/达上限/非中断仍终止/成功归零),
连带更新 DeleteBrandStaleTaskServiceTest 与 TaskModuleCoverageTest 的构造参数。
This commit is contained in:
2026-09-18 15:42:07 +08:00
parent 3634ea1d62
commit bd359411a9
11 changed files with 619 additions and 3 deletions
@@ -27,6 +27,7 @@ import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import com.nanri.aiimage.modules.task.service.TaskDistributedLockService;
import com.nanri.aiimage.modules.task.service.TaskFileJobService;
import com.nanri.aiimage.modules.task.service.TaskHeartbeatPositionService;
import com.nanri.aiimage.modules.task.service.TaskResumeService;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
@@ -92,6 +93,8 @@ public class DeleteBrandStaleTaskService {
private final ShopDataCrawlTaskService shopDataCrawlTaskService;
private final TaskScopeStateMapper taskScopeStateMapper;
private final TaskHeartbeatPositionService taskHeartbeatPositionService;
/** 客户端中断任务的自动续跑(保留失败记录 + 重排队续跑任务,V129)。 */
private final TaskResumeService taskResumeService;
@Value("${aiimage.temp-dir.retention-hours:24}")
private long tempDirRetentionHours;
@@ -114,13 +117,16 @@ public class DeleteBrandStaleTaskService {
ShopMatchStaleCheckStats queryAsinStats = failStaleQueryAsinTasks();
ShopMatchStaleCheckStats withdrawStats = failStaleWithdrawTasks();
NoResultUploadStaleCheckStats noUploadStats = failNoResultUploadTasks();
// 客户端重启中断的任务:除了标失败(客户端上报,用户能看到原因),还要重新排队
// 一条 PENDING 续跑任务交给客户端兜底拉取执行,否则长任务一遇客户端更新就整个白跑
TaskResumeService.ResumeStats resumeStats = resumeInterruptedSafely();
for (Map.Entry<String, Runnable> delegated : delegatedStaleChecks().entrySet()) {
runModuleStaleCheck(delegated.getKey(), delegated.getValue());
}
// 周期每 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={}) 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={}",
stats.scannedTaskCount, stats.finalizedTaskCount, stats.failedTaskCount, stats.skippedTaskCount,
priceTrackStats.scannedTaskCount, priceTrackStats.finalizedTaskCount, priceTrackStats.failedTaskCount, priceTrackStats.skippedTaskCount,
shopMatchStats.scannedTaskCount, shopMatchStats.finalizedTaskCount, shopMatchStats.failedTaskCount, shopMatchStats.skippedTaskCount,
@@ -128,11 +134,26 @@ public class DeleteBrandStaleTaskService {
queryAsinStats.scannedTaskCount, queryAsinStats.finalizedTaskCount, queryAsinStats.failedTaskCount, queryAsinStats.skippedTaskCount,
withdrawStats.scannedTaskCount, withdrawStats.finalizedTaskCount, withdrawStats.failedTaskCount, withdrawStats.skippedTaskCount,
noUploadStats.scannedTaskCount, noUploadStats.failedTaskCount, noUploadStats.skippedTaskCount,
resumeStats.scannedTaskCount, resumeStats.resumedTaskCount, resumeStats.skippedTaskCount,
System.currentTimeMillis() - startedAt,
Thread.currentThread().getName());
}
}
/**
* 自动续跑巡检的隔离壳:续跑失败绝不能拖垮判死主流程。
* 判死是任务状态的兜底(不做会留下永远 RUNNING 的孤儿),优先级高于续跑;
* 这里失败只记日志并返回空统计,2 分钟后的下一轮自然重试。
*/
private TaskResumeService.ResumeStats resumeInterruptedSafely() {
try {
return taskResumeService.resumeInterruptedTasks();
} catch (Exception ex) {
log.warn("[task-resume] 自动续跑巡检失败(不影响本轮判死): {}", ex.getMessage(), ex);
return new TaskResumeService.ResumeStats();
}
}
/**
* 委派式陈旧判死:moduleType → 处理动作。
* {@code TaskModuleRegistry} 中 delegatedStaleCheck=true 的模块都必须在这里登记,
@@ -29,6 +29,8 @@ public class PriceTrackLoopRunEntity {
@TableField(value = "active_task_id", updateStrategy = FieldStrategy.ALWAYS)
private Long activeTaskId;
private Boolean stopRequested;
/** 因客户端中断自动重派当前轮的次数:封顶用,避免会话持续不可用时无限重派(V129)。 */
private Integer resumeAttempt;
private String errorMessage;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
@@ -39,6 +39,14 @@ public class PriceTrackLoopRunService {
private static final String STATUS_STOPPED = "STOPPED";
private static final String EXECUTION_MODE_FINITE = "FINITE";
private static final String EXECUTION_MODE_INFINITE = "INFINITE";
/** 客户端异常中断的任务 errorMessage 前缀(TaskHeartbeatService.markInterrupted 写入)。 */
private static final String CLIENT_INTERRUPT_ERROR_PREFIX = "客户端异常中断";
/**
* 中断后自动重派当前轮的次数上限。
* 会话持续不可用(账号被风控、紫鸟未就绪)时,不封顶会无限重派、无限重开浏览器——
* 2026-09-18 任务 28587 就是这样白烧了 5 小时。
*/
private static final int MAX_AUTO_RESUME_ATTEMPT = 3;
private final PriceTrackLoopRunMapper loopRunMapper;
private final FileTaskMapper fileTaskMapper;
@@ -259,6 +267,28 @@ public class PriceTrackLoopRunService {
entity.setActiveTaskId(null);
entity.setUpdatedAt(LocalDateTime.now());
if (STATUS_FAILED.equals(task.getStatus())) {
// 客户端重启中断(markInterrupted 写入的前缀)不终止循环:清空 active_task_id 后保持
// RUNNING,客户端下次 dispatchNext 会拿到**同一店铺、同一轮次**的 childTaskRequest
// 等于原地续跑——页内已处理的 ASIN 由服务端 skip_asins 去重,不会重复改价。
if (isClientInterrupt(task) && currentResumeAttempt(entity) < MAX_AUTO_RESUME_ATTEMPT) {
int attempt = currentResumeAttempt(entity) + 1;
entity.setStatus(STATUS_RUNNING);
entity.setErrorMessage(null);
entity.setFinishedAt(null);
entity.setResumeAttempt(attempt);
loopRunMapper.updateById(entity);
log.warn("[price-track-loop] 子任务因客户端中断失败,自动重派当前轮 loopRunId={} childTaskId={} "
+ "round={} shopIndex={} resumeAttempt={}/{} error={}",
entity.getId(), childTaskId, entity.getCurrentRound(), entity.getCurrentShopIndex(),
attempt, MAX_AUTO_RESUME_ATTEMPT, task.getErrorMessage());
return;
}
if (isClientInterrupt(task)) {
log.warn("[price-track-loop] 客户端中断续跑已达上限 {} 次,终止循环 loopRunId={} childTaskId={} "
+ "round={} shopIndex={}",
MAX_AUTO_RESUME_ATTEMPT, entity.getId(), childTaskId,
entity.getCurrentRound(), entity.getCurrentShopIndex());
}
entity.setStatus(STATUS_FAILED);
entity.setErrorMessage(task.getErrorMessage() == null || task.getErrorMessage().isBlank()
? "子任务执行失败"
@@ -269,6 +299,8 @@ public class PriceTrackLoopRunService {
entity.getId(), childTaskId, entity.getCurrentRound(), entity.getCurrentShopIndex(), entity.getErrorMessage());
return;
}
// 子任务成功 → 中断续跑计数归零(否则历史上的中断会一直占用封顶额度)
entity.setResumeAttempt(0);
List<PriceTrackMatchShopsVo.PriceTrackShopQueueItem> items = parseShops(entity);
if (items.isEmpty()) {
entity.setStatus(STATUS_FAILED);
@@ -300,6 +332,16 @@ public class PriceTrackLoopRunService {
entity.getId(), childTaskId, entity.getStatus(), entity.getCurrentRound(), entity.getCurrentShopIndex(), entity.getActiveTaskId());
}
/** 子任务失败原因是否为「客户端异常中断」(客户端重启上报,可自动重派续跑)。 */
private static boolean isClientInterrupt(FileTaskEntity task) {
String message = task == null ? null : task.getErrorMessage();
return message != null && message.startsWith(CLIENT_INTERRUPT_ERROR_PREFIX);
}
private static int currentResumeAttempt(PriceTrackLoopRunEntity entity) {
return entity.getResumeAttempt() == null ? 0 : entity.getResumeAttempt();
}
private void reconcileWithTerminalChild(PriceTrackLoopRunEntity entity) {
if (entity == null || entity.getActiveTaskId() == null || isTerminal(entity.getStatus())) {
return;
@@ -26,6 +26,10 @@ public class FileTaskEntity {
private String createdBy;
private Long userId;
private String ownerInstanceId;
/** 续跑来源任务 ID:本行是客户端中断后由服务端自动重排队的续跑任务时非空(V129)。 */
private Long resumeOfTaskId;
/** 续跑代数:0=原始任务,N=第 N 次自动续跑(封顶用,V129)。 */
private Integer resumeAttempt;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
private LocalDateTime finishedAt;
@@ -0,0 +1,164 @@
package com.nanri.aiimage.modules.task.service;
import cn.hutool.core.util.IdUtil;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.nanri.aiimage.modules.task.mapper.FileTaskMapper;
import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import com.nanri.aiimage.modules.task.spi.ClientTaskPullSpi;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
/**
* 客户端中断任务的自动续跑:**保留失败记录 + 重新排队一条 PENDING 续跑任务**。
*
* <p>场景:客户端被更新脚本/任务管理器杀掉时,{@code TaskHeartbeatService.markInterrupted} 会把在跑的任务
* 标 FAILED(原因「客户端异常中断: …」),用户能看到发生过什么;但长任务(采集/相似ASIN/外观专利/单次跟价)
* 一遇客户端重启就整个白跑,需要有人把它重新放回队列。本服务负责这件事。
*
* <p>续跑方式刻意复用已有链路:新建的任务落 PENDING,由 {@link TaskClientPullService} 的兜底拉取
* (客户端每分钟 {@code GET /api/tasks/pull-pending})领走执行——不再新造第二套派发机制。
* 只有注册了 {@link ClientTaskPullSpi} 的模块才会被续跑(这些模块的载荷能由服务端自行组装);
* 上架/改价等写操作模块不在其中,避免盲目重跑。
*
* <p>封顶 {@code max-attempt}:会话持续不可用(账号被风控、紫鸟未就绪)时,不封顶会无限重排队。
* 计数写在 {@code biz_file_task.resume_attempt}V129),随续跑任务代代递增。
*
* <p>幂等:同一原任务已存在续跑任务({@code resume_of_task_id} 反查)则跳过,重复扫描安全。
*/
@Slf4j
@Service
public class TaskResumeService {
private static final String STATUS_FAILED = "FAILED";
private static final String STATUS_PENDING = "PENDING";
/** 客户端重启中断的失败原因前缀,由 TaskHeartbeatService.markInterrupted 写入。 */
private static final String CLIENT_INTERRUPT_PREFIX = "客户端异常中断";
private final FileTaskMapper fileTaskMapper;
/** moduleType → 模块兜底载荷实现;只有注册了 SPI 的模块才能被自动续跑。 */
private final Map<String, ClientTaskPullSpi> resumeHandlers;
@Value("${aiimage.task-resume.enabled:true}")
private boolean enabled;
/** 续跑代数上限:达到后不再重排队,任务保持 FAILED 等人工介入。 */
@Value("${aiimage.task-resume.max-attempt:3}")
private int maxAttempt;
/** 只处理最近这段时间内中断的任务,避免开机扫描历史积压。 */
@Value("${aiimage.task-resume.window-minutes:30}")
private long windowMinutes;
@Value("${aiimage.task-resume.limit:20}")
private int limit;
public TaskResumeService(FileTaskMapper fileTaskMapper, List<ClientTaskPullSpi> pullSpiHandlers) {
this.fileTaskMapper = fileTaskMapper;
Map<String, ClientTaskPullSpi> index = new LinkedHashMap<>();
if (pullSpiHandlers != null) {
for (ClientTaskPullSpi handler : pullSpiHandlers) {
String moduleType = handler.moduleType();
if (moduleType == null || moduleType.isBlank()) {
continue;
}
index.put(moduleType.trim().toUpperCase(Locale.ROOT), handler);
}
}
this.resumeHandlers = Map.copyOf(index);
log.info("[task-resume] 可自动续跑模块注册完成 count={} modules={}", index.size(), index.keySet());
}
/** 扫描一轮:把最近因客户端中断而失败、且未续跑过的任务重新排队。 */
public ResumeStats resumeInterruptedTasks() {
ResumeStats stats = new ResumeStats();
if (!enabled) {
return stats;
}
if (resumeHandlers.isEmpty()) {
log.warn("[task-resume] 没有注册任何 ClientTaskPullSpi,跳过本轮续跑");
return stats;
}
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(Math.max(1L, windowMinutes));
int safeLimit = Math.max(1, limit);
List<FileTaskEntity> candidates = fileTaskMapper.selectList(new LambdaQueryWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getStatus, STATUS_FAILED)
.likeRight(FileTaskEntity::getErrorMessage, CLIENT_INTERRUPT_PREFIX)
.in(FileTaskEntity::getModuleType, resumeHandlers.keySet())
.lt(FileTaskEntity::getResumeAttempt, maxAttempt)
.ge(FileTaskEntity::getFinishedAt, cutoff)
.orderByAsc(FileTaskEntity::getId)
.last("limit " + safeLimit));
if (candidates.isEmpty()) {
return stats;
}
stats.scannedTaskCount = candidates.size();
for (FileTaskEntity original : candidates) {
if (original.getUserId() == null || original.getUserId() <= 0) {
log.warn("[task-resume] 原任务没有归属用户,跳过 taskId={}", original.getId());
stats.skippedTaskCount++;
continue;
}
// 幂等:同一原任务已经排过一次续跑就跳过(重复扫描 / 双实例竞态下不会重复建单)
if (hasResumeChild(original.getId())) {
stats.skippedTaskCount++;
continue;
}
FileTaskEntity resume = buildResumeTask(original);
try {
fileTaskMapper.insert(resume);
stats.resumedTaskCount++;
log.warn("[task-resume] 已自动重新排队 taskId={} moduleType={} userId={} 续跑任务={} 代数={}/{} 中断原因={}",
original.getId(), original.getModuleType(), original.getUserId(),
resume.getId(), resume.getResumeAttempt(), maxAttempt, original.getErrorMessage());
} catch (Exception ex) {
stats.skippedTaskCount++;
log.error("[task-resume] 重新排队失败 taskId={} moduleType={} err={}",
original.getId(), original.getModuleType(), ex.getMessage(), ex);
}
}
return stats;
}
/** 该原任务是否已经有续跑任务(反查 resume_of_task_id)。 */
private boolean hasResumeChild(Long originalTaskId) {
Long count = fileTaskMapper.selectCount(new LambdaQueryWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getResumeOfTaskId, originalTaskId));
return count != null && count > 0;
}
private FileTaskEntity buildResumeTask(FileTaskEntity original) {
LocalDateTime now = LocalDateTime.now();
FileTaskEntity resume = new FileTaskEntity();
resume.setTaskNo(original.getModuleType() + "-" + IdUtil.getSnowflakeNextIdStr());
resume.setModuleType(original.getModuleType());
resume.setTaskMode(original.getTaskMode());
resume.setStatus(STATUS_PENDING);
resume.setSourceFileCount(original.getSourceFileCount());
resume.setSuccessFileCount(0);
resume.setFailedFileCount(0);
// 原任务的请求参数整体带走:各模块的 ClientTaskPullSpi 从 request_json / 关联表还原执行参数
resume.setRequestJson(original.getRequestJson());
resume.setCreatedBy(original.getCreatedBy());
resume.setUserId(original.getUserId());
resume.setResumeOfTaskId(original.getId());
int attempt = original.getResumeAttempt() == null ? 0 : original.getResumeAttempt();
resume.setResumeAttempt(attempt + 1);
resume.setCreatedAt(now);
resume.setUpdatedAt(now);
return resume;
}
/** 单轮统计(合并进 stale-check 的 summary 日志,避免定期刷屏)。 */
public static final class ResumeStats {
public int scannedTaskCount;
public int resumedTaskCount;
public int skippedTaskCount;
}
}