|
|
@@ -0,0 +1,477 @@
|
|
|
+package com.zsjz.ai.module.agent.rag;
|
|
|
+
|
|
|
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
|
|
|
+import com.baomidou.mybatisplus.core.toolkit.Wrappers;
|
|
|
+import com.zsjz.ai.common.enums.HasInnerEnum;
|
|
|
+import com.zsjz.ai.common.exception.ServerException;
|
|
|
+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 com.zsjz.ai.module.agent.config.EtlProperties;
|
|
|
+import com.zsjz.ai.module.agent.entity.RagIndexStatus;
|
|
|
+import com.zsjz.ai.module.agent.mapper.RagIndexStatusMapper;
|
|
|
+import com.zsjz.ai.module.agent.mapper.RagVectorMapper;
|
|
|
+import com.zsjz.ai.module.agent.rag.EmbeddingModelFactory.EmbeddingSpec;
|
|
|
+
|
|
|
+import com.zsjz.ai.module.agent.vo.RagIndexStatusVO;
|
|
|
+import com.zsjz.ai.module.agent.vo.RagSearchResultVO;
|
|
|
+import com.zsjz.ai.module.plat.mapper.TableFieldMapper;
|
|
|
+import com.zsjz.ai.module.plat.mapper.TableInfoMapper;
|
|
|
+import io.agentscope.core.message.TextBlock;
|
|
|
+import lombok.extern.slf4j.Slf4j;
|
|
|
+import org.springframework.context.event.EventListener;
|
|
|
+import org.springframework.stereotype.Service;
|
|
|
+import org.springframework.util.StringUtils;
|
|
|
+
|
|
|
+import java.nio.charset.StandardCharsets;
|
|
|
+import java.security.MessageDigest;
|
|
|
+import java.time.Duration;
|
|
|
+import java.time.LocalDateTime;
|
|
|
+import java.util.*;
|
|
|
+import java.util.concurrent.ConcurrentHashMap;
|
|
|
+import java.util.concurrent.ExecutorService;
|
|
|
+import java.util.concurrent.Executors;
|
|
|
+import java.util.stream.Collectors;
|
|
|
+
|
|
|
+/**
|
|
|
+ * RAG 表结构向量服务:模板元数据 → 嵌入 → 主数据源 pgvector 存储与检索
|
|
|
+ *
|
|
|
+ * <p>元数据来源(全部 master PG,不依赖 DuckDB 系统表):
|
|
|
+ * <ul>
|
|
|
+ * <li>{@code meta_raw_sheet}:按 workspace_id 过滤,得到工作空间内已匹配模板的表
|
|
|
+ * (template_id 非空),行数按同模板多 sheet 求和</li>
|
|
|
+ * <li>{@code table_info}:表英文名(DuckDB 物理表名)+ 表中文名</li>
|
|
|
+ * <li>{@code table_field}:字段英文名 + 字段中文名(注释)+ 类型 + 是否必填</li>
|
|
|
+ * </ul>
|
|
|
+ *
|
|
|
+ * <p>嵌入仍使用未弃用的 {@code io.agentscope.core.embedding.*},向量存取由本服务通过
|
|
|
+ * JDBC 直接管理 pgvector 表(表名 rag_ws{workspaceId}_d{dims},列为
|
|
|
+ * doc_id/content/payload jsonb/embedding vector(dims))。
|
|
|
+ *
|
|
|
+ * <p>索引策略:结构指纹(SHA-256,不含行数)相同且嵌入模型未变时跳过重建;
|
|
|
+ * 重建采用 DROP + 全量写入。
|
|
|
+ */
|
|
|
+@Slf4j
|
|
|
+@Service
|
|
|
+public class RagSchemaService {
|
|
|
+
|
|
|
+ private final EmbeddingModelFactory embeddingModelFactory;
|
|
|
+ private final RagIndexStatusMapper ragIndexStatusMapper;
|
|
|
+ private final TableInfoMapper tableInfoMapper;
|
|
|
+ private final TableFieldMapper tableFieldMapper;
|
|
|
+ private final RagVectorMapper ragVectorMapper;
|
|
|
+ private final EtlProperties etlProperties;
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 索引重建防重入标记
|
|
|
+ */
|
|
|
+ private final Set<Long> indexing = ConcurrentHashMap.newKeySet();
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 后台索引单线程池(守护线程,不阻塞应用关闭)
|
|
|
+ */
|
|
|
+ private final ExecutorService indexExecutor = Executors.newSingleThreadExecutor(r -> {
|
|
|
+ Thread t = new Thread(r, "rag-index");
|
|
|
+ t.setDaemon(true);
|
|
|
+ return t;
|
|
|
+ });
|
|
|
+
|
|
|
+ public RagSchemaService(EmbeddingModelFactory embeddingModelFactory,
|
|
|
+ RagIndexStatusMapper ragIndexStatusMapper,
|
|
|
+ TableInfoMapper tableInfoMapper,
|
|
|
+ TableFieldMapper tableFieldMapper,
|
|
|
+ RagVectorMapper ragVectorMapper,
|
|
|
+ EtlProperties etlProperties) {
|
|
|
+ this.embeddingModelFactory = embeddingModelFactory;
|
|
|
+ this.ragIndexStatusMapper = ragIndexStatusMapper;
|
|
|
+ this.tableInfoMapper = tableInfoMapper;
|
|
|
+ this.tableFieldMapper = tableFieldMapper;
|
|
|
+ this.ragVectorMapper = ragVectorMapper;
|
|
|
+ this.etlProperties = etlProperties;
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 工作空间切换事件:后台异步重建索引(指纹相同自动跳过)
|
|
|
+ */
|
|
|
+ @EventListener
|
|
|
+ public void onWorkspaceSwitched(WorkspaceSwitchedEvent event) {
|
|
|
+ if (!etlProperties.getRag().isEnabled() || !etlProperties.getRag().isAutoIndexOnSwitch()) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ Long workspaceId = event.workspaceId();
|
|
|
+ log.info("工作空间切换,触发后台 RAG 索引: workspaceId={}", workspaceId);
|
|
|
+ indexExecutor.submit(() -> {
|
|
|
+ try {
|
|
|
+ reindexWorkspace(workspaceId);
|
|
|
+ } catch (Exception e) {
|
|
|
+ log.warn("后台 RAG 索引失败: workspaceId={}, error={}", workspaceId, e.getMessage());
|
|
|
+ }
|
|
|
+ });
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * RAG 能力是否可用(配置开关开启且存在 embedding 模型)
|
|
|
+ */
|
|
|
+ public boolean isAvailable() {
|
|
|
+ return etlProperties.getRag().isEnabled() && embeddingModelFactory.resolveDefault() != null;
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 查询工作空间索引状态
|
|
|
+ */
|
|
|
+ public RagIndexStatusVO getStatus(Long workspaceId) {
|
|
|
+ RagIndexStatus status = ragIndexStatusMapper.selectById(workspaceId);
|
|
|
+ if (status == null) {
|
|
|
+ RagIndexStatusVO vo = new RagIndexStatusVO();
|
|
|
+ vo.setWorkspaceId(workspaceId);
|
|
|
+ vo.setStatus(RagIndexStatus.STATUS_IDLE);
|
|
|
+ vo.setTableCount(0);
|
|
|
+ return vo;
|
|
|
+ }
|
|
|
+ return convertStatusToVO(status);
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 重建工作空间表结构索引(同步,全量)
|
|
|
+ *
|
|
|
+ * <p>结构指纹与嵌入模型均未变化时直接返回;否则 DROP 向量表后全量重建。
|
|
|
+ * 同一工作空间重建互斥。只读 master 元数据表,不影响当前 slave(DuckDB)指向。
|
|
|
+ */
|
|
|
+ public RagIndexStatusVO reindexWorkspace(Long workspaceId) {
|
|
|
+ EmbeddingSpec spec = embeddingModelFactory.resolveDefault();
|
|
|
+ if (spec == null) {
|
|
|
+ throw new ServerException(400, "未配置嵌入模型(model 表需存在 type='embedding' 的活跃记录)");
|
|
|
+ }
|
|
|
+ if (!indexing.add(workspaceId)) {
|
|
|
+ throw new ServerException(409, "该工作空间正在构建索引,请稍后再试");
|
|
|
+ }
|
|
|
+ try {
|
|
|
+ return doReindex(workspaceId, spec);
|
|
|
+ } finally {
|
|
|
+ indexing.remove(workspaceId);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ private RagIndexStatusVO doReindex(Long workspaceId, EmbeddingSpec spec) {
|
|
|
+ List<WorkspaceTable> tables = tableInfoMapper.selectList(Wrappers.<TableInfo>lambdaQuery().eq(TableInfo::getHasTable, HasInnerEnum.YES))
|
|
|
+ .stream()
|
|
|
+ .map(t -> new WorkspaceTable(t.getId(), t.getTableNameEn(), t.getTableNameCn()))
|
|
|
+ .toList();
|
|
|
+
|
|
|
+ List<TableField> fields = tableFieldMapper.selectList(new LambdaQueryWrapper<TableField>()
|
|
|
+ .in(TableField::getTableId, tables.stream().map(WorkspaceTable::id).toList())
|
|
|
+ .orderByAsc(TableField::getFieldSort));
|
|
|
+ String fingerprint = fingerprint(tables, fields);
|
|
|
+
|
|
|
+ upsertStatus(workspaceId, RagIndexStatus.STATUS_INDEXING, spec, fingerprint, tables.size(), null);
|
|
|
+ try {
|
|
|
+ String vectorTable = vectorTableName(workspaceId, spec.dimensions());
|
|
|
+ dropVectorTable(vectorTable);
|
|
|
+ List<SchemaDoc> docs = buildDocuments(workspaceId, tables, fields);
|
|
|
+ writeToVectorTable(vectorTable, spec, docs);
|
|
|
+ RagIndexStatus done = upsertStatus(workspaceId, RagIndexStatus.STATUS_SUCCESS, spec,
|
|
|
+ fingerprint, tables.size(), null);
|
|
|
+ log.info("RAG 表结构索引完成: workspaceId={}, tables={}, dims={}, vectorTable={}",
|
|
|
+ workspaceId, tables.size(), spec.dimensions(), vectorTable);
|
|
|
+ return convertStatusToVO(done);
|
|
|
+ } catch (Exception e) {
|
|
|
+ log.error("RAG 表结构索引失败: workspaceId={}, error={}", workspaceId, e.getMessage(), e);
|
|
|
+ upsertStatus(workspaceId, RagIndexStatus.STATUS_FAILED, spec, fingerprint, tables.size(),
|
|
|
+ e.getMessage());
|
|
|
+ throw new ServerException(500, "RAG 索引构建失败: " + e.getMessage());
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 语义检索工作空间表结构:嵌入查询文本后按 cosine 相似度取 top-k
|
|
|
+ *
|
|
|
+ * <p>嵌入模型取索引状态记录的模型(保证与已建索引一致);索引未就绪时返回 409。
|
|
|
+ */
|
|
|
+ public List<RagSearchResultVO> search( String query, Integer topK) {
|
|
|
+ int limit = (topK != null && topK > 0) ? topK : etlProperties.getRag().getTopK();
|
|
|
+
|
|
|
+ RagIndexStatus status = ragIndexStatusMapper.selectOne(Wrappers.emptyWrapper(), false);
|
|
|
+ if (status == null || !RagIndexStatus.STATUS_SUCCESS.equals(status.getStatus())) {
|
|
|
+ throw new ServerException(409, "该工作空间表结构索引不可用(状态: "
|
|
|
+ + (status == null ? "未构建" : status.getStatus()) + "),请先重建或等待构建完成");
|
|
|
+ }
|
|
|
+ EmbeddingSpec spec = embeddingModelFactory.createById(status.getEmbeddingModelId());
|
|
|
+
|
|
|
+ double[] queryVector = embed(spec, query);
|
|
|
+ String literal = toVectorLiteral(queryVector);
|
|
|
+ try {
|
|
|
+ List<Map<String, Object>> rows = ragVectorMapper.searchVector(status.getVectorTable(),
|
|
|
+ literal, etlProperties.getRag().getScoreThreshold(), limit);
|
|
|
+ List<RagSearchResultVO> results = new ArrayList<>(rows.size());
|
|
|
+ for (Map<String, Object> row : rows) {
|
|
|
+ Object payload = row.get("payload");
|
|
|
+ Object score = row.get("score");
|
|
|
+ results.add(convertSearchResult(payload == null ? null : payload.toString(),
|
|
|
+ score instanceof Number n ? n.doubleValue() : 0.0));
|
|
|
+ }
|
|
|
+ return results;
|
|
|
+ } catch (Exception e) {
|
|
|
+ throw new ServerException(500, "表结构向量检索失败: " + e.getMessage());
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ // ==================== 元数据读取(master:meta_raw_sheet / table_info / table_field) ====================
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 工作空间表元信息:物理表名 + 中英文表名 + 行数合计
|
|
|
+ */
|
|
|
+ public record WorkspaceTable(Integer id, String tableNameEn, String tableNameCn) {
|
|
|
+ }
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+ // ==================== 文档构建 ====================
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 待向量的表结构文档:docId + 嵌入文本 + 结构化 payload JSON
|
|
|
+ */
|
|
|
+ record SchemaDoc(String docId, String content, String payloadJson) {
|
|
|
+ }
|
|
|
+
|
|
|
+ private List<SchemaDoc> buildDocuments(Long workspaceId, List<WorkspaceTable> tables,
|
|
|
+ List<TableField> fields) {
|
|
|
+ Map<String, List<TableField>> fieldsByTable = fields.stream()
|
|
|
+ .filter(f -> StringUtils.hasText(f.getTableNameEn()))
|
|
|
+ .collect(Collectors.groupingBy(TableField::getTableNameEn, LinkedHashMap::new,
|
|
|
+ Collectors.toList()));
|
|
|
+
|
|
|
+ List<SchemaDoc> documents = new ArrayList<>(tables.size());
|
|
|
+ for (WorkspaceTable table : tables) {
|
|
|
+ documents.add(buildDocument(workspaceId, table,
|
|
|
+ fieldsByTable.getOrDefault(table.tableNameEn(), List.of())));
|
|
|
+ }
|
|
|
+ return documents;
|
|
|
+ }
|
|
|
+
|
|
|
+ private SchemaDoc buildDocument(Long workspaceId, WorkspaceTable table, List<TableField> fields) {
|
|
|
+ StringBuilder sb = new StringBuilder();
|
|
|
+ sb.append("表名: ").append(table.tableNameEn());
|
|
|
+ if (StringUtils.hasText(table.tableNameCn())) {
|
|
|
+ sb.append("(中文名: ").append(table.tableNameCn()).append(")");
|
|
|
+ }
|
|
|
+ sb.append("\n字段:");
|
|
|
+ for (TableField field : fields) {
|
|
|
+ sb.append("\n- ").append(field.getFieldNameEn()).append(" (").append(field.getFieldType());
|
|
|
+ if (StringUtils.hasText(field.getFieldNameCn())) {
|
|
|
+ sb.append(", 注释: ").append(field.getFieldNameCn());
|
|
|
+ }
|
|
|
+ if (field.getRequired() != null && field.getRequired() == 1) {
|
|
|
+ sb.append(", 必填");
|
|
|
+ }
|
|
|
+ sb.append(")");
|
|
|
+ }
|
|
|
+
|
|
|
+ List<Map<String, Object>> columnPayload = new ArrayList<>(fields.size());
|
|
|
+ for (TableField field : fields) {
|
|
|
+ Map<String, Object> cm = new LinkedHashMap<>();
|
|
|
+ cm.put("name", field.getFieldNameEn());
|
|
|
+ cm.put("type", field.getFieldType());
|
|
|
+ if (StringUtils.hasText(field.getFieldNameCn())) {
|
|
|
+ cm.put("comment", field.getFieldNameCn());
|
|
|
+ }
|
|
|
+ if (field.getRequired() != null && field.getRequired() == 1) {
|
|
|
+ cm.put("required", true);
|
|
|
+ }
|
|
|
+ columnPayload.add(cm);
|
|
|
+ }
|
|
|
+ Map<String, Object> payload = new LinkedHashMap<>();
|
|
|
+ payload.put("workspaceId", workspaceId);
|
|
|
+ payload.put("tableName", table.tableNameEn());
|
|
|
+ if (StringUtils.hasText(table.tableNameCn())) {
|
|
|
+ payload.put("tableComment", table.tableNameCn());
|
|
|
+ }
|
|
|
+ payload.put("columns", columnPayload);
|
|
|
+
|
|
|
+ return new SchemaDoc("ws" + workspaceId + ":" + table.tableNameEn(),
|
|
|
+ sb.toString(), Json.toStr(payload));
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 结构指纹:对排序后的(表/字段/类型/注释/必填)元组做 SHA-256(不含行数,避免数据量变化触发重建)
|
|
|
+ */
|
|
|
+ private String fingerprint(List<WorkspaceTable> tables, List<TableField> fields) {
|
|
|
+ StringBuilder sb = new StringBuilder();
|
|
|
+ tables.stream()
|
|
|
+ .sorted(Comparator.comparing(WorkspaceTable::tableNameEn))
|
|
|
+ .forEach(t -> sb.append("T|").append(t.tableNameEn()).append('|')
|
|
|
+ .append(nullToEmpty(t.tableNameCn())).append('\n'));
|
|
|
+ fields.stream()
|
|
|
+ .sorted(Comparator.comparing(TableField::getTableNameEn,
|
|
|
+ Comparator.nullsLast(Comparator.naturalOrder()))
|
|
|
+ .thenComparing(TableField::getFieldNameEn,
|
|
|
+ Comparator.nullsLast(Comparator.naturalOrder())))
|
|
|
+ .forEach(f -> sb.append("C|").append(nullToEmpty(f.getTableNameEn()))
|
|
|
+ .append('|').append(nullToEmpty(f.getFieldNameEn()))
|
|
|
+ .append('|').append(nullToEmpty(f.getFieldType().getValue()))
|
|
|
+ .append('|').append(nullToEmpty(f.getFieldNameCn()))
|
|
|
+ .append('|').append(f.getRequired() == null ? "" : f.getRequired())
|
|
|
+ .append('\n'));
|
|
|
+ try {
|
|
|
+ MessageDigest digest = MessageDigest.getInstance("SHA-256");
|
|
|
+ return HexFormat.of().formatHex(digest.digest(sb.toString().getBytes(StandardCharsets.UTF_8)));
|
|
|
+ } catch (Exception e) {
|
|
|
+ throw new ServerException(500, "计算结构指纹失败: " + e.getMessage());
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ private static String nullToEmpty(String s) {
|
|
|
+ return s == null ? "" : s;
|
|
|
+ }
|
|
|
+
|
|
|
+ // ==================== 向量存储管理(应用层 JDBC 直连 pgvector) ====================
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 向量表名:rag_ws{workspaceId}_d{dims}
|
|
|
+ */
|
|
|
+ private String vectorTableName(Long workspaceId, int dims) {
|
|
|
+ return "rag_ws" + workspaceId + "_d" + dims;
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 建表并逐条写入文档向量(每文档一次嵌入调用)
|
|
|
+ */
|
|
|
+ private void writeToVectorTable(String tableName, EmbeddingSpec spec, List<SchemaDoc> docs) {
|
|
|
+ try {
|
|
|
+ ragVectorMapper.createVectorTable(tableName, spec.dimensions());
|
|
|
+ } catch (Exception e) {
|
|
|
+ throw new ServerException(500, "创建向量表失败(请检查 PG 是否已安装 pgvector 扩展): " + e.getMessage());
|
|
|
+ }
|
|
|
+ try {
|
|
|
+ for (SchemaDoc doc : docs) {
|
|
|
+ double[] vector = embed(spec, doc.content());
|
|
|
+ ragVectorMapper.insertVectorDoc(tableName, doc.docId(), doc.content(),
|
|
|
+ doc.payloadJson(), toVectorLiteral(vector));
|
|
|
+ }
|
|
|
+ } catch (Exception e) {
|
|
|
+ throw new ServerException(500, "写入向量数据失败: " + e.getMessage());
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 嵌入文本并校验维度一致性
|
|
|
+ */
|
|
|
+ private double[] embed(EmbeddingSpec spec, String content) {
|
|
|
+ double[] vector = spec.model()
|
|
|
+ .embed(TextBlock.builder().text(content).build())
|
|
|
+ .block(Duration.ofSeconds(60));
|
|
|
+ if (vector == null || vector.length != spec.dimensions()) {
|
|
|
+ throw new ServerException(500, "嵌入向量维度异常(期望 " + spec.dimensions()
|
|
|
+ + ",实际 " + (vector == null ? 0 : vector.length)
|
|
|
+ + "),请检查 model.config 中 dimensions 配置");
|
|
|
+ }
|
|
|
+ return vector;
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * double[] → pgvector 字面量(如 [0.1,0.2,...])
|
|
|
+ */
|
|
|
+ private static String toVectorLiteral(double[] vector) {
|
|
|
+ StringBuilder sb = new StringBuilder(vector.length * 10).append('[');
|
|
|
+ for (int i = 0; i < vector.length; i++) {
|
|
|
+ if (i > 0) {
|
|
|
+ sb.append(',');
|
|
|
+ }
|
|
|
+ sb.append(vector[i]);
|
|
|
+ }
|
|
|
+ return sb.append(']').toString();
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 删除向量表
|
|
|
+ */
|
|
|
+ private void dropVectorTable(String tableName) {
|
|
|
+ try {
|
|
|
+ ragVectorMapper.dropVectorTable(tableName);
|
|
|
+ log.info("已删除向量表: {}", tableName);
|
|
|
+ } catch (Exception e) {
|
|
|
+ log.warn("删除向量表失败: {}, error={}", tableName, e.getMessage());
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ // ==================== 状态与转换 ====================
|
|
|
+
|
|
|
+ private RagIndexStatus upsertStatus(Long workspaceId, String status, EmbeddingSpec spec,
|
|
|
+ String fingerprint, int tableCount, String error) {
|
|
|
+ RagIndexStatus entity = ragIndexStatusMapper.selectById(workspaceId);
|
|
|
+ boolean exists = entity != null;
|
|
|
+ if (!exists) {
|
|
|
+ entity = new RagIndexStatus();
|
|
|
+ entity.setWorkspaceId(workspaceId);
|
|
|
+ }
|
|
|
+ entity.setEmbeddingModelId(spec.modelId());
|
|
|
+ entity.setEmbeddingModelName(spec.modelName());
|
|
|
+ entity.setDimensions(spec.dimensions());
|
|
|
+ entity.setVectorTable(vectorTableName(workspaceId, spec.dimensions()));
|
|
|
+ entity.setSchemaFingerprint(fingerprint);
|
|
|
+ entity.setTableCount(tableCount);
|
|
|
+ entity.setStatus(status);
|
|
|
+ entity.setLastError(error);
|
|
|
+ entity.setUpdateAt(LocalDateTime.now());
|
|
|
+ if (RagIndexStatus.STATUS_SUCCESS.equals(status)) {
|
|
|
+ entity.setIndexedAt(LocalDateTime.now());
|
|
|
+ }
|
|
|
+ if (exists) {
|
|
|
+ ragIndexStatusMapper.updateById(entity);
|
|
|
+ } else {
|
|
|
+ ragIndexStatusMapper.insert(entity);
|
|
|
+ }
|
|
|
+ return entity;
|
|
|
+ }
|
|
|
+
|
|
|
+ private RagIndexStatusVO convertStatusToVO(RagIndexStatus entity) {
|
|
|
+ RagIndexStatusVO vo = new RagIndexStatusVO();
|
|
|
+ vo.setWorkspaceId(entity.getWorkspaceId());
|
|
|
+ vo.setEmbeddingModelId(entity.getEmbeddingModelId());
|
|
|
+ vo.setEmbeddingModelName(entity.getEmbeddingModelName());
|
|
|
+ vo.setDimensions(entity.getDimensions());
|
|
|
+ vo.setVectorTable(entity.getVectorTable());
|
|
|
+ vo.setSchemaFingerprint(entity.getSchemaFingerprint());
|
|
|
+ vo.setTableCount(entity.getTableCount());
|
|
|
+ vo.setStatus(entity.getStatus());
|
|
|
+ vo.setLastError(entity.getLastError());
|
|
|
+ vo.setIndexedAt(entity.getIndexedAt());
|
|
|
+ vo.setUpdateAt(entity.getUpdateAt());
|
|
|
+ return vo;
|
|
|
+ }
|
|
|
+
|
|
|
+ @SuppressWarnings("unchecked")
|
|
|
+ private RagSearchResultVO convertSearchResult(String payloadJson, double score) {
|
|
|
+ RagSearchResultVO vo = new RagSearchResultVO();
|
|
|
+ vo.setScore(score);
|
|
|
+ if (!StringUtils.hasText(payloadJson)) {
|
|
|
+ return vo;
|
|
|
+ }
|
|
|
+ try {
|
|
|
+ Map<String, Object> payload = Json.objectMapper().readValue(payloadJson, Map.class);
|
|
|
+ vo.setTableName((String) payload.get("tableName"));
|
|
|
+ vo.setTableComment((String) payload.get("tableComment"));
|
|
|
+ Object rowCount = payload.get("rowCount");
|
|
|
+ if (rowCount instanceof Number n) {
|
|
|
+ vo.setRowCount(n.longValue());
|
|
|
+ }
|
|
|
+ Object columns = payload.get("columns");
|
|
|
+ if (columns instanceof List<?> list) {
|
|
|
+ List<RagSearchResultVO.ColumnMeta> metas = new ArrayList<>(list.size());
|
|
|
+ for (Object item : list) {
|
|
|
+ if (item instanceof Map<?, ?> m) {
|
|
|
+ RagSearchResultVO.ColumnMeta meta = new RagSearchResultVO.ColumnMeta();
|
|
|
+ meta.setName((String) m.get("name"));
|
|
|
+ meta.setType((String) m.get("type"));
|
|
|
+ meta.setComment((String) m.get("comment"));
|
|
|
+ metas.add(meta);
|
|
|
+ }
|
|
|
+ }
|
|
|
+ vo.setColumns(metas);
|
|
|
+ }
|
|
|
+ } catch (Exception e) {
|
|
|
+ log.warn("解析向量 payload 失败: {}", e.getMessage());
|
|
|
+ }
|
|
|
+ return vo;
|
|
|
+ }
|
|
|
+
|
|
|
+}
|