cc 3 недель назад
Родитель
Сommit
de317fb977

+ 197 - 0
.trae/documents/PreDataListener-DuckDB动态建表存储方案.md

@@ -0,0 +1,197 @@
+# PreDataListener 新增 DuckDB 动态建表存储
+
+## 概述
+
+在 [PreDataListener.java](../../../ai-server/src/main/java/com/zsjz/ai/module/dm/pre/PreDataListener.java) 中新增逻辑:解析 Excel/CSV 时,**每个 sheet 在 DuckDB(slave 数据源,即当前案件库)动态创建一张表**,用 `DuckDBAppender` 流式追加**全部数据行**;表结构为 `file_id`、`sheet_id` 两个关联列 + `col_0`、`col_1`… 数据列,**全部 VARCHAR**。表名统一使用 **`raw_` 前缀**(`raw_{sheetId}`)。现有 JSON 预览逻辑(前 2000 行写临时文件)保持不变。
+
+用户已确认的两个决策:
+- 列名按列索引生成:`col_0`、`col_1`…(与表头内容无关)
+- 数据范围:全部行(放开 2000 行截断,仅对 DuckDB 写入放开;JSON 预览仍限前 2000 行)
+
+## 现状分析(Phase 1 探索结论)
+
+| 事项 | 结论 |
+|---|---|
+| 调用链 | [DmService.java](../../../ai-server/src/main/java/com/zsjz/ai/module/dm/service/DmService.java) L344/L363 → `PreTask.call()` → `PreDataListener`;多文件虚拟线程并行(≤20 并发),每文件独立 listener 实例,sheet 串行读取 |
+| DuckDB 接入范式 | [AbstractDataLoader.java](../../../ai-server/src/main/java/com/zsjz/ai/module/dm/clean/AbstractDataLoader.java) L23-26:`SpringUtil.getBean(DataSource.class)` → 强转 `DynamicRoutingDataSource` → `getDataSource(DS_KEY_SLAVE)` → `(DuckDBConnection) getConnection()` → `createAppender(DEFAULT_SCHEMA, tableName)`;关闭顺序:先 appender 后 connection |
+| slave 数据源 | [DuckdbUnpooledDataSource.java](../../../ai-server/src/main/java/com/zsjz/ai/common/config/DuckdbUnpooledDataSource.java):`getConnection()` 返回 `duplicate()` 的新连接,共享同一案件库;开案时由 CaseInfoService 切换。未开案时 slave 回退 master(PG),强转会失败 → 需降级 |
+| 驱动 | duckdb_jdbc(pom 声明 1.5.5.1,本地 1.5.0.0 源码已核实):`append(String)` 对 null 安全(内部 appendNull);appender 创建时绑定表 schema,要求表已存在、每行列数与表列数一致 |
+| 现有截断 | `invoke()` 中 `if (BATCH_COUNT < rowIndex) return;` 同时截断了 JSON 预览和(未来的)DB 写入,需拆分 |
+| 空表行为 | sheet 无任何数据行时 `invoke()` 不会执行、`fileSheet` 为 null(系统本就不为空 sheet 生成 FileInfo/雪花 sheetId) |
+
+## 改动方案
+
+只改 2 个既有文件,不新建文件。
+
+### 1. PreDataListener.java(核心改动)
+
+**新增字段**
+
+```java
+/** 当前 sheet 的 DuckDB 连接(每 sheet 独立,sheet 结束即关闭) */
+private DuckDBConnection rawConn;
+/** 当前 sheet 的 Appender(初始化失败或异常关闭后为 null,本 sheet 跳过 DB 写入) */
+private DuckDBAppender rawAppender;
+/** 动态表当前数据列数(col_0 ~ col_{n-1}) */
+private int rawColCount;
+/** 动态表表名:raw_{sheetId} */
+private String rawTableName;
+/** 当前 sheet 是否已尝试初始化动态表(成功失败都只试一次) */
+private boolean rawTableInited;
+```
+
+**invoke() 重构**(去掉整段早退,JSON 预览改为按 `rowIndex <= BATCH_COUNT` 门控,行为与原来完全一致;`buildFileSheet` 幂等,放开截断后无副作用)
+
+```java
+@Override
+public void invoke(Map<Integer, String> rowData, AnalysisContext analysisContext) {
+    Integer rowIndex = analysisContext.readSheetHolder().getRowIndex();
+    // 构建当前工作表的文件信息对象
+    buildFileSheet(analysisContext);
+    // 前 20 行且未匹配到模板时,尝试匹配表头
+    if (rowIndex <= 20 && !matched.get()) {
+        findTableHead(rowData, rowIndex);
+    }
+    // 全部行写入 DuckDB 动态表
+    appendRawRow(rowData);
+    // JSON 预览仍只保留前 2000 行(与原逻辑一致)
+    if (rowIndex <= BATCH_COUNT) {
+        preDataList.add(Json.toStr(rowData.values()));
+        if (preDataList.size() == BATCH_COUNT) {
+            writePreTmp(analysisContext.readSheetHolder().getSheetName());
+        }
+    }
+}
+```
+
+**新增方法**
+
+```java
+/** 追加一行到当前 sheet 的动态表;首行时懒初始化建表,行宽超过现有列数时动态扩列 */
+private void appendRawRow(Map<Integer, String> rowData) {
+    if (!rawTableInited) {
+        initRawTable(rowData);
+    }
+    if (rawAppender == null) {
+        return; // 初始化失败,本 sheet 跳过 DB 写入
+    }
+    try {
+        // 稀疏行:列数按最大 key+1 计算,后续更宽的行通过 ALTER 扩列
+        int width = rowData.isEmpty() ? 0 : Collections.max(rowData.keySet()) + 1;
+        if (width > rawColCount) {
+            growRawColumns(width);
+        }
+        rawAppender.beginRow();
+        rawAppender.append(String.valueOf(fileId));
+        rawAppender.append(String.valueOf(fileSheet.getId()));
+        for (int i = 0; i < rawColCount; i++) {
+            rawAppender.append(rowData.get(i)); // null 安全,缺失列自动写 NULL
+        }
+        rawAppender.endRow();
+    } catch (Exception e) {
+        log.error("动态表写入失败,本 sheet 剩余行跳过, table={}", rawTableName, e);
+        closeRawTable(); // 关闭并置空,避免 appender 行状态损坏影响后续
+    }
+}
+
+/** 首行懒初始化:取 slave 连接、CREATE TABLE、createAppender;失败仅记日志不影响预览 */
+private void initRawTable(Map<Integer, String> rowData) {
+    rawTableInited = true;
+    rawTableName = "raw_" + fileSheet.getId();
+    try {
+        DynamicRoutingDataSource ds = (DynamicRoutingDataSource) SpringUtil.getBean(DataSource.class);
+        rawConn = (DuckDBConnection) ds.getDataSource(StrConsts.DS_KEY_SLAVE).getConnection();
+        rawColCount = Math.max(1, rowData.isEmpty() ? 1 : Collections.max(rowData.keySet()) + 1);
+        // 全 VARCHAR:file_id、sheet_id 关联列 + col_0..col_{n-1}
+        StringJoiner cols = new StringJoiner(", ");
+        cols.add("file_id VARCHAR").add("sheet_id VARCHAR");
+        for (int i = 0; i < rawColCount; i++) {
+            cols.add("col_" + i + " VARCHAR");
+        }
+        try (Statement st = rawConn.createStatement()) {
+            st.execute("CREATE TABLE IF NOT EXISTS " + rawTableName + " (" + cols + ")");
+        }
+        rawAppender = rawConn.createAppender(DuckDBConnection.DEFAULT_SCHEMA, rawTableName);
+        log.info("动态表初始化成功, table={}, colCount={}", rawTableName, rawColCount);
+    } catch (Exception e) {
+        // 典型场景:未开案时 slave 回退为 PG,强转失败 → 降级为不写 DB,预览流程不受影响
+        log.error("动态表初始化失败,本 sheet 跳过 DuckDB 存储, table={}", rawTableName, e);
+        closeRawTable();
+    }
+}
+
+/** 行宽超过现有列数:先关 appender(flush 已缓冲行)→ ALTER TABLE 补列 → 重建 appender */
+private void growRawColumns(int width) throws SQLException {
+    rawAppender.close();
+    try (Statement st = rawConn.createStatement()) {
+        for (int i = rawColCount; i < width; i++) {
+            st.execute("ALTER TABLE " + rawTableName + " ADD COLUMN col_" + i + " VARCHAR");
+        }
+    }
+    rawColCount = width;
+    rawAppender = rawConn.createAppender(DuckDBConnection.DEFAULT_SCHEMA, rawTableName);
+}
+
+/** 释放当前 sheet 的 DuckDB 资源(幂等,先 appender 后 connection,与 AbstractDataLoader 顺序一致) */
+private void closeRawTable() {
+    if (rawAppender != null) {
+        try { rawAppender.close(); } catch (Exception e) { log.error("关闭动态表 Appender 失败, table={}", rawTableName, e); }
+        rawAppender = null;
+    }
+    if (rawConn != null) {
+        try { rawConn.close(); } catch (Exception e) { log.error("关闭动态表连接失败, table={}", rawTableName, e); }
+        rawConn = null;
+    }
+}
+
+/** 对外兜底释放(PreTask 异常中断时调用;正常流程各 sheet 结束时已释放,此处为空操作) */
+public void close() {
+    closeRawTable();
+}
+```
+
+**doAfterAllAnalysed()**:在现有 `writePreTmp(sheetName)` 后、重置 `fileSheet = null` 前增加
+
+```java
+// 关闭当前 sheet 的动态表资源并重置状态
+closeRawTable();
+rawTableInited = false;
+```
+
+**新增 import**:`cn.hutool.extra.spring.SpringUtil`、`com.baomidou.dynamic.datasource.DynamicRoutingDataSource`、`com.zsjz.ai.common.constants.StrConsts`、`org.duckdb.DuckDBAppender`、`org.duckdb.DuckDBConnection`、`javax.sql.DataSource`、`java.sql.SQLException`、`java.sql.Statement`、`java.util.StringJoiner`
+
+### 2. PreTask.java(小改)
+
+读取流程加 finally 兜底,防止 reader 异常中断时连接泄漏(appender/连接为 duplicate 连接,泄漏会占用案件库资源):
+
+```java
+try (ExcelReader reader = FesodSheet.read(...).build()) {
+    ...
+    reader.readAll();
+} catch (Exception e) {
+    ... // 原有异常分支不变
+} finally {
+    // 兜底释放监听器持有的 DuckDB 资源(正常流程已在 doAfterAllAnalysed 释放,幂等)
+    preDataListener.close();
+}
+```
+
+## 假设与决策
+
+1. **表名** `raw_{sheetId}`(用户指定统一 `raw_` 前缀;sheetId 为雪花 ID,全局唯一,表名可由 sheetId 确定性推导,无需在 FileInfo 上加字段)
+2. **file_id / sheet_id 也用 VARCHAR**(用户要求"字段都用 varchar",按字符串写入)
+3. **动态扩列**:首行定初始列数,后续出现更宽的行时 `ALTER TABLE ADD COLUMN` + 重建 appender(DuckDB appender 创建时绑定 schema);已有行新列为 NULL
+4. **空 sheet 不建表**(无任何行的 sheet 本就没有 FileInfo/sheetId)
+5. **降级策略**:未开案 / slave 非 DuckDB / 建表失败 → 记 error 日志、本 sheet 跳过 DB 写入,预览主流程不受影响;行写入中途失败 → 关闭资源、剩余行跳过
+6. **并发安全**:每个 listener(每文件)独立 duplicate 连接;表名含唯一 sheetId 无跨文件冲突;DuckDB MVCC 支持多连接并发 append
+7. **孤儿表清理不在本次范围**(重新解析、删除文件、清批次会遗留 `raw_*` 表);如需要,后续可在 `deletedFiles`/`clearUploadBatch` 中按 sheetId DROP —— 本次不做
+8. **JSON 预览行为完全不变**(前 2000 行 + 原有 flush/覆盖语义)
+
+## 验证步骤
+
+1. **编译验证**(按工作区规则用 IDEA MCP):`get_file_problems` 检查 PreDataListener.java / PreTask.java,再 `build_project`(或 `mvn compile -f pom.xml`)确认无编译错误
+2. **运行时手动验证**:
+   - 开案后通过上传接口传一个多 sheet 的 xlsx 和一个 csv
+   - 用返回的 FileInfo 树取 sheetId,在案件 DuckDB 中执行 `SHOW TABLES` 确认出现 `raw_{sheetId}`;`SELECT * FROM raw_{sheetId} LIMIT 5` 核对 file_id/sheet_id/col_x 值
+   - 传一个 >2048 行的文件,`SELECT count(*)` 核对全量行数;核对 `SELECT max(col_n)` 是否覆盖最宽行
+   - 未开案时跑本地路径预览(`preProcessFile`),确认仅日志报错、预览正常返回

