fix(任务中断): 已终态任务也幂等刷新缓存/快照,修历史遗留的处理中卡死
第一次修复只在 DB 刚被改成 FAILED 时同步缓存;对 DB 已是终态、但 Redis 里 仍缓存 RUNNING 的历史任务不生效(真机:28094 反复调中断仍返回 progress=RUNNING)。 现在无论是否刚更新 DB,都按 DB 现值同步模块缓存与进度快照,可自愈。
This commit is contained in:
+24
-12
@@ -143,21 +143,19 @@ 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 读缓存/快照,不刷新会让前端
|
// 同步模块缓存与进度快照:progress/batch 读缓存/快照(Redis,跨实例共享、
|
||||||
// 一直按 RUNNING 渲染(任务面板永远"处理中"并阻塞该工具后续任务)。
|
// 不随 JVM 重启清空),不刷新会让前端一直按 RUNNING 渲染
|
||||||
try {
|
// (任务面板永远"处理中"并阻塞该工具后续任务)。
|
||||||
fileTask.setStatus("FAILED");
|
fileTask.setStatus("FAILED");
|
||||||
fileTask.setErrorMessage(safeReason);
|
fileTask.setErrorMessage(safeReason);
|
||||||
fileTask.setFinishedAt(LocalDateTime.now());
|
fileTask.setFinishedAt(LocalDateTime.now());
|
||||||
saveFileTaskCache(fileTask.getModuleType(), fileTask);
|
syncTerminalState(fileTask, "FAILED", safeReason);
|
||||||
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);
|
||||||
|
// 已终态但缓存/快照仍残留 RUNNING 时(历史遗留或上一次中断未刷缓存)自愈:
|
||||||
|
// 幂等按 DB 现值刷新,避免前端被旧缓存永久卡住。
|
||||||
|
syncTerminalState(fileTask, status, fileTask.getErrorMessage());
|
||||||
return TaskHeartbeatVo.notAlive(fileTask.getModuleType(), status, "task is not running");
|
return TaskHeartbeatVo.notAlive(fileTask.getModuleType(), status, "task is not running");
|
||||||
}
|
}
|
||||||
BrandCrawlTaskEntity brandTask = selectBrandTask(taskId);
|
BrandCrawlTaskEntity brandTask = selectBrandTask(taskId);
|
||||||
@@ -355,6 +353,20 @@ public class TaskHeartbeatService {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** 把终态同步到模块缓存与进度快照;失败只告警,不影响中断接口本身的成功语义。 */
|
||||||
|
private void syncTerminalState(FileTaskEntity task, String status, String message) {
|
||||||
|
if (task == null || task.getId() == null) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
saveFileTaskCache(task.getModuleType(), task);
|
||||||
|
taskProgressSnapshotService.markTerminal(task.getId(), task.getModuleType(), status, message);
|
||||||
|
} catch (Exception ex) {
|
||||||
|
log.warn("[task-interrupted] refresh cache/snapshot failed taskId={} err={}",
|
||||||
|
task.getId(), ex.getMessage());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private void putIfPresent(Map<String, String> values, String key, Object value) {
|
private void putIfPresent(Map<String, String> values, String key, Object value) {
|
||||||
if (value != null) {
|
if (value != null) {
|
||||||
values.put(key, String.valueOf(value));
|
values.put(key, String.valueOf(value));
|
||||||
|
|||||||
Reference in New Issue
Block a user