Преглед изворни кода

fix: 智能清洗链路加固(评审 P0/P1/P2 全项落地)

- P0-1 出口B流式写库:AiRowReader.forEachRow(带真实行号) + AiDuckDb.BatchWriter
  (2000行攒批),loadViaAiTemplate 重写为建表→流式边洗边写→成功才登记 ONLINE,
  失败删空表不登记 —— 消除全量内存物化与幽灵模板写序矛盾
- P0-2 AiCleanWatchdog:启动恢复(非终态批次标FAILED) + 运行时心跳超时看门狗
  (Celery visibility timeout 同款),恢复路径=重置计数开新一轮
- P0-3 upsertBatch 撞唯一键回落查询(insertQuiet 同款幂等,等价 ON CONFLICT),
  修复并发首点 500 竞态;新批次补 update_time
- P1-4 LlmRetry(LiteLLM 同款:429/5xx/网络层重试,指数退避+抖动) 套 askLlm 与
  文件识别;批次 token 预算护栏(zsjz.ai-clean.batch-token-budget,超限转人工)
- P1-5 遥测落库:LlmService.chatAsDetail 带回用量 + AiMatchOutcome.llm 组件,
  judge 写 model_name/token_in/token_out/draft/detail/header_row(此前全部未赋值)
- P1-6 过程状态机接通:mark(PROBING/MATCHING/VERIFYING/LOADING) 随判定推进落库
- P2-7 NEED_REVIEW 改块级判据:SheetEvidence.primaryMergeCount(主块内合并区),
  远处无关合并格不再否决整表
- 测试同步:MatcherTest 桩改 chatAsDetail、LiveTest 适配 LlmAnswer、yaml 增
  zsjz.ai-clean 配置段
cc пре 1 недеља
родитељ
комит
36e167e324

+ 87 - 0
ai-server/src/main/java/com/zsjz/ai/module/agent/llm/LlmRetry.java

@@ -0,0 +1,87 @@
+package com.zsjz.ai.module.agent.llm;
+
+import com.zsjz.ai.common.exception.ServerException;
+import lombok.extern.slf4j.Slf4j;
+
+import java.util.concurrent.ThreadLocalRandom;
+import java.util.function.Supplier;
+
+/**
+ * LLM 调用重试:LiteLLM/Portkey 网关同款策略。
+ *
+ * <p>只重试「值得重试」的失败:429(限流)、5xx(厂商抖动/网关超时)、网络层异常。
+ * 400/404 是参数与配置错误,重试只会同样失败 —— 直接抛。
+ * 退避曲线用 AWS 架构博客的推荐形态:指数退避(1s → 2s → 4s)+ 抖动防雪崩。
+ *
+ * <p><b>为什么是独立工具而不是塞进 {@link LlmService}</b>:那里的契约是「不做重试」
+ * (类注释明说)。重试是<b>调用方策略</b> —— 批处理链路(智能清洗、文件识别)值得为
+ * 成功率多等两轮,在线对话链路(意图/追问)失败即降级,加重试只会拖慢失败路径。
+ *
+ * <p>注意:温度 0 的判定类调用重试大概率复现同样的失败(解析失败重试收益有限),
+ * 但厂商偶发的限流/截断/抽风占比不低,三次封顶的代价有界。
+ */
+@Slf4j
+public final class LlmRetry {
+
+    /** 首次 + 2 次重试,对齐 LiteLLM num_retries 默认档 */
+    private static final int MAX_ATTEMPTS = 3;
+
+    private static final long BASE_DELAY_MS = 1000;
+
+    private static final long MAX_DELAY_MS = 30_000;
+
+    private LlmRetry() {
+    }
+
+    /**
+     * 执行一次可重试的 LLM 调用(默认 3 次)。
+     *
+     * @param what 动作名,仅用于日志定位
+     * @param call 实际调用(闭包持有全部入参)
+     * @return 首次成功的返回值;全部尝试失败时抛最后一次的异常
+     */
+    public static <T> T call(String what, Supplier<T> call) {
+        return call(what, call, MAX_ATTEMPTS);
+    }
+
+    /**
+     * 执行一次可重试的 LLM 调用。
+     */
+    public static <T> T call(String what, Supplier<T> call, int attempts) {
+        RuntimeException last = null;
+        for (int attempt = 1; attempt <= attempts; attempt++) {
+            try {
+                return call.get();
+            } catch (ServerException e) {
+                if (!retryable(e.getCode())) {
+                    throw e;
+                }
+                last = e;
+            } catch (RuntimeException e) {
+                // 网络层异常(连接重置、流中断等)值得再试
+                last = e;
+            }
+            if (attempt < attempts) {
+                sleep(what, attempt, last);
+            }
+        }
+        throw last;
+    }
+
+    private static boolean retryable(int code) {
+        return code == 429 || (code >= 500 && code <= 599);
+    }
+
+    private static void sleep(String what, int attempt, RuntimeException cause) {
+        long backoff = Math.min(MAX_DELAY_MS, BASE_DELAY_MS << (attempt - 1));
+        long jitter = (long) (backoff * 0.2 * (ThreadLocalRandom.current().nextDouble() * 2 - 1));
+        long delay = Math.max(250, backoff + jitter);
+        log.warn("{} 第 {} 次调用失败,{}ms 后重试: {}", what, attempt, delay, cause.getMessage());
+        try {
+            Thread.sleep(delay);
+        } catch (InterruptedException e) {
+            Thread.currentThread().interrupt();
+            throw new IllegalStateException(what + " 重试等待被中断", e);
+        }
+    }
+}

+ 24 - 2
ai-server/src/main/java/com/zsjz/ai/module/agent/llm/LlmService.java