+ 146 - 44
ai-server/src/main/java/com/zsjz/ai/module/dm/pre/PreDataListener.java

@@ -2,26 +2,29 @@ package com.zsjz.ai.module.dm.pre;
 
 import cn.hutool.core.collection.CollUtil;
 import cn.hutool.core.collection.ListUtil;
-import cn.hutool.core.io.FileUtil;
 import cn.hutool.core.map.MapUtil;
-import cn.hutool.core.util.CharsetUtil;
 import cn.hutool.core.util.IdUtil;
 import cn.hutool.core.util.StrUtil;
+import cn.hutool.extra.spring.SpringUtil;
+import com.baomidou.dynamic.datasource.DynamicRoutingDataSource;
 import com.google.common.collect.Lists;
 import com.zsjz.ai.common.cache.GlobalCache;
-import com.zsjz.ai.common.constants.PathConst;
+import com.zsjz.ai.common.constants.StrConsts;
 import com.zsjz.ai.common.enums.FileSheetStateEnum;
 import com.zsjz.ai.common.model.dm.dto.TableRuleDTO;
 import com.zsjz.ai.common.model.dm.entity.FileInfo;
 import com.zsjz.ai.common.model.dm.mapstruct.TableFieldConvert;
 import com.zsjz.ai.common.model.plat.entity.TableField;
 import com.zsjz.ai.common.model.plat.entity.TableInfo;
