fix(任务中断): 标失败时同步模块缓存与进度快照,修"处理中"卡死

客户端重启上报中断只改了 file_task,模块缓存/进度快照仍是 RUNNING,
导致 progress/batch 继续回报 RUNNING:前端任务面板永远"处理中"并阻塞该工具后续任务
(2026-09-13 真机:店铺数据采集 28094 history=FAILED 但 progress=RUNNING,界面卡住)。
markInterrupted 现在同时刷新模块缓存(saveFileTaskCache)与进度快照(新增 markTerminal)。
This commit is contained in:
2026-09-13 15:16:14 +08:00
parent a4f60ef21c
commit 3d208ea0e5
2 changed files with 38 additions and 0 deletions
@@ -54,6 +54,7 @@ public class TaskHeartbeatService {
private static final String MODULE_BRAND = "BRAND"; private static final String MODULE_BRAND = "BRAND";
private final FileTaskMapper fileTaskMapper; private final FileTaskMapper fileTaskMapper;
private final TaskProgressSnapshotService taskProgressSnapshotService;
private final BrandCrawlTaskMapper brandCrawlTaskMapper; private final BrandCrawlTaskMapper brandCrawlTaskMapper;
private final ProductRiskTaskCacheService productRiskTaskCacheService; private final ProductRiskTaskCacheService productRiskTaskCacheService;
private final PublishTaskService publishTaskService; private final PublishTaskService publishTaskService;
@@ -142,6 +143,18 @@ public class TaskHeartbeatService {
if (updated > 0) { if (updated > 0) {
log.warn("[task-interrupted] file task marked failed by client restart taskId={} moduleType={} reason={}", log.warn("[task-interrupted] file task marked failed by client restart taskId={} moduleType={} reason={}",
taskId, fileTask.getModuleType(), safeReason); taskId, fileTask.getModuleType(), safeReason);
// 同步模块缓存与进度快照:progress/batch 读缓存/快照,不刷新会让前端
// 一直按 RUNNING 渲染(任务面板永远"处理中"并阻塞该工具后续任务)。
try {
fileTask.setStatus("FAILED");
fileTask.setErrorMessage(safeReason);
fileTask.setFinishedAt(LocalDateTime.now());
saveFileTaskCache(fileTask.getModuleType(), fileTask);
taskProgressSnapshotService.markTerminal(taskId, fileTask.getModuleType(), "FAILED", safeReason);
} catch (Exception cacheEx) {
log.warn("[task-interrupted] refresh cache/snapshot failed taskId={} err={}",
taskId, cacheEx.getMessage());
}
return TaskHeartbeatVo.notAlive(fileTask.getModuleType(), "FAILED", "marked failed"); return TaskHeartbeatVo.notAlive(fileTask.getModuleType(), "FAILED", "marked failed");
} }
log.info("[task-interrupted] file task not in RUNNING, skipped taskId={} status={}", taskId, status); log.info("[task-interrupted] file task not in RUNNING, skipped taskId={} status={}", taskId, status);
@@ -157,6 +157,31 @@ public class TaskProgressSnapshotService {
lastWriteAtMillis.remove(cacheKey(taskId, moduleType)); lastWriteAtMillis.remove(cacheKey(taskId, moduleType));
} }
/**
* 任务被外部置为终态(客户端中断上报 / stale 兜底修复)时同步进度快照。
* 不同步的话,progress/batch 仍按快照里的 RUNNING 回放,前端任务面板永远"处理中"
* 并阻塞该工具的后续任务(2026-09-13 真机:客户端重启后店铺数据采集一直显示 28094 处理中)。
*/
@Transactional
public void markTerminal(Long taskId, String moduleType, String status, String message) {
if (taskId == null || taskId <= 0 || isBlank(moduleType) || isBlank(status)) {
return;
}
TaskProgressSnapshotEntity existing = find(taskId, moduleType);
if (existing == null) {
return;
}
if (Objects.equals(existing.getStatus(), status) && Objects.equals(existing.getMessage(), message)) {
return;
}
taskProgressSnapshotMapper.update(null, new LambdaUpdateWrapper<TaskProgressSnapshotEntity>()
.eq(TaskProgressSnapshotEntity::getId, existing.getId())
.set(TaskProgressSnapshotEntity::getStatus, status)
.set(TaskProgressSnapshotEntity::getMessage, message)
.set(TaskProgressSnapshotEntity::getUpdatedAt, LocalDateTime.now()));
lastWriteAtMillis.remove(cacheKey(taskId, moduleType));
}
private String writeJson(Object value) { private String writeJson(Object value) {
try { try {
return objectMapper.writeValueAsString(value); return objectMapper.writeValueAsString(value);