fix(task): 新增陈旧任务自动修复兜底——历史任务永久停留中间状态的排查与兜底

生产库排查发现三类脏状态:陈旧 PENDING/SCHEDULED 无调度器再捞(243条)、
心跳残留的万年 RUNNING、SUCCESS 但结果报错(上报契约问题)。
已手工修复 269 条;本提交补常驻兜底:每 10 分钟检查 PENDING/SCHEDULED/RUNNING
心跳(updated_at)超 120 分钟即标 FAILED,分布式锁防双实例重复。
只用心跳线、不限制总时长——超长真活任务不会被误杀。
This commit is contained in:
2026-09-11 16:03:01 +08:00
parent 32985c78f3
commit bb74c9a192
@@ -0,0 +1,128 @@
package com.nanri.aiimage.modules.task.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.nanri.aiimage.common.service.DistributedJobLockService;
import com.nanri.aiimage.modules.brand.mapper.BrandCrawlTaskMapper;
import com.nanri.aiimage.modules.brand.model.entity.BrandCrawlTaskEntity;
import com.nanri.aiimage.modules.task.mapper.FileTaskMapper;
import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.util.List;
/**
* 陈旧任务自动修复。
*
* 修复型巡检(区别于 task-161 的只读报告巡检):把超过阈值时间无心跳(updated_at 未刷新)的
* PENDING / SCHEDULED / RUNNING 任务自动标终态失败,避免历史任务列表永久停留在中间状态。
* 分布式锁防双实例重复修复;注意 RUNNING 的存活判定基于 updatedAt 心跳,真活任务会不断刷新不会被误杀。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class StaleTaskRepairService {
private static final String STATUS_PENDING = "PENDING";
private static final String STATUS_SCHEDULED = "SCHEDULED";
private static final String STATUS_RUNNING = "RUNNING";
private static final String STATUS_FAILED = "FAILED";
private static final String STATUS_BRAND_PENDING = "pending";
private static final String STATUS_BRAND_RUNNING = "running";
private static final String STATUS_BRAND_FAILED = "failed";
private static final String STATUS_BRAND_CANCELLED = "cancelled";
/** 中间态(无人接单)或 RUNNING(心跳残留)超过该时长未更新即判死。真活任务心跳会不断刷新 updatedAt,不会被误杀。 */
private static final long STALE_IDLE_MINUTES = 120;
private final FileTaskMapper fileTaskMapper;
private final BrandCrawlTaskMapper brandCrawlTaskMapper;
private final DistributedJobLockService distributedJobLockService;
/**
* 每 10 分钟执行一次;双实例下用分布式锁保证只有一台执行修复。
* 注意:这里直接用 @Scheduled 而不走 InspectionScheduler(那是默认关闭的报表巡检),兜底修复必须常开。
*/
@Scheduled(fixedDelay = 10 * 60 * 1000, initialDelay = 5 * 60 * 1000)
public void repairStaleTasks() {
var lock = distributedJobLockService.tryLock("stale-task-repair", java.time.Duration.ofMinutes(10));
if (lock == null) {
log.info("[stale-task-repair] 另一实例持有修复锁,跳过");
return;
}
try (lock) {
repairFileTaskStaleIdle();
repairFileTaskStaleRunning();
repairBrandStale();
} catch (Exception ex) {
log.warn("[stale-task-repair] 修护失败 msg={}", ex.getMessage(), ex);
}
}
/** 中间态 PENDING/SCHEDULED 超时未接单 → FAILED。 */
private void repairFileTaskStaleIdle() {
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(STALE_IDLE_MINUTES);
List<FileTaskEntity> stale = fileTaskMapper.selectList(new LambdaQueryWrapper<FileTaskEntity>()
.select(FileTaskEntity::getId, FileTaskEntity::getModuleType)
.in(FileTaskEntity::getStatus, STATUS_PENDING, STATUS_SCHEDULED)
.lt(FileTaskEntity::getUpdatedAt, cutoff)
.last("limit 500"));
if (stale.isEmpty()) {
return;
}
String reason = "任务长期未被领取,已自动失败";
LocalDateTime now = LocalDateTime.now();
int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper<FileTaskEntity>()
.in(FileTaskEntity::getId, stale.stream().map(FileTaskEntity::getId).toList())
.in(FileTaskEntity::getStatus, STATUS_PENDING, STATUS_SCHEDULED)
.set(FileTaskEntity::getStatus, STATUS_FAILED)
.set(FileTaskEntity::getErrorMessage, reason)
.set(FileTaskEntity::getFinishedAt, now));
if (updated > 0) {
log.warn("[stale-task-repair] 陈旧中间态已标失败 count={} ids={}",
updated, stale.stream().map(FileTaskEntity::getId).toList());
}
}
/** RUNNING 但心跳超过 2h 未刷新 → FAILED(心跳残留续命即僵尸)。 */
private void repairFileTaskStaleRunning() {
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(STALE_IDLE_MINUTES);
int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getStatus, STATUS_RUNNING)
.lt(FileTaskEntity::getUpdatedAt, cutoff)
.set(FileTaskEntity::getStatus, STATUS_FAILED)
.set(FileTaskEntity::getErrorMessage, "任务心跳超时,已自动失败")
.set(FileTaskEntity::getFinishedAt, LocalDateTime.now()));
if (updated > 0) {
log.warn("[stale-task-repair] 心跳超时 RUNNING 已标失败 count={}", updated);
}
}
/** brand_crawl_tasks 的 pending/running 陈旧 → failed。 */
private void repairBrandStale() {
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(STALE_IDLE_MINUTES);
List<BrandCrawlTaskEntity> stale = brandCrawlTaskMapper.selectList(new LambdaQueryWrapper<BrandCrawlTaskEntity>()
.select(BrandCrawlTaskEntity::getId)
.in(BrandCrawlTaskEntity::getStatus, STATUS_BRAND_PENDING, STATUS_BRAND_RUNNING)
.lt(BrandCrawlTaskEntity::getUpdatedAt, cutoff)
.last("limit 500"));
if (stale.isEmpty()) {
return;
}
int updated = brandCrawlTaskMapper.update(null, new LambdaUpdateWrapper<BrandCrawlTaskEntity>()
.in(BrandCrawlTaskEntity::getId, stale.stream().map(BrandCrawlTaskEntity::getId).toList())
.in(BrandCrawlTaskEntity::getStatus, STATUS_BRAND_PENDING, STATUS_BRAND_RUNNING)
.lt(BrandCrawlTaskEntity::getUpdatedAt, cutoff)
.set(BrandCrawlTaskEntity::getStatus, STATUS_BRAND_FAILED)
.set(BrandCrawlTaskEntity::getErrorMessage, "任务长期无心跳,已自动失败"));
if (updated > 0) {
log.warn("[stale-task-repair] brand 陈旧任务已标失败 count={} ids={}",
updated, stale.stream().map(BrandCrawlTaskEntity::getId).toList());
}
}
}