-import com.zsjz.ai.common.utils.Json;
 import lombok.extern.slf4j.Slf4j;
 import org.apache.fesod.sheet.context.AnalysisContext;
 import org.apache.fesod.sheet.event.AnalysisEventListener;
+import org.duckdb.DuckDBAppender;
+import org.duckdb.DuckDBConnection;
 
-import java.nio.file.Path;
+import javax.sql.DataSource;
+import java.sql.SQLException;
+import java.sql.Statement;
 import java.time.LocalDateTime;
 import java.util.*;
 import java.util.concurrent.atomic.AtomicBoolean;
@@ -30,7 +33,8 @@ import java.util.stream.Collectors;
 
 /**
  * 预览数据监听器
- * 监听 Excel/CSV 文件读取过程,解析文件结构、匹配模板、提取预览数据
+ * 监听 Excel/CSV 文件读取过程,解析文件结构、匹配模板
+ * 为每个 sheet 在 DuckDB 动态创建 raw_{sheetId} 表(全 VARCHAR),用 DuckDBAppender 流式追加全部数据行
  * 实现流式数据处理,避免大文件内存溢出
  */
 @Slf4j
@@ -57,24 +61,39 @@ public class PreDataListener extends AnalysisEventListener<Map<Integer, String>>
     private final String fileName;
 
     /**
-     * 批量写入阈值,每读取 2000 行写入一次临时文件
+     * 是否已匹配到模板的标识(原子操作,线程安全)
      */