@@ -204,12 +204,34 @@ public class LlmService {
      * @throws ServerException 模型未返回合法 JSON 或字段不匹配
      */
     public <T> T chatAs(LlmRequest request, Class<T> type) {
+        return chatAsDetail(request, type).value();
+    }
+
+    /**
+     * {@link #chatAs} 的带用量版本:结构化输出之外,把本次调用的模型名 / token / 耗时一并带出。
+     *
+     * <p>批处理链路(智能清洗、文件识别)用它把成本记进各自的埋点列 ——
+     * 没有用量遥测,「降本有没有成立」就只能靠猜。
+     *
+     * @param request 请求参数;{@code system} 会被追加 schema 说明
+     * @param type    目标类型(POJO,字段用 public 或带 getter/setter)
+     * @param <T>     目标类型
+     * @return 解析结果 + 原始调用用量
+     * @throws ServerException 模型未返回合法 JSON 或字段不匹配
+     */
+    public <T> StructuredResult<T> chatAsDetail(LlmRequest request, Class<T> type) {
         if (type == null) {
             throw new ServerException(400, "type 不能为空");
         }
         String schema = PromptHelper.generateSchema(type);
-        String raw = chat(request.toBuilder().system(appendSchema(request.getSystem(), schema)).build());
-        return parseStructured(raw, type);
+        LlmResult raw = chatDetail(request.toBuilder().system(appendSchema(request.getSystem(), schema)).build());
+        return new StructuredResult<>(parseStructured(raw.text(), type), raw);
+    }
+
+    /**
+     * 结构化调用的产出:解析后的对象 + 原始用量。
+     */
+    public record StructuredResult<T>(T value, LlmResult raw) {
     }
 
     /**

+ 143 - 20
ai-server/src/main/java/com/zsjz/ai/module/aiclean/exec/AiDuckDb.java

@@ -113,45 +113,156 @@ public final class AiDuckDb {
     }
 
     /**
-     * 批量写入。
+     * 批量写入(收集齐一批后一次性落库)。
      *
      * @param columns 业务列(顺序与 {@code rows} 里每行的取值顺序一致)
-     * @return 实际写入行数
+     * @return 实际写入行数;失败返回 -1
      */
     public static long insertBatch(String table, List<String> columns,
                                    List<Map<String, Object>> rows) {
         if (rows.isEmpty()) {
             return 0L;
         }
-        List<String> cols = new ArrayList<>(List.of("id", "case_id", "file_id", "sheet_id",
-                "block_no", "row_no", "ai_job_id", "ai_rule_ver", "category"));
-        cols.addAll(columns);
-        String sql = "INSERT INTO " + table + " (" + String.join(",", cols) + ") VALUES ("
-                + String.join(",", cols.stream().map(x -> "?").toList()) + ")";
-        long written = 0;
-        try (Connection c = caseDb().getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
-            int pending = 0;
+        try (BatchWriter writer = openWriter(table, columns)) {
             for (Map<String, Object> row : rows) {
+                writer.add(row);
+            }
+            return writer.finish();
+        } catch (IllegalStateException e) {
+            // 与旧实现同契约:返回 -1 而不是「已写行数」—— 吞掉异常会让调用方把
+            // 「一行都没写进去」当成清洗成功,而全自动链路上没有人工关口去发现这件事
+            log.warn("AI 批量写入失败: table={}, err={}", table, e.getMessage());
+            return -1L;
+        }
+    }
+
+    /**
+     * 打开一个流式批量写入器:Fesod 回调里逐行 {@link BatchWriter#add},攒满
+     * {@value #BATCH} 行自动 flush,堆占用恒定 —— 大表全量入库唯一安全的姿势
+     * (EasyExcel 官方建议与 Spark/Flink 管道同款铁律:读数回调里绝不全量 collect)。
+     *
+     * <p>一个 writer 独占一条连接(DuckDB 单写者模型),{@link BatchWriter#finish()} /
+     * {@link BatchWriter#close()} 归还。行里的固定列(id/case_id/…/category)由调用方
+     * 在回调里填好,键名与 {@link BatchWriter} 的列装配一致。
+     *
+     * @throws IllegalStateException 连接/预编译失败(调用方据此放弃本轮入库)
+     */
+    public static BatchWriter openWriter(String table, List<String> columns) {
+        try {
+            return new BatchWriter(table, columns);
+        } catch (Exception e) {
+            throw new IllegalStateException("AI 批量写入器创建失败: table=" + table + ", err=" + e.getMessage(), e);
+        }
+    }
+
+    private static String insertSql(String table, List<String> cols) {
+        return "INSERT INTO " + table + " (" + String.join(",", cols) + ") VALUES ("
+                + String.join(",", cols.stream().map(x -> "?").toList()) + ")";
+    }
+
+    /**
+     * 流式批量写入器(内部类):持有一条连接 + 一个 {@link PreparedStatement}。
+     *
+     * <p><b>失败语义</b>:任何 SQL 失败都会立刻中断写入并让 {@code add} 抛出 ——
+     * 半截数据绝不能被当成成功。已 flush 的行会留在表里,但每行都带着 {@code ai_job_id},
+     * job 标 FAIL 后可用 {@link #deleteByJob} 精确清理 —— 这正是撤销接口存在的意义。
+     */
+    public static final class BatchWriter implements AutoCloseable {
+
+        private final String table;
+        private final List<String> cols;
+        private final Connection connection;
+        private final PreparedStatement ps;
+
+        private int pending;
+        private long written;
+        private boolean failed;
+        private boolean closed;
+
+        private BatchWriter(String table, List<String> columns) throws Exception {
+            this.table = table;
+            this.cols = new ArrayList<>(List.of("id", "case_id", "file_id", "sheet_id",
+                    "block_no", "row_no", "ai_job_id", "ai_rule_ver", "category"));
+            this.cols.addAll(columns);
+            this.connection = caseDb().getConnection();
+            this.ps = this.connection.prepareStatement(insertSql(table, this.cols));
+        }
+
+        /**
+         * 追加一行;攒满 {@value #BATCH} 行自动 flush。
+         *
+         * @throws IllegalStateException SQL 失败(写入链路应立刻中断本轮入库)
+         */
+        public void add(Map<String, Object> row) {
+            if (failed || closed) {
+                throw new IllegalStateException("AI 批量写入器已关闭或已失败,不能再写: " + table);
+            }
+            try {
                 int idx = 1;
                 for (String col : cols) {
                     ps.setObject(idx++, row.get(col));
                 }
                 ps.addBatch();
                 if (++pending >= BATCH) {
-                    written += java.util.Arrays.stream(ps.executeBatch()).filter(x -> x >= 0).count();
-                    pending = 0;
+                    flush();
                 }
+            } catch (Exception e) {
+                failed = true;
+                closeQuietly();
+                throw new IllegalStateException("AI 批量写入失败: table=" + table + ", err=" + e.getMessage(), e);
             }
-            if (pending > 0) {
-                written += java.util.Arrays.stream(ps.executeBatch()).filter(x -> x >= 0).count();
+        }
+
+        /**
+         * 冲刷剩余行并关闭。
+         *
+         * @return 累计写入行数;过程中失败过返回 -1
+         */
+        public long finish() {
+            if (closed) {
+                return failed ? -1L : written;
+            }
+            try {
+                if (!failed && pending > 0) {
+                    flush();
+                }
+                return failed ? -1L : written;
+            } catch (Exception e) {
+                failed = true;
+                log.warn("AI 批量写入收尾失败: table={}, err={}", table, e.getMessage());
+                return -1L;
+            } finally {
+                closeQuietly();
+            }
+        }
+
+        private void flush() throws Exception {
+            written += java.util.Arrays.stream(ps.executeBatch()).filter(x -> x >= 0).count();
+            pending = 0;
+        }
+
+        @Override
+        public void close() {
+            closeQuietly();
+        }
+
+        private void closeQuietly() {
+            closed = true;
+            try {
+                if (ps != null) {
+                    ps.close();
+                }
+            } catch (Exception ignore) {
+                // 归还连接优先,关闭失败不再放大
+            }
+            try {
+                if (connection != null) {
+                    connection.close();
+                }
+            } catch (Exception ignore) {
+                // 同上
             }
-        } catch (Exception e) {
-            // 返回 -1 而不是「已写行数」:吞掉异常会让调用方把「一行都没写进去」
-            // 当成清洗成功,而全自动链路上没有人工关口去发现这件事
-            log.warn("AI 批量写入失败: table={}, err={}", table, e.getMessage());
-            return -1L;
         }
-        return written;
     }
 
     /** 按 job 精确删除:AI 链路的撤销单位,不借用既有按文件维度的物理删除 */
@@ -166,6 +277,18 @@ public final class AiDuckDb {
         }
     }
 
+    /**
+     * 删除 AI 动态表:入库失败后的补偿清理(流式写数中断/零行时调用),
+     * 不在案件库里留孤儿空表。表名由服务端生成(ai_t + 雪花 id),无注入面。
+     */
+    public static void dropTable(String table) {
+        try (Connection c = caseDb().getConnection(); Statement st = c.createStatement()) {
+            st.executeUpdate("DROP TABLE IF EXISTS " + table);
+        } catch (Exception e) {
+            log.warn("AI 动态表删除失败(可留待人工清理): table={}, err={}", table, e.getMessage());
+        }
+    }
+
     /** 表不存在或出错时返回 -1,调用方据此跳过治理树节点而不是让整个任务失败 */
     public static long count(String table) {
         try (Connection c = caseDb().getConnection(); Statement st = c.createStatement();

+ 30 - 17
ai-server/src/main/java/com/zsjz/ai/module/aiclean/match/AiTemplateMatcher.java

@@ -11,7 +11,9 @@ import com.zsjz.ai.module.aiclean.model.RuleBind;
 import com.zsjz.ai.module.aiclean.model.SheetEvidence;
 import com.zsjz.ai.module.aiclean.model.TemplateCandidate;
 import com.zsjz.ai.module.agent.llm.LlmRequest;
+import com.zsjz.ai.module.agent.llm.LlmResult;
 import com.zsjz.ai.module.agent.llm.LlmService;
+import com.zsjz.ai.module.agent.llm.LlmRetry;
 import lombok.RequiredArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.stereotype.Component;
@@ -87,7 +89,7 @@ public class AiTemplateMatcher {
         for (AiTemplateWithFields awf : aiOnline) {
             if (md5 != null && md5.equals(awf.template().getHeaderMd5())) {
                 return outcome(AiMatchTierEnum.STRONG, AiCleanExitEnum.B,
-                        toCandidate(awf, 0, 1.0), List.of(), AiMatchOutcome.BY_CACHE, md5, false, null);
+                        toCandidate(awf, 0, 1.0), List.of(), AiMatchOutcome.BY_CACHE, md5, false, null, null);
             }
         }
 
@@ -105,7 +107,7 @@ public class AiTemplateMatcher {
         if (best != null && best.score() >= HeaderScorer.STRONG) {
             // 这一腿是确定性名字匹配,不掺模型猜测,所以允许直接出口 A
             return outcome(AiMatchTierEnum.STRONG, exitOf(AiMatchTierEnum.STRONG, best),
-                    best, scored, AiMatchOutcome.BY_SCORE, md5, false, null);
+                    best, scored, AiMatchOutcome.BY_SCORE, md5, false, null, null);
         }
 
         // 按名次重编号:提示词里的 #n 必须正好对应下面列表的下标 n-1。
@@ -120,7 +122,8 @@ public class AiTemplateMatcher {
         if (!llmEnabled || shortList.isEmpty()) {
             return AiMatchOutcome.none(md5, ranked, false);
         }
-        AiRuleDraft draft = askLlm(ev, headers, shortList);
+        LlmAnswer answer = askLlm(ev, headers, shortList);
+        AiRuleDraft draft = answer == null ? null : answer.draft();
         if (draft == null || draft.getTemplateNo() == null || draft.getTemplateNo() <= 0) {
             return AiMatchOutcome.none(md5, ranked, true);
         }
@@ -137,25 +140,35 @@ public class AiTemplateMatcher {
         AiMatchTierEnum tier = conf >= LLM_MIN_CONFIDENCE ? AiMatchTierEnum.STRONG : AiMatchTierEnum.WEAK;
         log.info("AI 匹配第三腿:选中 #{} {},自评置信 {},判档 {}",
                 draft.getTemplateNo(), picked.nameCn(), conf, tier);
-        return outcome(tier, exitOf(tier, picked), picked, ranked, AiMatchOutcome.BY_LLM, md5, true, draft);
+        return outcome(tier, exitOf(tier, picked), picked, ranked, AiMatchOutcome.BY_LLM, md5, true,
+                draft, answer.usage());
     }
 
     /**
-     * 调模型出草稿;失败或不可解析返回 null(调用方按未命中处理,不抛异常)。
+     * 单次模型调用的产出:规则草稿 + 用量(遥测落库用)。
      */
-    public AiRuleDraft askLlm(SheetEvidence ev, List<String> headers, List<TemplateCandidate> ranked) {
+    public record LlmAnswer(AiRuleDraft draft, LlmResult usage) {
+    }
+
+    /**
+     * 调模型出草稿(含 {@link LlmRetry} 指数退避重试);彻底失败返回 null(调用方按未命中处理,不抛异常)。
+     */
+    public LlmAnswer askLlm(SheetEvidence ev, List<String> headers, List<TemplateCandidate> ranked) {
         try {
-            LlmRequest request = LlmRequest.builder()
-                    .system(MatchPromptBuilder.system())
-                    .user(MatchPromptBuilder.user(ev, headers, ranked))
-                    // 绑定是判定不是创作,温度压到 0 才可能重放一致
-                    .temperature(0.0)
-                    .timeout(Duration.ofSeconds(90))
-                    .build();
-            return llmService.chatAs(request, AiRuleDraft.class);
+            return LlmRetry.call("AI 模板匹配", () -> {
+                LlmRequest request = LlmRequest.builder()
+                        .system(MatchPromptBuilder.system())
+                        .user(MatchPromptBuilder.user(ev, headers, ranked))
+                        // 绑定是判定不是创作,温度压到 0 才可能重放一致
+                        .temperature(0.0)
+                        .timeout(Duration.ofSeconds(90))
+                        .build();
+                LlmService.StructuredResult<AiRuleDraft> out = llmService.chatAsDetail(request, AiRuleDraft.class);
+                return new LlmAnswer(out.value(), out.raw());
+            });
         } catch (Exception e) {
             // 记全异常类型:只留 message 时,超时和「模型返回的不是合法 JSON」看起来一模一样
-            log.warn("AI 匹配模型调用失败,按未命中处理: {} {}", e.getClass().getSimpleName(), e.getMessage());
+            log.warn("AI 匹配模型调用失败(含重试),按未命中处理: {} {}", e.getClass().getSimpleName(), e.getMessage());
             return null;
         }
     }
@@ -198,8 +211,8 @@ public class AiTemplateMatcher {
 
     private static AiMatchOutcome outcome(AiMatchTierEnum tier, AiCleanExitEnum exit, TemplateCandidate cand,
                                           List<TemplateCandidate> ranked, String by, String md5, boolean llm,
-                                          AiRuleDraft draft) {
-        return new AiMatchOutcome(tier, exit, cand, ranked, by, md5, llm, draft);
+                                          AiRuleDraft draft, LlmResult usage) {
+        return new AiMatchOutcome(tier, exit, cand, ranked, by, md5, llm, draft, usage);
     }
 
     private static TemplateCandidate toCandidate(AiTemplateWithFields awf, int no, double score) {

+ 5 - 2
ai-server/src/main/java/com/zsjz/ai/module/aiclean/model/AiMatchOutcome.java

@@ -2,6 +2,7 @@ package com.zsjz.ai.module.aiclean.model;
 
 import com.zsjz.ai.module.aiclean.enums.AiCleanExitEnum;
 import com.zsjz.ai.module.aiclean.enums.AiMatchTierEnum;
+import com.zsjz.ai.module.agent.llm.LlmResult;
 
 import java.util.List;
 
@@ -17,6 +18,7 @@ import java.util.List;
  * @param llmCalled  这一步是否真调了模型(直通时为 false,是降本看板的分子)
  * @param draft      模型已经给出的规则草稿。走第三腿时顺带产出,
  *                   <b>必须复用</b> —— 否则同一个 job 要打两次模型(实测单次就要几十秒,翻倍)
+ * @param llm        第三腿那次模型调用的用量(模型名/token),遥测落库用;前两腿为 null
  * @author cc
  * @since 2026/9/21
  */
@@ -27,7 +29,8 @@ public record AiMatchOutcome(AiMatchTierEnum tier,
                              String matchedBy,
                              String headerMd5,
                              boolean llmCalled,
-                             com.zsjz.ai.module.aiclean.model.AiRuleDraft draft) {
+                             com.zsjz.ai.module.aiclean.model.AiRuleDraft draft,
+                             LlmResult llm) {
 
     /** 命中途径:表头指纹缓存,零模型调用 */
     public static final String BY_CACHE = "HEADER_CACHE";
@@ -43,6 +46,6 @@ public record AiMatchOutcome(AiMatchTierEnum tier,
      */
     public static AiMatchOutcome none(String headerMd5, List<TemplateCandidate> candidates, boolean llmCalled) {
         return new AiMatchOutcome(AiMatchTierEnum.NONE, AiCleanExitEnum.B, null,
-                candidates, BY_NONE, headerMd5, llmCalled, null);
+                candidates, BY_NONE, headerMd5, llmCalled, null, null);
     }
 }

+ 4 - 1
ai-server/src/main/java/com/zsjz/ai/module/aiclean/model/SheetEvidence.java

@@ -21,7 +21,9 @@ import java.util.List;
  * @param blocks      探测出的表块(单表头时只有一个)
  * @param headWindow  前 N 行原文,供模型判断表头位置与说明行
  * @param colStats    选定表块后的逐列统计
- * @param mergeCount  整个 sheet 的合并区数量(脏结构的一路信号)
+ * @param mergeCount  整个 sheet 的合并区数量(给模型的上下文信号)
+ * @param primaryMergeCount 主表块范围内的合并区数量 —— 结构归一判据只认它:
+ *                   Fesod 读数只会被本块范围内的合并区弄脏,远处无关合并格不该否决整表
  * @param headerMd5   归一化表头指纹,命中缓存就完全不用调模型
  * @author cc
  * @since 2026/9/21
@@ -37,6 +39,7 @@ public record SheetEvidence(Long caseId,
                             List<List<String>> headWindow,
                             List<ColStat> colStats,
                             int mergeCount,
+                            int primaryMergeCount,
                             String headerMd5) {
 
     /**

+ 45 - 10
ai-server/src/main/java/com/zsjz/ai/module/aiclean/probe/AiRowReader.java

@@ -36,13 +36,50 @@ public final class AiRowReader {
     }
 
     /**
-     * @param headerRow  表头行,1 基;数据从它的下一行开始
-     * @param limit      最多取多少行,&lt;=0 表示全量
+     * 流式行回调:{@code row} = 列号→值(已拷贝的新 Map,可安全持有),
+     * {@code lineNo} = 原文件物理行号(1 基,空行/表头之外的行也按物理行号计)。
+     */
+    @FunctionalInterface
+    public interface RowHandler {
+
+        void accept(Map<Integer, String> row, int lineNo);
+    }
+
+    /**
+     * 收集型读取(样本/校验用,容忍中途失败):缺几行样本不影响判定结论,
+     * 所以 {@link #forEachRow} 的返回值在这里被刻意忽略。
+     *
+     * @param headerRow 表头行,1 基;数据从它的下一行开始
+     * @param limit     最多取多少行,&lt;=0 表示全量
      */
     public static List<Map<Integer, String>> read(File file, String sheetName, int headerRow, int limit) {
         List<Map<Integer, String>> rows = new ArrayList<>();
+        forEachRow(file, sheetName, headerRow, (row, lineNo) -> {
+            if (limit > 0 && rows.size() >= limit) {
+                return;
+            }
+            rows.add(row);
+        });
+        return rows;
+    }
+
+    /**
+     * 流式读取:逐行回调,不在内存攒列表。
+     *
+     * <p>为什么必须流式:出口 B 的全量入库曾经是「collect 完再写」,几十万行的 sheet
+     * 会在堆里同时驻留原始行 + 清洗行两份 Map(评审 P0-1)。Fesod 本来就是 SAX 逐行回调,
+     * 边读边交给调用方,配 {@code AiDuckDb.BatchWriter} 攒批落库,堆占用恒定。
+     *
+     * <p>装配与 {@code read()} 完全一致({@code headRowNumber(0)} 自管表头 +
+     * 同款两个转换器),保证「同一格两条链路读出来的值一致」——
+     * 这条对齐很重要,否则抽样校验的结论不能外推到正式入库。
+     *
+     * @return false = 读取失败(文件打不开/解析异常/回调抛出);true = 完整读到位。
+     *         半途失败<b>不算成功</b>,调用方必须据此放弃本轮入库而不是拿半截数据当全量
+     */
+    public static boolean forEachRow(File file, String sheetName, int headerRow, RowHandler handler) {
         if (file == null || !file.exists()) {
-            return rows;
+            return false;
         }
         ExcelTypeEnum type = file.getName().toLowerCase().endsWith(".xls")
                 ? ExcelTypeEnum.XLS : ExcelTypeEnum.XLSX;
@@ -50,9 +87,6 @@ public final class AiRowReader {
             var reader = FesodSheet.read(in, new AnalysisEventListener<Map<Integer, String>>() {
                 @Override
                 public void invoke(Map<Integer, String> data, AnalysisContext context) {
-                    if (limit > 0 && rows.size() >= limit) {
-                        return;
-                    }
                     int lineNo = context.readRowHolder().getRowIndex() + 1;
                     if (lineNo <= headerRow) {
                         return;
@@ -60,12 +94,12 @@ public final class AiRowReader {
                     if (data.values().stream().allMatch(v -> v == null || v.isBlank())) {
                         return;
                     }
-                    rows.add(new LinkedHashMap<>(data));
+                    handler.accept(new LinkedHashMap<>(data), lineNo);
                 }
 
                 @Override
                 public void doAfterAllAnalysed(AnalysisContext context) {
-                    // 收集型读取,收尾无需处理
+                    // 流式回调,收尾由调用方负责
                 }
             }).excelType(type).headRowNumber(0).autoTrim(true)
                     .registerConverter(new Str2NullConverter())
@@ -76,9 +110,10 @@ public final class AiRowReader {
             } else {
                 reader.sheet(sheetName).doRead();
             }
+            return true;
         } catch (Exception e) {
-            log.warn("AI 读表失败: file={}, sheet={}, err={}", file.getName(), sheetName, e.getMessage());
+            log.warn("AI 流式读表失败: file={}, sheet={}, err={}", file.getName(), sheetName, e.getMessage());
+            return false;
         }
-        return rows;
     }
 }

+ 3 - 1
ai-server/src/main/java/com/zsjz/ai/module/aiclean/probe/SheetProbe.java

@@ -117,10 +117,12 @@ public final class SheetProbe {
         List<List<String>> dataRows = dataOf(rows, primary);
         List<ColStat> stats = ColStatBuilder.build(dataRows, colCount);
         String md5 = HeaderFingerprint.of(primary.headers());
+        // 主块自己的合并区数:以前算出来了却没用上,needsStructure 误用整个 sheet 的计数 ——
+        // 远处一个无关合并格就能把整张表打入人工(评审 P2-7)
         return new SheetEvidence(caseId, fileId, sheetId, fileName, sheetName,
                 file == null ? null : file.getAbsolutePath(), totalRows,
                 blocks, rows.stream().limit(BlockDetector.HEAD_SCAN_ROWS).toList(),
-                stats, merged.size(), md5);
+                stats, merged.size(), inBlock.size(), md5);
     }
 
     /**

+ 136 - 33
ai-server/src/main/java/com/zsjz/ai/module/aiclean/service/AiCleanService.java

@@ -6,6 +6,7 @@ import com.baomidou.mybatisplus.core.toolkit.Wrappers;
 import com.zsjz.ai.common.base.BasicColumn;
 import com.zsjz.ai.common.cache.GlobalCache;
 import com.zsjz.ai.common.enums.FileCategoryEnum;
+import com.zsjz.ai.common.utils.Json;
 import com.zsjz.ai.common.model.dm.dto.TableRuleDTO;
 import com.zsjz.ai.common.model.dm.entity.FileInfo;
 import com.zsjz.ai.common.model.plat.entity.TableInfo;
@@ -34,6 +35,7 @@ import com.zsjz.ai.module.aiclean.rule.AiRulePlan;
 import com.zsjz.ai.module.aiclean.store.AiCleanStore;
 import com.zsjz.ai.module.aiclean.verify.AiRuleVerifier;
 import com.zsjz.ai.module.dm.mapper.FileInfoMapper;
+import com.zsjz.ai.module.agent.llm.LlmResult;
 import lombok.RequiredArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.stereotype.Service;
@@ -219,18 +221,29 @@ public class AiCleanService {
         if (root == null || StrUtil.isBlank(root.getFilePath()) || !new File(root.getFilePath()).exists()) {
             return new AiJudgeResult.Skip("源文件不存在或已被清理");
         }
+        mark(job, AiCleanStatusEnum.PROBING);
         SheetEvidence ev = SheetProbe.probe(job.getCaseId(), job.getFileId(), job.getSheetId(),
                 root.getFileName(), sheet.getSheetName(), new File(root.getFilePath()));
         if (ev == null) {
             return new AiJudgeResult.Skip("结构探测不可用");
         }
+        mark(job, AiCleanStatusEnum.MATCHING);
         List<String> headers = ev.primaryBlock().headers();
         List<AiTemplateMatcher.AiTemplateWithFields> online = loadOnlineTemplates();
         AiMatchOutcome outcome = matcher.match(ev, online, true);
         job.setHeaderMd5(ev.headerMd5());
+        job.setHeaderRow(ev.primaryBlock().headerRow());
         job.setMatchedBy(outcome.matchedBy());
         job.setMatchTier(outcome.tier().name());
-        job.setLlmCalls(outcome.llmCalled() ? 1 : 0);
+        // 用量遥测:模型名/token 随判定写回 job(finish 时落库)—— 成本看板的原料(评审 P1-5)
+        int llmCalls = outcome.llmCalled() ? 1 : 0;
+        LlmResult llm = outcome.llm();
+        if (llm != null) {
+            job.setModelName(llm.modelName());
+            job.setTokenIn(llm.inputTokens());
+            job.setTokenOut(llm.outputTokens());
+        }
+        job.setLlmCalls(llmCalls);
         if (outcome.candidate() == null) {
             return new AiJudgeResult.NeedReview("没有合适的模板(含 AI 模板库)");
         }
@@ -240,23 +253,40 @@ public class AiCleanService {
         }
         // 选模板那次调用已经把列绑定一起返回了,复用它:单次模型调用实测就要几十秒,
         // 为了拿 binds 再打一次等于把延迟和 token 直接翻倍
-        AiRuleDraft draft = outcome.draft() != null ? outcome.draft()
+        AiTemplateMatcher.LlmAnswer answer = outcome.draft() != null
+                ? new AiTemplateMatcher.LlmAnswer(outcome.draft(), outcome.llm())
                 : matcher.askLlm(ev, headers, List.of(outcome.candidate()));
+        AiRuleDraft draft = answer == null ? null : answer.draft();
         if (draft == null) {
-            job.setLlmCalls(1);
+            // 调用发生但没拿到草稿(含重试耗尽):这次尝试也计入调用数
+            job.setLlmCalls(llmCalls + 1);
             return new AiJudgeResult.Failed("模型未返回可解析的规则草稿");
         }
-        job.setLlmCalls(outcome.draft() != null ? 1 : 2);
+        if (outcome.draft() == null) {
+            // 第二次模型调用(腿 1/2 确定性命中但没有 binds 的场景):次数与用量都要补记
+            llmCalls++;
+            if (answer.usage() != null) {
+                job.setModelName(answer.usage().modelName());
+                job.setTokenIn(nzToken(job.getTokenIn()) + answer.usage().inputTokens());
+                job.setTokenOut(nzToken(job.getTokenOut()) + answer.usage().outputTokens());
+            }
+        }
+        job.setLlmCalls(llmCalls);
+        // 模型原始草稿落库:复盘「哪一版 prompt 判的」以及给后续提示词回归集当语料(评审 P1-5)
+        job.setDraft(StrUtil.maxLength(Json.toStr(draft), 4000));
         List<RuleBind> binds = AiTemplateMatcher.filterBinds(draft.getBinds(), outcome.candidate(),
                 headers.size());
         if (binds.isEmpty()) {
             return new AiJudgeResult.NeedReview("模型未给出任何可用列绑定");
         }
 
+        mark(job, AiCleanStatusEnum.VERIFYING);
         List<Map<Integer, String>> sample = AiRowReader.read(new File(root.getFilePath()),
                 sheet.getSheetName(), ev.primaryBlock().headerRow(), SAMPLE);
         VerifyReport report = verify(binds, outcome.candidate(), headers, sample, outcome);
         job.setConfidence(java.math.BigDecimal.valueOf(report.confidence().overall()));
+        // 抽样报告落库(预览行截到 5 行防膨胀):排障时能看到「哪个列错误率多少、为什么被拒」
+        job.setDetail(detailJson(report));
         if (report.rejectedCols() != null && !report.rejectedCols().isEmpty()) {
             binds = binds.stream().filter(b -> !report.rejectedCols().contains(b.getFileColIndex())).toList();
         }
@@ -309,37 +339,29 @@ public class AiCleanService {
     /**
      * 出口 B:建 AI 模板 + 动态表并写数。
      *
+     * <p><b>流式管道</b>(评审 P0-1):Fesod 逐行回调里边清洗边写 DuckDB(攒批 flush),
+     * 不再把整个 sheet collect 成两份 Map 驻留堆内 —— 几十万行的流水表以前会在这里打爆内存。
+     *
+     * <p><b>写序与 {@link AiCleanStore} 的注释对齐</b>(评审发现的幽灵模板问题):
+     * 模板 id 先行生成(物理表名 {@code ai_t{id}} 依赖它)→ 建表 → 流式写数 →
+     * <b>成功才登记 ONLINE</b>。写数失败/零行就不登记,并删掉刚建的空表 ——
+     * 绝不留下「有元数据没物理表」的模板占着 {@code header_md5} 唯一键。
+     *
      * <p>公开给融合链路用 —— 它自己决定何时刷批次计数,这里只负责「把这块数据落进 AI 表」
      * 并把结果写回 job 的字段(不 finish,由调用方收尾)。
      *
      * @return 写入行数;失败返回 -1
      */
     public long loadViaAiTemplate(AiCleanJob job, FileInfo root, FileInfo sheet, AiJudgeResult.AiPlan plan) {
+        mark(job, AiCleanStatusEnum.LOADING);
         SheetEvidence ev = plan.ev();
         List<RuleBind> binds = plan.binds();
         List<String> headers = plan.headers();
         AiRulePlan.Plan aiPlan = AiRulePlan.forAi(binds, headers, "ai_t");
-        List<Map<String, Object>> rows = new ArrayList<>();
-        for (Map<Integer, String> r : AiRowReader.read(new File(root.getFilePath()),
-                sheet.getSheetName(), ev.primaryBlock().headerRow(), 0)) {
-            Map<String, String> errors = new LinkedHashMap<>();
-            Map<String, Object> cleaned = AiRecordCleaner.cleanRow(r, aiPlan.columns(), errors);
-            if (cleaned == null) {
-                continue;
-            }
-            cleaned.put("id", IdUtil.getSnowflakeNextId());
-            cleaned.put("case_id", job.getCaseId());
-            cleaned.put("file_id", job.getFileId());
-            cleaned.put("sheet_id", job.getSheetId());
-            cleaned.put("block_no", ev.primaryBlock().blockNo());
-            // row_no 存原文件行号,是事后定位与撤销的抓手:读表器已经按表头行截断,这里用序号回填
-            cleaned.put("row_no", rows.size() + ev.primaryBlock().headerRow() + 1);
-            rows.add(cleaned);
-        }
-        if (rows.isEmpty()) {
-            return -1L;
-        }
+        // 模板 id 先行:物理表名依赖它;元数据登记挪到写数成功之后
         AiCleanTemplate tpl = new AiCleanTemplate();
+        tpl.setId(IdUtil.getSnowflakeNextId());
+        tpl.setTableNameEn("ai_t" + tpl.getId());
         tpl.setTableNameCn(StrUtil.blankToDefault(plan.draft().getReason(),
                 plan.candidate() == null ? "AI 识别表" : plan.candidate().nameCn()));
         tpl.setHeaders(String.join(",", headers));
@@ -348,7 +370,6 @@ public class AiCleanService {
         tpl.setBlockNo(ev.primaryBlock().blockNo());
         tpl.setCategory(FileCategoryEnum.parse(plan.draft().getCategory()).name());
         tpl.setMatchTier(job.getMatchTier());
-        tpl.setStatus("ONLINE");
         tpl.setVer(1);
         tpl.setNeedsStruct(0);
         tpl.setCaseIdSrc(job.getCaseId());
@@ -363,16 +384,56 @@ public class AiCleanService {
             r.setVer(1);
             ruleRows.add(r);
         }
-        Long tplId = store.register(tpl, aiPlan.fields(), ruleRows);
-        long written = store.writeRows(tpl, aiPlan.fields(), rows, job.getId());
+        String err = AiDuckDb.createAiTable(tpl.getTableNameEn(), aiPlan.fields());
+        if (err != null) {
+            log.warn("AI 动态表建表失败,跳过写入: table={}, err={}", tpl.getTableNameEn(), err);
+            return -1L;
+        }
+        List<String> cols = aiPlan.fields().stream().map(AiCleanField::getColumnName).toList();
+        long written;
+        java.util.concurrent.atomic.AtomicLong total = new java.util.concurrent.atomic.AtomicLong();
+        try (AiDuckDb.BatchWriter writer = AiDuckDb.openWriter(tpl.getTableNameEn(), cols)) {
+            boolean ok = AiRowReader.forEachRow(new File(root.getFilePath()), sheet.getSheetName(),
+                    ev.primaryBlock().headerRow(), (row, lineNo) -> {
+                        Map<String, String> errors = new LinkedHashMap<>();
+                        Map<String, Object> cleaned = AiRecordCleaner.cleanRow(row, aiPlan.columns(), errors);
+                        if (cleaned == null) {
+                            return;
+                        }
+                        cleaned.put("id", IdUtil.getSnowflakeNextId());
+                        cleaned.put("case_id", job.getCaseId());
+                        cleaned.put("file_id", job.getFileId());
+                        cleaned.put("sheet_id", job.getSheetId());
+                        cleaned.put("block_no", ev.primaryBlock().blockNo());
+                        // row_no 存真实原文件行号(流式回调按物理行回填),是事后定位与撤销的抓手;
+                        // 旧实现用「清洗后序号 + 表头行」近似,sheet 里有空行时对不上号
+                        cleaned.put("row_no", lineNo);
+                        cleaned.put("ai_job_id", job.getId());
+                        cleaned.put("ai_rule_ver", tpl.getVer() == null ? 1 : tpl.getVer());
+                        cleaned.put("category", tpl.getCategory());
+                        // 写失败会在这里抛出,中断流式读取 —— 半截数据绝不能被当成成功
+                        writer.add(cleaned);
+                        total.incrementAndGet();
+                    });
+            written = ok ? writer.finish() : -1L;
+        } catch (Exception e) {
+            log.warn("AI 流式入库失败: jobId={}, table={}, err={}", job.getId(), tpl.getTableNameEn(), e.getMessage());
+            return -1L;
+        }
         if (written <= 0) {
-            // 一行没落库不能算成功:全自动链路没有人看日志,写成 0 行会被上层当成清洗完成
+            // 一行没落库不能算成功(含「读到 0 行」):全自动链路没有人看日志,
+            // 写成 0 行会被上层当成清洗完成。空表删掉,模板不登记 —— 不留幽灵
+            AiDuckDb.dropTable(tpl.getTableNameEn());
             return -1L;
         }
-        job.setAiTemplateId(tplId);
+        // 数据已落盘,现在才登记元数据为 ONLINE 并刷新治理节点(写序收口)
+        tpl.setStatus("ONLINE");
+        store.register(tpl, aiPlan.fields(), ruleRows);
+        store.refreshTree();
+        job.setAiTemplateId(tpl.getId());
         job.setExitCode(AiCleanExitEnum.B.name());
-        job.setRowTotal(rows.size());
-        job.setRowLoaded((int) written);
+        job.setRowTotal((int) Math.min(total.get(), Integer.MAX_VALUE));
+        job.setRowLoaded((int) Math.min(written, Integer.MAX_VALUE));
         return written;
     }
 
@@ -402,9 +463,16 @@ public class AiCleanService {
                 cand.score(), total);
     }
 
-    /** 有合并区、或表头上方存在标题/说明行时认为需要结构归一 */
+    /**
+     * 主表块内有合并区、或表头上方存在标题/说明行时认为需要结构归一。
+     *
+     * <p>合并区只看<b>主块自己的</b>({@link SheetEvidence#primaryMergeCount()}):
+     * 旧实现用整个 sheet 的合并区数一票否决,远处一个无关合并格就能把整张表打入人工
+     * (评审 P2-7);而 Fesod 读数只会被本块范围内的合并区弄脏,块外的不影响主块。
+     * STRUCTURING 状态留给 P1 的结构归一出口(A2/A3),本版本仍不执行。
+     */
     private boolean needsStructure(SheetEvidence ev) {
-        return ev.mergeCount() > 0 || ev.primaryBlock().headerRow() > 1;
+        return ev.primaryMergeCount() > 0 || ev.primaryBlock().headerRow() > 1;
     }
 
     private List<AiTemplateMatcher.AiTemplateWithFields> loadOnlineTemplates() {
@@ -472,4 +540,39 @@ public class AiCleanService {
                 job.getId(), status, job.getExitCode(), job.getCostMs(), reason);
         return job;
     }
+
+    /**
+     * 过程状态落库:让批内能观测「哪张表卡在哪个阶段」——
+     * 以前 PROBING/MATCHING/VERIFYING/LOADING 四个中间状态从未被赋值,
+     * 单表判定(含几十秒的模型调用)期间 job 一直挂在 PENDING(评审 P1-6)。
+     * STRUCTURING 留给 P1 的结构归一出口,本版本不产生。
+     */
+    private void mark(AiCleanJob job, AiCleanStatusEnum status) {
+        job.setStatus(status);
+        job.setUpdateTime(LocalDateTime.now());
+        jobMapper.updateById(job);
+    }
+
+    /**
+     * 抽样报告 → 有界 JSON:预览行只留前 5 行(整份报告是给人复盘的,不是给程序回放的),
+     * 否则 300 行抽样全量进 {@code detail} 列,jobs 列表一次查询会被撑爆。
+     */
+    private static String detailJson(VerifyReport report) {
+        Map<String, Object> detail = new LinkedHashMap<>();
+        detail.put("sampled", report.sampled());
+        detail.put("valid", report.valid());
+        detail.put("invalid", report.invalid());
+        detail.put("estTotalRows", report.estTotalRows());
+        detail.put("perColErrorRate", report.perColErrorRate());
+        detail.put("rejectedCols", report.rejectedCols());
+        detail.put("columnErrors", report.columnErrors());
+        detail.put("confidence", report.confidence());
+        detail.put("previewRows", report.previewRows() == null ? List.of()
+                : report.previewRows().stream().limit(5).toList());
+        return StrUtil.maxLength(Json.toStr(detail), 8000);
+    }
+
+    private static int nzToken(Integer v) {
+        return v == null ? 0 : v;
+    }
 }

+ 95 - 0
ai-server/src/main/java/com/zsjz/ai/module/aiclean/service/AiCleanWatchdog.java

@@ -0,0 +1,95 @@
+package com.zsjz.ai.module.aiclean.service;
+
+import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
+import com.baomidou.mybatisplus.core.toolkit.Wrappers;
+import com.zsjz.ai.module.aiclean.entity.AiCleanBatch;
+import com.zsjz.ai.module.aiclean.enums.AiCleanStageEnum;
+import com.zsjz.ai.module.aiclean.mapper.AiCleanBatchMapper;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.boot.context.event.ApplicationReadyEvent;
+import org.springframework.context.event.EventListener;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+
+import java.time.LocalDateTime;
+import java.util.List;
+
+/**
+ * 智能清洗批次看门狗(评审 P0-2)。
+ *
+ * <p>批次线程是裸虚拟线程,<b>进程重启后批次行会永远停在非终态</b> —— 以前没有任何
+ * 机制能重新拉起,用户再点「智能」只会拿回旧 stage({@code upsertBatch} 见非终态就复用)。
+ * 对照 Celery 的 visibility timeout 与 Temporal 的任务超时:卡死的批次必须有人宣布死亡,
+ * 用户才能重试。标 {@link AiCleanStageEnum#FAILED}(终态)而不是删除 ——
+ * 终态会被 {@code upsertBatch} 重置计数开新一轮,<b>这就是恢复路径本身</b>。
+ *
+ * <h3>两道 sweep</h3>
+ * <ul>
+ *   <li><b>启动恢复</b>:{@code ApplicationReadyEvent} 时把全部非终态批次标 FAILED ——
+ *       虚拟线程不跨重启,这些批次的 worker 一定已经死了。
+ *       安全前提是<b>单实例部署</b>(与 SSE 按 userId 单 emitter、reviewItems 内存 map
+ *       同一部署假设);将来多实例时这里必须改成租约/心跳制。</li>
+ *   <li><b>运行时心跳超时</b>:非终态且 {@code update_time} 超过 stale-minutes 未推进 →
+ *       FAILED。判定/入库每完成一个 sheet 都会 bump {@code update_time}
+ *       ({@code bump()} / {@code stage()}),所以超时无推进 = 线程卡死
+ *       (模型流挂死、DB 连接挂起)。</li>
+ * </ul>
+ *
+ * @author cc
+ * @since 2026/9/29
+ */
+@Slf4j
+@Component
+@RequiredArgsConstructor
+public class AiCleanWatchdog {
+
+    /** 视为「在跑」的非终态阶段 */
+    private static final List<AiCleanStageEnum> RUNNING_STAGES = List.of(
+            AiCleanStageEnum.JUDGING, AiCleanStageEnum.CLEANING_EXISTING, AiCleanStageEnum.LOADING_AI);
+
+    private final AiCleanBatchMapper batchMapper;
+
+    /**
+     * 心跳超时分钟数。判定/入库每完成一个 sheet 都会刷新 {@code update_time},
+     * 单个超大 sheet 的流式入库也可能持续数分钟 —— 30 分钟无推进才认定卡死。
+     */
+    @Value("${zsjz.ai-clean.stale-minutes:30}")
+    private int staleMinutes;
+
+    /**
+     * 启动恢复:进程重启后,上一轮没跑完的批次全部宣布失败。
+     */
+    @EventListener(ApplicationReadyEvent.class)
+    public void failOrphansOnStartup() {
+        int n = batchMapper.update(null, failAll("服务重启,批次中断,请重新发起智能清洗"));
+        if (n > 0) {
+            log.warn("启动恢复:{} 个上一轮遗留的智能清洗批次已标失败,用户可重新发起", n);
+        }
+    }
+
+    /**
+     * 运行时看门狗:心跳超时的批次标 FAILED。
+     * fixedDelay(非 fixedRate)让慢查询自然串行;启动延迟 2 分钟,与启动恢复错开。
+     */
+    @Scheduled(fixedDelayString = "${zsjz.ai-clean.watchdog-interval-ms:60000}",
+            initialDelayString = "${zsjz.ai-clean.watchdog-initial-delay-ms:120000}")
+    public void failStale() {
+        int minutes = Math.max(5, staleMinutes);
+        int n = batchMapper.update(null, failAll("批次心跳超时(超过 " + minutes
+                        + " 分钟无进展),请重新发起智能清洗")
+                .lt(AiCleanBatch::getUpdateTime, LocalDateTime.now().minusMinutes(minutes)));
+        if (n > 0) {
+            log.warn("看门狗:{} 个心跳超时的智能清洗批次已标失败", n);
+        }
+    }
+
+    private LambdaUpdateWrapper<AiCleanBatch> failAll(String reason) {
+        return Wrappers.<AiCleanBatch>lambdaUpdate()
+                .in(AiCleanBatch::getStage, RUNNING_STAGES)
+                .set(AiCleanBatch::getStage, AiCleanStageEnum.FAILED)
+                .set(AiCleanBatch::getLastError, reason)
+                .set(AiCleanBatch::getUpdateTime, LocalDateTime.now());
+    }
+}

+ 51 - 2
ai-server/src/main/java/com/zsjz/ai/module/aiclean/service/AiSmartCleanService.java

@@ -29,6 +29,7 @@ import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Semaphore;
 import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
 
 /**
  * 智能清洗融合编排:一次请求把「既有模板清洗」和「AI 清洗未匹配文件」跑完。
@@ -64,6 +65,14 @@ public class AiSmartCleanService {
     @Value("${zsjz.ai-clean.max-concurrent:3}")
     private int maxConcurrent;
 
+    /**
+     * 单批次 AI 判定的 token 预算上限(0 = 不限,LiteLLM budget 同款护栏)。
+     * 判定每完成一个 sheet 就把真实用量累加进 {@link Stats#tokens},超限后剩余表
+     * 直接转人工 —— 近似闸门:并发在途的判定会略微冲过线,但绝不会无限烧钱。
+     */
+    @Value("${zsjz.ai-clean.batch-token-budget:0}")
+    private long batchTokenBudget;
+
     /** 每任务一个虚拟线程,真正的并发上限由闸门控制(与 FileRecognitionService 同一套) */
     private final ExecutorService executor = Executors.newThreadPerTaskExecutor(
             Thread.ofVirtual().name("ai-clean-", 0).factory());
@@ -161,6 +170,17 @@ public class AiSmartCleanService {
             Map<Long, AiJudgeResult.AiPlan> plans = new ConcurrentHashMap<>();
             List<CompletableFuture<Void>> futures = pending.stream()
                     .map(p -> CompletableFuture.runAsync(() -> {
+                        // 预算护栏(LiteLLM budget 同款思路):近似检查,超限就把剩余表转人工,不再烧钱。
+                        // 放在 acquire 之前 —— 不为一张注定不判的表占闸门;精确性不重要,兜底才重要
+                        if (batchTokenBudget > 0 && st.tokens.get() >= batchTokenBudget) {
+                            aiCleanService.finish(p.job(), AiCleanStatusEnum.NEED_REVIEW,
+                                    "AI token 预算已耗尽,转人工处理", System.currentTimeMillis());
+                            addReview(batchId, p, "AI token 预算已耗尽,转人工处理");
+                            st.needReview.incrementAndGet();
+                            bump(batchId, "needReview");
+                            bump(batchId, "aiJudged");
+                            return;
+                        }
                         // acquire 必须在 try 之外:acquire 本身被中断时 finally 会 release 一个没拿到的许可,闸门被放水
                         try {
                             gate.acquire();
@@ -217,6 +237,8 @@ public class AiSmartCleanService {
     private void handle(Long batchId, PendingAi p, AiJudgeResult result,
                         List<AiJudgeResult.ExistingHit> existingHits,
                         Map<Long, AiJudgeResult.AiPlan> plans, Stats st) {
+        // 批次 token 计入预算:judge 把真实用量写回了 job(空值 = 指纹直通没调模型)
+        st.tokens.addAndGet(nzToken(p.job().getTokenIn()) + nzToken(p.job().getTokenOut()));
         switch (result) {
             case AiJudgeResult.ExistingHit hit -> {
                 existingHits.add(hit);
@@ -350,8 +372,29 @@ public class AiSmartCleanService {
         batch.setNeedReviewCount(0);
         batch.setFailedCount(0);
         batch.setCreateTime(LocalDateTime.now());
-        batchMapper.insert(batch);
-        return new BatchRef(batch, true);
+        batch.setUpdateTime(LocalDateTime.now());
+        try {
+            batchMapper.insert(batch);
+            return new BatchRef(batch, true);
+        } catch (Exception e) {
+            // 并发首点的另一方已经把批次插进去了(唯一键 uk_ai_clean_batch 撞上)。
+            // 与 AiCleanService#insertQuiet 同一套幂等手法,等价于 PG 的
+            // INSERT ... ON CONFLICT DO NOTHING:回落查询,拿对方那行按「进行中」复用,
+            // 绝不重起线程、绝不重复发 doClean(评审 P0-3:原先这里直接 500 抛给用户)
+            exist = batchMapper.selectOne(batchKey(caseId, batchId));
+            if (exist == null) {
+                throw new ServerException(500,
+                        "智能清洗批次登记失败: " + StrUtil.maxLength(e.getMessage(), 200));
+            }
+            boolean running = exist.getStage() != null && !exist.getStage().isFinal();
+            if (!running) {
+                resetCounters(exist, fileCount, sheetTotal, matched);
+                batchMapper.updateById(exist);
+                return new BatchRef(exist, true);
+            }
+            log.info("智能清洗批次已被并发请求登记,复用: batchId={}, stage={}", batchId, exist.getStage());
+            return new BatchRef(exist, false);
+        }
     }
 
     /** 批次行 + 是否该由本次调用启动后台流程 */
@@ -433,6 +476,12 @@ public class AiSmartCleanService {
         private final AtomicInteger aiLoad = new AtomicInteger();
         private final AtomicInteger needReview = new AtomicInteger();
         private final AtomicInteger failed = new AtomicInteger();
+        /** 本轮已消耗的模型 token(判定回写),供批次预算护栏比对 */
+        private final AtomicLong tokens = new AtomicLong();
+    }
+
+    private static int nzToken(Integer v) {
+        return v == null ? 0 : v;
     }
 
     @PreDestroy

+ 3 - 24
ai-server/src/main/java/com/zsjz/ai/module/aiclean/store/AiCleanStore.java

@@ -72,32 +72,11 @@ public class AiCleanStore {
     }
 
     /**
-     * 建表(若不存在)并写入清洗好的行;成功后刷新治理节点。
-     *
-     * @return 实际写入行数;建表失败返回 -1
+     * 刷新 AI 治理节点:数据落库成功、模板登记 ONLINE 之后由入库链路调用 ——
+     * 数据一落盘就补节点,用户不必等下一轮治理才能在树上看到。
      */
-    public long writeRows(AiCleanTemplate tpl, List<AiCleanField> fields,
-                          List<Map<String, Object>> rows, Long jobId) {
-        String err = AiDuckDb.createAiTable(tpl.getTableNameEn(), fields);
-        if (err != null) {
-            log.warn("AI 动态表建表失败,跳过写入: table={}, err={}", tpl.getTableNameEn(), err);
-            return -1L;
-        }
-        List<String> cols = fields.stream().map(AiCleanField::getColumnName).toList();
-        for (Map<String, Object> row : rows) {
-            row.put("ai_job_id", jobId);
-            row.put("ai_rule_ver", tpl.getVer() == null ? 1 : tpl.getVer());
-            row.put("category", tpl.getCategory());
-        }
-        long written = AiDuckDb.insertBatch(tpl.getTableNameEn(), cols, rows);
-        if (written < 0) {
-            // 写失败就不刷树:否则节点会指向一张行数对不上的表
-            log.warn("AI 动态表写入失败,跳过治理节点刷新: table={}", tpl.getTableNameEn());
-            return -1L;
-        }
-        // 数据一落盘就补节点,用户不必等下一轮治理才能在树上看到
+    public void refreshTree() {
         aiGovernService.calcAiTree();
-        return written;
     }
 
     /**

+ 4 - 1
ai-server/src/main/java/com/zsjz/ai/module/dm/ai/FileRecognitionService.java

@@ -8,6 +8,7 @@ import com.zsjz.ai.common.enums.AiParseStatusEnum;
 import com.zsjz.ai.common.enums.FileCategoryEnum;
 import com.zsjz.ai.common.model.dm.entity.FileAiProfile;
 import com.zsjz.ai.module.agent.llm.LlmRequest;
+import com.zsjz.ai.module.agent.llm.LlmRetry;
 import com.zsjz.ai.module.agent.llm.LlmService;
 import com.zsjz.ai.module.dm.mapper.FileAiProfileMapper;
 import jakarta.annotation.PreDestroy;
@@ -240,7 +241,9 @@ public class FileRecognitionService {
                     .temperature(0.1)
                     .timeout(timeout)
                     .build();
-            FileRecognitionResult result = llmService.chatAs(request, FileRecognitionResult.class);
+            // LlmRetry:限流/厂商抖动指数退避重试(评审 P1-4)—— 批处理链路值得为成功率多等两轮
+            FileRecognitionResult result = LlmRetry.call("文件 AI 识别",
+                    () -> llmService.chatAs(request, FileRecognitionResult.class));
             long cost = System.currentTimeMillis() - start;
             FileCategoryEnum category = FileCategoryEnum.parse(result.getCategory());
             writeBack(profile, null, null, AiParseStatusEnum.SUCCESS, null, result,

+ 10 - 0
ai-server/src/main/resources/application.yaml

@@ -101,6 +101,16 @@ zsjz:
     # 启动时后台预热加载模型(false = 懒加载,首次业务调用时才载入)
     eager-load: false
 
+  # 智能清洗(AI 兜底链路,module/aiclean)。
+  # max-concurrent:AI 判定并发闸门(未配置走代码默认 3);模型调用重试由 LlmRetry 固定 3 次封顶。
+  # batch-token-budget:单批次 AI 判定的 token 预算上限,0 = 不限 —— 超限后剩余表自动转人工,
+  # 防止一次异常上传把 token 烧穿(LiteLLM budget 同款护栏,近似检查)。
+  # stale-minutes:批次心跳超时分钟数 —— 每完成一个 sheet 都会刷新 update_time,
+  # 超时无推进的批次由看门狗标 FAILED,用户重新点「智能」即可恢复(进程重启同理,启动时自动清理)。
+  ai-clean:
+    batch-token-budget: 0
+    stale-minutes: 30
+
 mybatis-plus:
   # 实体别名包(各业务域实体)。已核实 10 个域的实体类简单名无重复,可安全启用;
   # 若后续新增实体出现重名,MyBatis 会因别名冲突启动失败,此时应删除本配置(XML 全用全限定名)。

+ 7 - 1
ai-server/src/test/java/com/zsjz/ai/module/aiclean/AiCleanLlmLiveTest.java

@@ -115,9 +115,15 @@ class AiCleanLlmLiveTest {
                 + " (fromAiLib=" + cand.fromAiLib() + ", score=" + cand.score() + ")");
 
         // 复用选模板那次的草稿:单次调用实测几十秒起,再打一次纯属浪费
-        AiRuleDraft draft = outcome.draft() != null ? outcome.draft()
+        AiTemplateMatcher.LlmAnswer answer = outcome.draft() != null
+                ? new AiTemplateMatcher.LlmAnswer(outcome.draft(), outcome.llm())
                 : matcher.askLlm(ev, ev.primaryBlock().headers(), List.of(cand));
+        AiRuleDraft draft = answer == null ? null : answer.draft();
         assertNotNull(draft, "模型未返回可解析结果");
+        if (answer.usage() != null) {
+            System.out.println("  用量: model=" + answer.usage().modelName()
+                    + " in=" + answer.usage().inputTokens() + "t out=" + answer.usage().outputTokens() + "t");
+        }
         System.out.println("  分类=" + draft.getCategory() + " 置信=" + draft.getConfidence()
                 + " 理由=" + draft.getReason());
 

+ 12 - 5
ai-server/src/test/java/com/zsjz/ai/module/aiclean/match/AiTemplateMatcherTest.java

@@ -12,6 +12,7 @@ import com.zsjz.ai.module.aiclean.model.RuleBind;
 import com.zsjz.ai.module.aiclean.model.SheetEvidence;
 import com.zsjz.ai.module.aiclean.support.HeaderFingerprint;
 import com.zsjz.ai.module.agent.llm.LlmRequest;
+import com.zsjz.ai.module.agent.llm.LlmResult;
 import com.zsjz.ai.module.agent.llm.LlmService;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.DisplayName;
@@ -35,7 +36,7 @@ import static org.mockito.Mockito.when;
  * 如果它偷偷走到第三腿,功能照样正确、成本却翻倍,而且没有任何报错会提示 ——
  * 所以这里断言的是 {@code llmCalled=false},不只是结果对。
  *
- * <p>② <b>模型幻觉必须被兜住</b>。{@code chatAs} 只是提示词约束 + 宽松解析,
+ * <p>② <b>模型幻觉必须被兜住</b>。{@code chatAsDetail} 只是提示词约束 + 宽松解析,
  * 没有 schema 强校验,模型完全可能回一个不存在的候选序号;越界一律按未命中处理,
  * 绝不能拿它去落库。
  *
@@ -104,7 +105,7 @@ class AiTemplateMatcherTest {
         SheetEvidence ev = evidenceOf(odd);
         AiRuleDraft draft = new AiRuleDraft();
         draft.setTemplateNo(99);
-        when(llm.chatAs(any(LlmRequest.class), any(Class.class))).thenReturn(draft);
+        when(llm.chatAsDetail(any(LlmRequest.class), any(Class.class))).thenReturn(structured(draft));
 
         TableInfo info = new TableInfo();
         info.setId(78);
@@ -143,7 +144,7 @@ class AiTemplateMatcherTest {
         AiRuleDraft draft = new AiRuleDraft();
         draft.setTemplateNo(1);
         draft.setConfidence(0.3);
-        when(llm.chatAs(any(LlmRequest.class), any(Class.class))).thenReturn(draft);
+        when(llm.chatAsDetail(any(LlmRequest.class), any(Class.class))).thenReturn(structured(draft));
 
         var outcome = matcher.match(evidenceOf(List.of("字段甲", "字段乙")), List.of(), true);
         assertEquals(AiMatchOutcome.BY_LLM, outcome.matchedBy());
@@ -162,7 +163,7 @@ class AiTemplateMatcherTest {
         AiRuleDraft draft = new AiRuleDraft();
         draft.setTemplateNo(1);
         draft.setConfidence(0.92);
-        when(llm.chatAs(any(LlmRequest.class), any(Class.class))).thenReturn(draft);
+        when(llm.chatAsDetail(any(LlmRequest.class), any(Class.class))).thenReturn(structured(draft));
 
         var outcome = matcher.match(evidenceOf(List.of("字段甲", "字段乙")), List.of(), true);
         assertEquals(com.zsjz.ai.module.aiclean.enums.AiMatchTierEnum.STRONG, outcome.tier());
@@ -188,6 +189,12 @@ class AiTemplateMatcherTest {
         return b;
     }
 
+    /** chatAsDetail 的桩:草稿 + 一份虚构用量(askLlm 内部走 LlmRetry 包装它) */
+    private static LlmService.StructuredResult<AiRuleDraft> structured(AiRuleDraft draft) {
+        return new LlmService.StructuredResult<>(draft,
+                new LlmResult("{}", 1L, "mock:test-model", 10, 5, 0L));
+    }
+
     private static TableField field(String cn, String en) {
         TableField f = new TableField();
         f.setFieldNameCn(cn);
@@ -199,6 +206,6 @@ class AiTemplateMatcherTest {
     private static SheetEvidence evidenceOf(List<String> heads) {
         var block = new com.zsjz.ai.module.aiclean.model.BlockHint(0, 1, heads, 2, 10, List.of());
         return new SheetEvidence(1L, 2L, 3L, "f.xlsx", "Sheet1", "/tmp/f.xlsx", 10,
-                List.of(block), List.of(heads), List.of(), 0, HeaderFingerprint.of(heads));
+                List.of(block), List.of(heads), List.of(), 0, 0, HeaderFingerprint.of(heads));
     }
 }