|
|
@@ -7,12 +7,18 @@ import com.zsjz.ai.common.context.CaseContextHolder;
|
|
|
import com.zsjz.ai.common.exception.ServerException;
|
|
|
import com.zsjz.ai.common.model.dm.dto.DoCleanDTO;
|
|
|
import com.zsjz.ai.common.model.dm.entity.FileInfo;
|
|
|
+import com.zsjz.ai.common.model.plat.dto.SeeDataDTO;
|
|
|
+import com.zsjz.ai.common.model.plat.dto.SseDTO;
|
|
|
import com.zsjz.ai.module.aiclean.entity.AiCleanBatch;
|
|
|
import com.zsjz.ai.module.aiclean.entity.AiCleanJob;
|
|
|
import com.zsjz.ai.module.aiclean.enums.AiCleanStageEnum;
|
|
|
import com.zsjz.ai.module.aiclean.enums.AiCleanStatusEnum;
|
|
|
import com.zsjz.ai.module.aiclean.mapper.AiCleanBatchMapper;
|
|
|
+import com.zsjz.ai.module.aiclean.mapper.AiCleanJobMapper;
|
|
|
+import com.zsjz.ai.module.dm.mapper.FileInfoMapper;
|
|
|
import com.zsjz.ai.module.dm.service.DmService;
|
|
|
+import com.zsjz.ai.module.govern.serivce.GovernConfService;
|
|
|
+import com.zsjz.ai.module.plat.service.SseService;
|
|
|
import jakarta.annotation.PreDestroy;
|
|
|
import lombok.RequiredArgsConstructor;
|
|
|
import lombok.extern.slf4j.Slf4j;
|
|
|
@@ -36,12 +42,10 @@ import java.util.concurrent.atomic.AtomicLong;
|
|
|
/**
|
|
|
* 智能清洗融合编排:一次请求把「既有模板清洗」和「AI 清洗未匹配文件」跑完。
|
|
|
*
|
|
|
- * <h3>为什么判完再统一发一次 doClean</h3>
|
|
|
- * {@code DmService.doClean} 是 fire-and-forget,且完成回调里会 {@code sseService.closeSee()}。
|
|
|
- * 一个批次里若让 AI 命中既有模板的那些 sheet 各自再发一次 doClean,就会出现
|
|
|
- * ①同一文件被洗两遍 ②多个 SSE 会话先后关闭、进度页中途断流。所以本类的顺序是固定的:
|
|
|
- * <b>判定(并发) → 全部判完 → 合并成一次 doClean → AI 自建表入库</b>。
|
|
|
- * 判定未归零就进清洗阶段是错的,会把 AI 命中的表漏掉。
|
|
|
+ * <h3>本类负责什么</h3>
|
|
|
+ * <b>判定(并发) → 按降级顺序串行的静默 doClean(1级→2级) → 3级 AI 自建表入库 →
|
|
|
+ * SSE 收口推送</b>。判定未归零就进清洗阶段是错的,会把 AI 命中的表漏掉;
|
|
|
+ * 收口消息(统计/终态快照/治理宣告)统一在最后推送,保证顺序正确。
|
|
|
*
|
|
|
* <h3>与手动清洗不重叠</h3>
|
|
|
* 只接管 {@code templateId == null} 的 sheet;既有能匹配的绝不插手。
|
|
|
@@ -64,6 +68,15 @@ public class AiSmartCleanService {
|
|
|
|
|
|
private final DmService dmService;
|
|
|
|
|
|
+ /** 进度与终态全部走 SSE 推送(cmd 103),前端不再轮询状态接口 */
|
|
|
+ private final SseService sseService;
|
|
|
+
|
|
|
+ private final FileInfoMapper fileInfoMapper;
|
|
|
+
|
|
|
+ private final AiCleanJobMapper jobMapper;
|
|
|
+
|
|
|
+ private final GovernConfService governConfService;
|
|
|
+
|
|
|
@Value("${zsjz.ai-clean.max-concurrent:3}")
|
|
|
private int maxConcurrent;
|
|
|
|
|
|
@@ -130,7 +143,7 @@ public class AiSmartCleanService {
|
|
|
Thread.startVirtualThread(body);
|
|
|
}
|
|
|
|
|
|
- /** 前端轮询 */
|
|
|
+ /** 前端轮询(兼容旧入口;实时进度已改为 SSE 推送,见 {@link #pushProgress}) */
|
|
|
public SmartCleanStatusVO status(Long caseId, Long batchId) {
|
|
|
AiCleanBatch batch = batchMapper.selectOne(Wrappers.<AiCleanBatch>lambdaQuery()
|
|
|
.eq(AiCleanBatch::getCaseId, caseId)
|
|
|
@@ -144,6 +157,74 @@ public class AiSmartCleanService {
|
|
|
batch.getAiLoadCount(), batch.getNeedReviewCount(), batch.getFailedCount(), items);
|
|
|
}
|
|
|
|
|
|
+ /**
|
|
|
+ * 推送一次批次状态快照(cmd 103):阶段、计数、待人工明细随 SSE 直达前端,
|
|
|
+ * 前端据此刷新进度环与处理日志 —— <b>不再轮询状态接口</b>(需求 2026-10-11)。
|
|
|
+ *
|
|
|
+ * <p>推送是辅助链路:{@code status} 或发送失败都只记日志不上抛,
|
|
|
+ * 绝不能因为 SSE 断了把清洗主流程带崩。
|
|
|
+ */
|
|
|
+ private void pushProgress(Long caseId, Long batchId) {
|
|
|
+ try {
|
|
|
+ Long userId = CaseContextHolder.getUserId();
|
|
|
+ if (userId == null) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ SmartCleanStatusVO vo = status(caseId, batchId);
|
|
|
+ if (vo == null) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ sseService.sendSee(userId, new SseDTO(103, 0, "AI 清洗进度", new SeeDataDTO(0, 0, vo)));
|
|
|
+ } catch (Exception e) {
|
|
|
+ log.debug("AI 清洗进度推送失败: batchId={}, err={}", batchId, e.getMessage());
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 批次终态的 SSE 收口(口径与 {@code DmService.doClean} 完成时一致):
|
|
|
+ * 统计(cmd 101,100%)→ 状态快照(cmd 103)→ 治理启动宣告(cmd 102,开启治理时)。
|
|
|
+ *
|
|
|
+ * <p>顺序有意:治理关闭时前端收到终态快照会立即收尾并关闭流,统计消息必须先到。
|
|
|
+ * 智能清洗的 doClean 各趟都是静默的(见 {@code dispatchExistingByLevel}),
|
|
|
+ * 因此这一步是全链路唯一的终态出口 —— 纯 AI/无既有清洗的批次也靠它让前端完成跳转。
|
|
|
+ */
|
|
|
+ private void pushTerminal(Long caseId, Long batchId) {
|
|
|
+ try {
|
|
|
+ Long userId = CaseContextHolder.getUserId();
|
|
|
+ if (userId == null) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ // 成功条数 = 系统清洗入库行数(dataNum-failNum)+ AI 动态表入库行数
|
|
|
+ int dataNum = 0;
|
|
|
+ int failNum = 0;
|
|
|
+ for (FileInfo row : fileInfoMapper.selectList(Wrappers.lambdaQuery(FileInfo.class)
|
|
|
+ .eq(FileInfo::getBatchId, batchId)
|
|
|
+ .isNull(FileInfo::getFileType))) {
|
|
|
+ dataNum += java.util.Objects.requireNonNullElse(row.getDataNum(), 0);
|
|
|
+ failNum += java.util.Objects.requireNonNullElse(row.getFailNum(), 0);
|
|
|
+ }
|
|
|
+ long aiRows = 0;
|
|
|
+ for (AiCleanJob job : jobMapper.selectList(Wrappers.lambdaQuery(AiCleanJob.class)
|
|
|
+ .eq(AiCleanJob::getBatchId, batchId))) {
|
|
|
+ aiRows += java.util.Objects.requireNonNullElse(job.getRowLoaded(), 0);
|
|
|
+ }
|
|
|
+ int completed = (int) Math.max(0, dataNum - failNum + aiRows);
|
|
|
+ sseService.sendSee(userId, new SseDTO(101, 100, "[数据清洗]所有文件清洗完成!",
|
|
|
+ new SeeDataDTO(completed, failNum)));
|
|
|
+ SmartCleanStatusVO vo = status(caseId, batchId);
|
|
|
+ if (vo != null) {
|
|
|
+ sseService.sendSee(userId, new SseDTO(103, 100, "AI 清洗流程完成",
|
|
|
+ new SeeDataDTO(0, 0, vo)));
|
|
|
+ }
|
|
|
+ if (governConfService.hasGovern()) {
|
|
|
+ sseService.sendSee(userId, new SseDTO(102, 100, "[数据治理]开始数据治理!",
|
|
|
+ new SeeDataDTO(0, 0)));
|
|
|
+ }
|
|
|
+ } catch (Exception e) {
|
|
|
+ log.warn("智能清洗终态推送失败: batchId={}, err={}", batchId, e.getMessage());
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
// ==================== 后台主流程 ====================
|
|
|
|
|
|
/** 后台主流程:判定(并发)→ 一次 doClean → AI 自建表入库 → 终态 */
|
|
|
@@ -168,6 +249,8 @@ public class AiSmartCleanService {
|
|
|
}
|
|
|
}
|
|
|
}
|
|
|
+ // 判定开始即推一次快照:前端据此建立初始进度(此后全靠 SSE 推送,不再轮询)
|
|
|
+ pushProgress(caseId, batchId);
|
|
|
// 判定并发:闸门 + 虚拟线程,结果按批次聚合。
|
|
|
// ★ 判定任务跑在独立虚拟线程上,ScopedValue 绑定不会跨线程传播,
|
|
|
// 必须用 runWith 显式带上案件上下文,否则 judge 里的 slave 查询
|
|
|
@@ -187,6 +270,7 @@ public class AiSmartCleanService {
|
|
|
st.needReview.incrementAndGet();
|
|
|
bump(batchId, "needReview");
|
|
|
bump(batchId, "aiJudged");
|
|
|
+ pushProgress(caseId, batchId);
|
|
|
return;
|
|
|
}
|
|
|
// acquire 必须在 try 之外:acquire 本身被中断时 finally 会 release 一个没拿到的许可,闸门被放水
|
|
|
@@ -207,28 +291,36 @@ public class AiSmartCleanService {
|
|
|
gate.release();
|
|
|
// 三种结论都算「判完」,异常也算:否则前端进度会永远差几条卡在 99%
|
|
|
bump(batchId, "aiJudged");
|
|
|
+ // 每判完一张表推一次快照:进度环与无法清洗明细实时更新
|
|
|
+ pushProgress(caseId, batchId);
|
|
|
}
|
|
|
}), executor)).toList();
|
|
|
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
|
|
|
|
|
|
// 判定全部归零之后才能进清洗阶段,否则会漏掉 AI 命中的表
|
|
|
stage(batchId, AiCleanStageEnum.CLEANING_EXISTING);
|
|
|
+ pushProgress(caseId, batchId);
|
|
|
// ★ 按降级顺序串行清洗(需求 2026-10-11):1级(上传已匹配模板)先洗并等它
|
|
|
// 全部落库,再洗 2级(AI 命中既有模板),然后才进 3级(AI 自建表入库)。
|
|
|
// 4级(无法清洗)在判定阶段就只标记、不产生清洗任务。
|
|
|
int toExisting = dispatchExistingByLevel(roots, existingHits);
|
|
|
|
|
|
stage(batchId, AiCleanStageEnum.LOADING_AI);
|
|
|
+ pushProgress(caseId, batchId);
|
|
|
for (PendingAi p : pending) {
|
|
|
AiJudgeResult.AiPlan plan = plans.get(p.job().getId());
|
|
|
if (plan == null) {
|
|
|
continue;
|
|
|
}
|
|
|
loadPlan(batchId, p, plan, st);
|
|
|
+ pushProgress(caseId, batchId);
|
|
|
}
|
|
|
// 终态看本地计数,不回读 DB:并发写的行刚被别人累加过,读回来的值不代表这一轮
|
|
|
AiCleanStageEnum finalStage = st.failed.get() > 0 ? AiCleanStageEnum.PARTIAL : AiCleanStageEnum.DONE;
|
|
|
stage(batchId, finalStage);
|
|
|
+ // 终态收口(全 SSE):统计 → 状态快照 → 治理宣告;纯 AI/无既有清洗的批次
|
|
|
+ // 没有 doClean 的完成回调,前端的「完成并跳转统计页」完全依赖这一推
|
|
|
+ pushTerminal(caseId, batchId);
|
|
|
log.info("智能清洗批次完成: batchId={}, 阶段={}, 判定={} 条, 交既有清洗={} 条, AI 入库={} 条, "
|
|
|
+ "待人工={} 条, 失败={} 条, 耗时={}ms",
|
|
|
batchId, finalStage, pending.size(), toExisting, st.aiLoad.get(),
|
|
|
@@ -291,8 +383,9 @@ public class AiSmartCleanService {
|
|
|
* 2级(AI 命中既有模板的表)。两级各自的根树只放本级的 sheet ——
|
|
|
* NEED_REVIEW / Skip 的表不进任何一趟,doClean 不再为它们记「未配置模板」噪音日志。
|
|
|
*
|
|
|
- * <p>首趟静默(不推终态、不关 SSE 会话、不触发同时治理),由第二趟收口;趟间用
|
|
|
- * {@link DmService#doClean(DoCleanDTO, boolean)} 返回的 Latch 严格排队。
|
|
|
+ * <p>趟间用 {@link DmService#doClean(DoCleanDTO, boolean)} 返回的 Latch 严格排队。
|
|
|
+ * 两趟都静默:收口消息(统计/终态快照/治理宣告)由本类在 3级入库完成后统一推送
|
|
|
+ * (见 {@link #pushTerminal}),保证「治理启动」出现在全部清洗动作之后。
|
|
|
*
|
|
|
* @return 并入既有链路的 AI 命中表数(仅日志用)
|
|
|
*/
|
|
|
@@ -334,11 +427,10 @@ public class AiSmartCleanService {
|
|
|
log.info("智能清洗:按降级顺序串行清洗 1级表={} 个 → 2级表={} 个",
|
|
|
l1ByRoot.values().stream().mapToInt(List::size).sum(),
|
|
|
l2ByRoot.values().stream().mapToInt(List::size).sum());
|
|
|
- // 1级:若后面还有 2级那趟则静默(终态与关会话留给最后一趟收口),
|
|
|
- // 只有 1级时就由它自己收口;2级:等 1级 Latch 归零后才开始
|
|
|
- boolean l2HasWork = l2ByRoot.values().stream().anyMatch(l -> !l.isEmpty());
|
|
|
- runCleanPass(roots, l1ByRoot, l2HasWork);
|
|
|
- runCleanPass(roots, l2ByRoot, false);
|
|
|
+ // 两趟都静默:收口消息由本类在 3级入库完成后统一推送(run → pushTerminal),
|
|
|
+ // 保证「所有文件清洗完成 / 开始数据治理」出现在全部清洗动作之后
|
|
|
+ runCleanPass(roots, l1ByRoot, true);
|
|
|
+ runCleanPass(roots, l2ByRoot, true);
|
|
|
return applied;
|
|
|
}
|
|
|
|