-    private static final int BATCH_COUNT = 2000;
+    private final AtomicBoolean matched = new AtomicBoolean(false);
 
     /**
-     * 预览数据列表,临时存储待写入的 JSON 格式数据
+     * 当前工作表的文件信息对象
      */
-    private final List<String> preDataList = ListUtil.toList();
+    private FileInfo fileSheet;
 
     /**
-     * 是否已匹配到模板的标识(原子操作,线程安全)
+     * 当前 sheet 的 DuckDB 连接(每 sheet 独立,sheet 结束即关闭)
      */
-    private final AtomicBoolean matched = new AtomicBoolean(false);
+    private DuckDBConnection rawConn;
 
     /**
-     * 当前工作表的文件信息对象
+     * 当前 sheet 的 Appender(初始化失败或异常关闭后为 null,本 sheet 跳过 DB 写入)
      */
-    private FileInfo fileSheet;
+    private DuckDBAppender rawAppender;
+
+    /**
+     * 动态表当前数据列数(col_0 ~ col_{n-1})
+     */
+    private int rawColCount;
+
+    /**
+     * 动态表表名:raw_{sheetId}
+     */
+    private String rawTableName;
+
+    /**
+     * 当前 sheet 是否已尝试初始化动态表(成功失败都只试一次)
+     */
+    private boolean rawTableInited;
 
     /**
      * 构造函数
@@ -93,51 +112,135 @@ public class PreDataListener extends AnalysisEventListener<Map<Integer, String>>
 
     /**
      * 处理每一行数据
-     * 读取行数据、匹配模板、收集预览数据并定期写入临时文件
+     * 读取行数据、匹配模板、全部行写入 DuckDB 动态表
      *
      * @param rowData         当前行的数据,Map 的 key 为列索引,value 为单元格值
      * @param analysisContext 分析上下文,包含工作表信息
      */
     @Override
     public void invoke(Map<Integer, String> rowData, AnalysisContext analysisContext) {
-        // 只处理前 2000 行用于预览
         Integer rowIndex = analysisContext.readSheetHolder().getRowIndex();
-        if (BATCH_COUNT < rowIndex) {
-            return;
-        }
         // 构建当前工作表的文件信息对象
         buildFileSheet(analysisContext);
         // 前 20 行且未匹配到模板时,尝试匹配表头
         if (rowIndex <= 20 && !matched.get()) {
             findTableHead(rowData, rowIndex);
         }
-        // 将行数据转换为 JSON 格式并添加到预览列表
-        String jsonStr = Json.toStr(rowData.values());
-        preDataList.add(jsonStr);
-        // 达到批量写入阈值时,写入临时文件
-        if (preDataList.size() == BATCH_COUNT) {
-            String sheetName = analysisContext.readSheetHolder().getSheetName();
-            writePreTmp(sheetName);
-        }
+        // 全部行写入 DuckDB 动态表
+        appendRawRow(rowData);
     }
 
     /**
-     * 将预览数据写入临时文件
-     * 按工作表分别存储预览数据,用于后续展示
+     * 追加一行到当前 sheet 的动态表
+     * 首行时懒初始化建表,行宽超过现有列数时动态扩列
      *
-     * @param sheetName 工作表名称
+     * @param rowData 当前行数据,key 为列索引
      */
-    private void writePreTmp(String sheetName) {
-        if (CollUtil.isEmpty(preDataList)) {
+    private void appendRawRow(Map<Integer, String> rowData) {
+        if (!rawTableInited) {
+            initRawTable(rowData);
+        }
+        if (rawAppender == null) {
             return;
         }
-        // 构建临时文件路径:tmpPath/batchId/fileId_sheetName.json
-        sheetName = StrUtil.emptyToDefault(sheetName, "csv");
-        Path resolve = PathConst.TMP_PATH.resolve(batchId.toString()).resolve(fileId + "_" + sheetName + ".json");
-        FileUtil.writeLines(preDataList, resolve.toFile(), CharsetUtil.CHARSET_UTF_8);
-        preDataList.clear();
-        // 更新文件信息中的临时文件路径
-        fileSheet.setFilePath(resolve.toString());
+        try {
+            // 稀疏行:列数按最大 key+1 计算,后续更宽的行通过 ALTER 扩列
+            int width = rowData.isEmpty() ? 0 : Collections.max(rowData.keySet()) + 1;
+            if (width > rawColCount) {
+                growRawColumns(width);
+            }
+            rawAppender.beginRow();
+            rawAppender.append(String.valueOf(fileId));
+            rawAppender.append(String.valueOf(fileSheet.getId()));
+            for (int i = 0; i < rawColCount; i++) {
+                // append(String) 对 null 安全,缺失列自动写 NULL
+                rawAppender.append(rowData.get(i));
+            }
+            rawAppender.endRow();
+        } catch (Exception e) {
+            log.error("动态表写入失败,本 sheet 剩余行跳过, table={}", rawTableName, e);
+            // 关闭并置空,避免 appender 行状态损坏影响后续
+            closeRawTable();
+        }
+    }
+
+    /**
+     * 首行懒初始化动态表
+     * 取 slave(DuckDB)连接、CREATE TABLE、创建 Appender;失败仅记日志,不影响预览主流程
+     *
+     * @param rowData 首行数据,用于确定初始列数
+     */
+    private void initRawTable(Map<Integer, String> rowData) {
+        rawTableInited = true;
+        rawTableName = "raw_" + fileSheet.getId();
+        try {
+            DynamicRoutingDataSource ds = (DynamicRoutingDataSource) SpringUtil.getBean(DataSource.class);
+            rawConn = (DuckDBConnection) ds.getDataSource(StrConsts.DS_KEY_SLAVE).getConnection();
+            rawColCount = Math.max(1, rowData.isEmpty() ? 1 : Collections.max(rowData.keySet()) + 1);
+            // 全 VARCHAR:file_id、sheet_id 关联列 + col_0..col_{n-1}
+            StringJoiner cols = new StringJoiner(", ");
+            cols.add("file_id VARCHAR").add("sheet_id VARCHAR");
+            for (int i = 0; i < rawColCount; i++) {
+                cols.add("col_" + i + " VARCHAR");
+            }
+            try (Statement st = rawConn.createStatement()) {
+                st.execute("CREATE TABLE IF NOT EXISTS " + rawTableName + " (" + cols + ")");
+            }
+            rawAppender = rawConn.createAppender(DuckDBConnection.DEFAULT_SCHEMA, rawTableName);
+            log.info("动态表初始化成功, table={}, colCount={}", rawTableName, rawColCount);
+        } catch (Exception e) {
+            // 典型场景:未开案时 slave 回退为 master(PG),强转失败 → 降级为不写 DB
+            log.error("动态表初始化失败,本 sheet 跳过 DuckDB 存储, table={}", rawTableName, e);
+            closeRawTable();
+        }
+    }
+
+    /**
+     * 行宽超过现有列数时动态扩列
+     * 先关 appender(flush 已缓冲行)→ ALTER TABLE 补列 → 重建 appender
+     *
+     * @param width 新的列数
+     */
+    private void growRawColumns(int width) throws SQLException {
+        rawAppender.close();
+        try (Statement st = rawConn.createStatement()) {
+            for (int i = rawColCount; i < width; i++) {
+                st.execute("ALTER TABLE " + rawTableName + " ADD COLUMN col_" + i + " VARCHAR");
+            }
+        }
+        rawColCount = width;
+        rawAppender = rawConn.createAppender(DuckDBConnection.DEFAULT_SCHEMA, rawTableName);
+    }
+
+    /**
+     * 释放当前 sheet 的 DuckDB 资源
+     * 幂等,先 appender 后 connection,与 AbstractDataLoader 关闭顺序一致
+     */
+    private void closeRawTable() {
+        if (rawAppender != null) {
+            try {
+                rawAppender.close();
+            } catch (Exception e) {
+                log.error("关闭动态表 Appender 失败, table={}", rawTableName, e);
+            }
+            rawAppender = null;
+        }
+        if (rawConn != null) {
+            try {
+                rawConn.close();
+            } catch (Exception e) {
+                log.error("关闭动态表连接失败, table={}", rawTableName, e);
+            }
+            rawConn = null;
+        }
+    }
+
+    /**
+     * 对外兜底释放 DuckDB 资源
+     * PreTask 读取异常中断时调用;正常流程各 sheet 结束时已释放,此处为空操作
+     */
+    public void close() {
+        closeRawTable();
     }
 
     /**
@@ -199,19 +302,18 @@ public class PreDataListener extends AnalysisEventListener<Map<Integer, String>>
 
     /**
      * 所有数据分析完成后的回调
-     * 写入剩余的预览数据并重置状态
+     * 关闭动态表资源并重置状态,准备处理下一个工作表
      *
      * @param analysisContext 分析上下文
      */
     @Override
     public void doAfterAllAnalysed(AnalysisContext analysisContext) {
-        String sheetName = analysisContext.readSheetHolder().getSheetName();
-        // 写入剩余的预览数据
-        writePreTmp(sheetName);
+        // 关闭当前 sheet 的动态表资源并重置状态
+        closeRawTable();
+        rawTableInited = false;
         // 重置状态,准备处理下一个工作表
         fileSheet = null;
         matched.set(false);
-        preDataList.clear();
     }
 
     /**