From a016f6b62a2dcb9f4898942bb5725a426932a471 Mon Sep 17 00:00:00 2001 From: mengnankkkk Date: Wed, 23 Sep 2026 19:07:18 +0800 Subject: [PATCH 1/7] fix: init --- .../github/malonetalk/entity/ColumnInfo.java | 1 + .../entity/LogicalTableRelation.java | 2 + .../mapper/ColumnSemanticInfoMapper.java | 9 +-- .../mapper/LogicalTableRelationMapper.java | 6 +- .../malonetalk/mapper/TableInfoMapper.java | 2 +- .../column/ColumnSemanticServiceImpl.java | 19 +++-- .../relation/RelationSemanticServiceImpl.java | 53 ++++++++----- .../sync/SemanticSyncApplyService.java | 31 +++++--- .../table/TableSemanticServiceImpl.java | 4 +- .../mapper/ColumnSemanticInfoMapper.xml | 61 ++++++++------- .../mapper/LogicalTableRelationMapper.xml | 67 ++++++++++------- .../main/resources/mapper/TableInfoMapper.xml | 11 ++- .../sync/SemanticSyncApplyServiceTest.java | 75 +++++++++++++++++++ sql/data_source.sql | 30 +++++--- sql/migration_primary_key_relations.sql | 67 +++++++++++++++++ 15 files changed, 314 insertions(+), 124 deletions(-) create mode 100644 data-agent-backend/src/test/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyServiceTest.java create mode 100644 sql/migration_primary_key_relations.sql diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/entity/ColumnInfo.java b/data-agent-backend/src/main/java/io/github/malonetalk/entity/ColumnInfo.java index 5f035d99..0995b23f 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/entity/ColumnInfo.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/entity/ColumnInfo.java @@ -26,6 +26,7 @@ public class ColumnInfo { private Integer id; private Integer datasourceId; + private Integer tableId; private String tableName; private String columnName; private String physicalColumnDescription; diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/entity/LogicalTableRelation.java b/data-agent-backend/src/main/java/io/github/malonetalk/entity/LogicalTableRelation.java index 70ae2964..77caed96 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/entity/LogicalTableRelation.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/entity/LogicalTableRelation.java @@ -25,9 +25,11 @@ public class LogicalTableRelation { private Integer id; private Integer datasourceId; + private Integer sourceTableId; private String sourceTableName; private String sourceColumnNamesJson; private String sourceColumnSignature; + private Integer targetTableId; private String targetTableName; private String targetColumnNamesJson; private String targetColumnSignature; diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java index b4ab0f66..53dd4c2e 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java @@ -53,13 +53,8 @@ ColumnInfo selectByDatasourceIdAndTableNameAndColumnName( int updatePhysicalCacheFields(ColumnInfo columnInfo); - int markPhysicalMissingByIds( - @Param("datasourceId") Integer datasourceId, - @Param("ids") List ids, - @Param("now") LocalDateTime now); - - int deleteByDatasourceId(@Param("datasourceId") Integer datasourceId); + int markPhysicalMissingByIds(@Param("ids") List ids, @Param("now") LocalDateTime now); - int deleteByDatasourceIdAndIds( + int resetSemanticFieldsByIds( @Param("datasourceId") Integer datasourceId, @Param("ids") List ids); } diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/LogicalTableRelationMapper.java b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/LogicalTableRelationMapper.java index b975a568..72c50a83 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/LogicalTableRelationMapper.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/LogicalTableRelationMapper.java @@ -50,18 +50,18 @@ List selectPageByDatasourceIdAndSourceTable( int updateEnabled( @Param("id") Integer id, @Param("datasourceId") Integer datasourceId, - @Param("sourceTableName") String sourceTableName, + @Param("sourceTableId") Integer sourceTableId, @Param("isEnabled") Boolean isEnabled, @Param("updateTime") LocalDateTime updateTime); int deleteById( @Param("id") Integer id, @Param("datasourceId") Integer datasourceId, - @Param("sourceTableName") String sourceTableName); + @Param("sourceTableId") Integer sourceTableId); int deleteByIdsAndSourceTable( @Param("datasourceId") Integer datasourceId, - @Param("sourceTableName") String sourceTableName, + @Param("sourceTableId") Integer sourceTableId, @Param("ids") List ids); int deleteByIds(@Param("datasourceId") Integer datasourceId, @Param("ids") List ids); diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/TableInfoMapper.java b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/TableInfoMapper.java index 3b70cafc..a8952855 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/TableInfoMapper.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/TableInfoMapper.java @@ -33,7 +33,7 @@ public interface TableInfoMapper { int updatePhysicalCacheFields(TableInfo tableInfo); - int deleteByDatasourceIdAndIds( + int resetSemanticFieldsByIds( @Param("datasourceId") Integer datasourceId, @Param("ids") List ids); List selectByDatasourceId(@Param("datasourceId") Integer datasourceId); diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java index aae7d34f..6bf1f067 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java @@ -28,8 +28,10 @@ import io.github.malonetalk.dto.semantic.ColumnSemanticUpdateRequest; import io.github.malonetalk.entity.ColumnInfo; import io.github.malonetalk.entity.Datasource; +import io.github.malonetalk.entity.TableInfo; import io.github.malonetalk.exception.BusinessException; import io.github.malonetalk.mapper.ColumnSemanticInfoMapper; +import io.github.malonetalk.mapper.TableInfoMapper; import io.github.malonetalk.service.DatasourceService; import io.github.malonetalk.service.semantic.SemanticMergeService; import io.github.malonetalk.utils.SemanticUtils; @@ -45,6 +47,7 @@ public class ColumnSemanticServiceImpl implements ColumnSemanticService { private final DatasourceService datasourceService; + private final TableInfoMapper tableInfoMapper; private final ColumnSemanticInfoMapper columnSemanticInfoMapper; private final SemanticMergeService semanticMergeService; private final SemanticConverter semanticConverter; @@ -92,9 +95,17 @@ public void updateColumnSemantic(String tableName, ColumnSemanticUpdateRequest r columnSemanticInfoMapper.selectByDatasourceIdAndTableNameAndColumnName( request.datasourceId(), normalizedTableName, normalizedColumnName); if (existing == null) { + TableInfo tableInfo = + tableInfoMapper.selectByDatasourceIdAndTableName( + request.datasourceId(), normalizedTableName); + if (tableInfo == null) { + throw BusinessException.of( + ErrorCode.RESOURCE_NOT_FOUND, + "Table semantic metadata does not exist: " + normalizedTableName); + } ColumnInfo columnInfo = new ColumnInfo(); columnInfo.setDatasourceId(request.datasourceId()); - columnInfo.setTableName(normalizedTableName); + columnInfo.setTableId(tableInfo.getId()); columnInfo.setColumnName(normalizedColumnName); columnInfo.setColumnDescription(SemanticUtils.trimToNull(request.columnDescription())); columnInfo.setSemanticType(request.semanticType()); @@ -105,7 +116,6 @@ public void updateColumnSemantic(String tableName, ColumnSemanticUpdateRequest r columnSemanticInfoMapper.insert(columnInfo); return; } - existing.setTableName(normalizedTableName); existing.setColumnName(normalizedColumnName); existing.setColumnDescription(SemanticUtils.trimToNull(request.columnDescription())); existing.setSemanticType(request.semanticType()); @@ -131,8 +141,7 @@ public void resetColumnSemantic(Integer datasourceId, String tableName, String c throw BusinessException.of( ErrorCode.RESOURCE_NOT_FOUND, "Column semantic metadata does not exist."); } - columnSemanticInfoMapper.deleteByDatasourceIdAndIds( - datasourceId, List.of(existing.getId())); + columnSemanticInfoMapper.resetSemanticFieldsByIds(datasourceId, List.of(existing.getId())); } @Override @@ -178,7 +187,7 @@ public int resetColumnSemantics( + normalizedTableName + "."); } - return columnSemanticInfoMapper.deleteByDatasourceIdAndIds(datasourceId, matchedIds); + return columnSemanticInfoMapper.resetSemanticFieldsByIds(datasourceId, matchedIds); } private void requireDatasource(Integer datasourceId) { diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java index f734e5a6..3041d8a2 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java @@ -211,7 +211,7 @@ public boolean updateRelationSemanticEnabled( return logicalTableRelationMapper.updateEnabled( request.relationId(), request.datasourceId(), - relation.getSourceTableName(), + relation.getSourceTableId(), request.enabled(), LocalDateTime.now()) > 0; @@ -224,7 +224,7 @@ public boolean deleteRelationSemantic( requireDatasource(datasourceId); LogicalTableRelation relation = requireRelation(datasourceId, tableName, relationId); return logicalTableRelationMapper.deleteById( - relationId, datasourceId, relation.getSourceTableName()) + relationId, datasourceId, relation.getSourceTableId()) > 0; } @@ -238,6 +238,7 @@ public int deleteRelationSemantics( if (relationIds == null || relationIds.isEmpty()) { return 0; } + TableInfo sourceTable = requireTable(datasourceId, normalizedTableName, "sourceTable"); List matchedIds = logicalTableRelationMapper .selectByDatasourceIdAndSourceTable(datasourceId, normalizedTableName) @@ -254,7 +255,7 @@ public int deleteRelationSemantics( + "."); } return logicalTableRelationMapper.deleteByIdsAndSourceTable( - datasourceId, normalizedTableName, relationIds); + datasourceId, sourceTable.getId(), relationIds); } private void requireDatasource(Integer datasourceId) { @@ -301,10 +302,14 @@ private void applyRelationUpdate( } private void ensureRelationEndpointsOperable(LogicalTableRelation relation) { - ensureTableOperable( - relation.getDatasourceId(), relation.getSourceTableName(), "sourceTable"); - ensureTableOperable( - relation.getDatasourceId(), relation.getTargetTableName(), "targetTable"); + TableInfo sourceTable = + ensureTableOperable( + relation.getDatasourceId(), relation.getSourceTableName(), "sourceTable"); + TableInfo targetTable = + ensureTableOperable( + relation.getDatasourceId(), relation.getTargetTableName(), "targetTable"); + relation.setSourceTableId(sourceTable.getId()); + relation.setTargetTableId(targetTable.getId()); List sourceColumns = logicalTableRelationHelper.fromJson( relation.getSourceColumnNamesJson(), "sourceColumnNames"); @@ -323,16 +328,11 @@ private void ensureRelationEndpointsOperable(LogicalTableRelation relation) { "targetColumnNames"); } - private void ensureTableOperable(Integer datasourceId, String tableName, String fieldName) { - TableInfo tableInfo = - tableInfoMapper.selectByDatasourceIdAndTableName(datasourceId, tableName); - if (tableInfo == null) { - throw BusinessException.of( - ErrorCode.RESOURCE_NOT_FOUND, - fieldName + " " + tableName + " semantic metadata does not exist."); - } + private TableInfo ensureTableOperable( + Integer datasourceId, String tableName, String fieldName) { + TableInfo tableInfo = requireTable(datasourceId, tableName, fieldName); if (SemanticAvailabilityHelper.isTableAvailable(tableInfo, UsageLevelEnum.USER_OPERATION)) { - return; + return tableInfo; } throw BusinessException.of( ErrorCode.DATA_CONFLICT, @@ -343,6 +343,17 @@ private void ensureTableOperable(Integer datasourceId, String tableName, String tableInfo, UsageLevelEnum.USER_OPERATION))); } + private TableInfo requireTable(Integer datasourceId, String tableName, String fieldName) { + TableInfo tableInfo = + tableInfoMapper.selectByDatasourceIdAndTableName(datasourceId, tableName); + if (tableInfo == null) { + throw BusinessException.of( + ErrorCode.RESOURCE_NOT_FOUND, + fieldName + " " + tableName + " semantic metadata does not exist."); + } + return tableInfo; + } + private void ensureColumnsOperable( Integer datasourceId, String tableName, List columnNames, String fieldName) { for (String columnName : columnNames) { @@ -406,12 +417,14 @@ private LogicalTableRelation requireRelation( Integer datasourceId, String tableName, Integer relationId) { RequestAssert.requireNonNull(relationId, "relationId cannot be null."); LogicalTableRelation relation = logicalTableRelationMapper.selectById(relationId); + TableInfo sourceTable = + requireTable( + datasourceId, + logicalTableRelationHelper.normalizeTableName(tableName, "tableName"), + "sourceTable"); if (relation == null || !datasourceId.equals(relation.getDatasourceId()) - || !relation.getSourceTableName() - .equals( - logicalTableRelationHelper.normalizeTableName( - tableName, "tableName"))) { + || !sourceTable.getId().equals(relation.getSourceTableId())) { throw BusinessException.of( ErrorCode.RESOURCE_NOT_FOUND, "Logical relation does not exist."); } diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyService.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyService.java index b00be53c..7a49a3ec 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyService.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyService.java @@ -91,7 +91,7 @@ public List applyTableSync( } markMissingTables(datasourceId, missingTableNames, now); - markMissingColumns(datasourceId, missingColumnIds, now); + markMissingColumns(missingColumnIds, now); return results; } @@ -148,7 +148,7 @@ public List refreshPhysicalStatus( LocalDateTime now = LocalDateTime.now(); markMissingTables(datasourceId, tableNamesToMarkMissing, now); - markMissingColumns(datasourceId, columnIdsToMarkMissing, now); + markMissingColumns(columnIdsToMarkMissing, now); return results; } @@ -162,7 +162,6 @@ private void batchSavePresentSchema( } List newTables = new ArrayList<>(); - List newColumns = new ArrayList<>(); for (TableSyncSource table : presentTables) { String tableKey = tableKey(table.tableName()); TableInfo tableInfo = buildPhysicalTableInfo(datasourceId, table); @@ -173,12 +172,24 @@ private void batchSavePresentSchema( tableInfo.setId(existingTable.getId()); tableInfoMapper.updatePhysicalCacheFields(tableInfo); } + } + if (!newTables.isEmpty()) { + tableInfoMapper.batchUpsertPhysicalCache(newTables); + } + List presentTableNames = + presentTables.stream().map(TableSyncSource::tableName).toList(); + Map persistedTableIndex = + loadSemanticTableIndex(datasourceId, presentTableNames); + List newColumns = new ArrayList<>(); + for (TableSyncSource table : presentTables) { + String tableKey = tableKey(table.tableName()); + TableInfo persistedTable = persistedTableIndex.get(tableKey); Map existingColumnIndex = loadColumnIndex(columnsByTableName.getOrDefault(tableKey, List.of())); for (ColumnSyncSource column : table.columns()) { ColumnInfo columnInfo = - buildPhysicalColumnInfo(datasourceId, table.tableName(), column); + buildPhysicalColumnInfo(datasourceId, persistedTable.getId(), column); ColumnInfo existingColumn = existingColumnIndex.get(columnKey(column.columnName())); if (existingColumn == null) { newColumns.add(columnInfo); @@ -188,9 +199,6 @@ private void batchSavePresentSchema( } } } - if (!newTables.isEmpty()) { - tableInfoMapper.batchUpsertPhysicalCache(newTables); - } if (!newColumns.isEmpty()) { columnSemanticInfoMapper.batchUpsertPhysicalCache(newColumns); } @@ -204,10 +212,9 @@ private void markMissingTables( } } - private void markMissingColumns( - Integer datasourceId, List missingColumnIds, LocalDateTime now) { + private void markMissingColumns(List missingColumnIds, LocalDateTime now) { if (!missingColumnIds.isEmpty()) { - columnSemanticInfoMapper.markPhysicalMissingByIds(datasourceId, missingColumnIds, now); + columnSemanticInfoMapper.markPhysicalMissingByIds(missingColumnIds, now); } } @@ -362,10 +369,10 @@ private TableInfo buildPhysicalTableInfo(Integer datasourceId, TableSyncSource t } private ColumnInfo buildPhysicalColumnInfo( - Integer datasourceId, String tableName, ColumnSyncSource column) { + Integer datasourceId, Integer tableId, ColumnSyncSource column) { ColumnInfo columnInfo = new ColumnInfo(); columnInfo.setDatasourceId(datasourceId); - columnInfo.setTableName(tableName); + columnInfo.setTableId(tableId); columnInfo.setColumnName(column.columnName()); columnInfo.setPhysicalColumnDescription(column.description()); columnInfo.setColumnDescription(column.description()); diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java index 2e07a265..421d3fcf 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java @@ -182,7 +182,7 @@ public void resetTableSemantic(Integer datasourceId, String tableName) { throw BusinessException.of( ErrorCode.RESOURCE_NOT_FOUND, "Table semantic metadata does not exist."); } - tableInfoMapper.deleteByDatasourceIdAndIds(datasourceId, List.of(existing.getId())); + tableInfoMapper.resetSemanticFieldsByIds(datasourceId, List.of(existing.getId())); } @Override @@ -222,7 +222,7 @@ public int resetTableSemantics(Integer datasourceId, List tableNames) { + datasourceId + "."); } - return tableInfoMapper.deleteByDatasourceIdAndIds(datasourceId, matchedIds); + return tableInfoMapper.resetSemanticFieldsByIds(datasourceId, matchedIds); } private void requireDatasource(Integer datasourceId) { diff --git a/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml b/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml index adfa2ef7..37224fed 100644 --- a/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml @@ -5,6 +5,7 @@ + @@ -20,24 +21,24 @@ - SELECT t.* - FROM column_info t - WHERE t.id = ( - SELECT t2.id - FROM column_info t2 - WHERE t2.datasource_id = t.datasource_id - AND LOWER(t2.table_name) = LOWER(t.table_name) - AND LOWER(t2.column_name) = LOWER(t.column_name) + SELECT c.*, table_meta.table_name + FROM column_info c + INNER JOIN table_info table_meta ON table_meta.id = c.table_id + WHERE c.id = ( + SELECT c2.id + FROM column_info c2 + WHERE c2.table_id = c.table_id + AND LOWER(c2.column_name) = LOWER(c.column_name) ORDER BY CASE - WHEN t2.column_description IS NOT NULL AND t2.column_description <> '' THEN 0 - WHEN t2.semantic_type IS NOT NULL AND t2.semantic_type <> '' THEN 0 - WHEN t2.is_visible = 0 THEN 0 + WHEN c2.column_description IS NOT NULL AND c2.column_description <> '' THEN 0 + WHEN c2.semantic_type IS NOT NULL AND c2.semantic_type <> '' THEN 0 + WHEN c2.is_visible = 0 THEN 0 ELSE 1 END ASC, - t2.update_time DESC, - t2.create_time DESC, - t2.id DESC + c2.update_time DESC, + c2.create_time DESC, + c2.id DESC LIMIT 1 ) @@ -102,11 +103,11 @@ INSERT INTO column_info ( - datasource_id, table_name, column_name, physical_column_description, type_name, + datasource_id, table_id, column_name, physical_column_description, type_name, primary_key, index_info, column_description, semantic_type, is_visible, physical_status, create_time, update_time ) VALUES ( - #{datasourceId}, #{tableName}, #{columnName}, #{physicalColumnDescription}, + #{datasourceId}, #{tableId}, #{columnName}, #{physicalColumnDescription}, #{typeName}, COALESCE(#{primaryKey}, 0), #{indexInfo}, #{columnDescription}, #{semanticType}, COALESCE(#{isVisible}, 1), COALESCE(#{physicalStatus}, 1), COALESCE(#{createTime}, NOW()), COALESCE(#{updateTime}, NOW()) @@ -115,13 +116,13 @@ INSERT INTO column_info ( - datasource_id, table_name, column_name, physical_column_description, type_name, + datasource_id, table_id, column_name, physical_column_description, type_name, primary_key, index_info, column_description, is_visible, physical_status, create_time, update_time ) VALUES ( - #{column.datasourceId}, #{column.tableName}, #{column.columnName}, + #{column.datasourceId}, #{column.tableId}, #{column.columnName}, #{column.physicalColumnDescription}, #{column.typeName}, COALESCE(#{column.primaryKey}, 0), #{column.indexInfo}, #{column.columnDescription}, @@ -142,7 +143,6 @@ UPDATE column_info - table_name = #{tableName}, column_name = #{columnName}, column_description = #{columnDescription}, semantic_type = #{semanticType}, @@ -155,7 +155,6 @@ UPDATE column_info - table_name = #{tableName}, column_name = #{columnName}, physical_column_description = #{physicalColumnDescription}, column_description = COALESCE(NULLIF(column_description, ''), #{columnDescription}), @@ -172,24 +171,24 @@ UPDATE column_info SET physical_status = 0, update_time = #{now} - WHERE datasource_id = #{datasourceId} - AND id IN + WHERE id IN #{id} AND (physical_status IS NULL OR physical_status <> 0) - - DELETE FROM column_info WHERE datasource_id = #{datasourceId} - - - - DELETE FROM column_info - WHERE datasource_id = #{datasourceId} - AND id IN + + UPDATE column_info c + INNER JOIN table_info t ON t.id = c.table_id + SET c.column_description = c.physical_column_description, + c.semantic_type = NULL, + c.is_visible = 1, + c.update_time = NOW() + WHERE t.datasource_id = #{datasourceId} + AND c.id IN #{id} - + diff --git a/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml b/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml index 8543939c..12c466da 100644 --- a/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml @@ -5,9 +5,11 @@ + + @@ -18,61 +20,70 @@ + + SELECT relation.*, source_table.table_name AS source_table_name, + target_table.table_name AS target_table_name + FROM logical_table_relation relation + INNER JOIN table_info source_table ON source_table.id = relation.source_table_id + INNER JOIN table_info target_table ON target_table.id = relation.target_table_id + + INSERT INTO logical_table_relation ( - datasource_id, source_table_name, source_column_names_json, source_column_signature, - target_table_name, target_column_names_json, target_column_signature, relation_type, + datasource_id, source_table_id, source_column_names_json, source_column_signature, + target_table_id, target_column_names_json, target_column_signature, relation_type, description, is_enabled, create_time, update_time ) VALUES ( - #{datasourceId}, #{sourceTableName}, #{sourceColumnNamesJson}, #{sourceColumnSignature}, - #{targetTableName}, #{targetColumnNamesJson}, #{targetColumnSignature}, #{relationType}, + #{datasourceId}, #{sourceTableId}, #{sourceColumnNamesJson}, #{sourceColumnSignature}, + #{targetTableId}, #{targetColumnNamesJson}, #{targetColumnSignature}, #{relationType}, #{description}, #{isEnabled}, #{createTime}, #{updateTime} ) @@ -80,10 +91,10 @@ UPDATE logical_table_relation - source_table_name = #{sourceTableName}, + source_table_id = #{sourceTableId}, source_column_names_json = #{sourceColumnNamesJson}, source_column_signature = #{sourceColumnSignature}, - target_table_name = #{targetTableName}, + target_table_id = #{targetTableId}, target_column_names_json = #{targetColumnNamesJson}, target_column_signature = #{targetColumnSignature}, relation_type = #{relationType}, @@ -100,20 +111,20 @@ update_time = #{updateTime} WHERE id = #{id} AND datasource_id = #{datasourceId} - AND LOWER(source_table_name) = LOWER(#{sourceTableName}) + AND source_table_id = #{sourceTableId} DELETE FROM logical_table_relation WHERE id = #{id} AND datasource_id = #{datasourceId} - AND LOWER(source_table_name) = LOWER(#{sourceTableName}) + AND source_table_id = #{sourceTableId} DELETE FROM logical_table_relation WHERE datasource_id = #{datasourceId} - AND LOWER(source_table_name) = LOWER(#{sourceTableName}) + AND source_table_id = #{sourceTableId} AND id IN #{id} diff --git a/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml b/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml index 5483f433..7b63d9ec 100644 --- a/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml @@ -174,13 +174,18 @@ WHERE LOWER(domain) = LOWER(#{domain}) - - DELETE FROM table_info + + UPDATE table_info + SET table_description = physical_table_description, + data_granularity = NULL, + domain = 'default', + is_visible = 1, + update_time = NOW() WHERE datasource_id = #{datasourceId} AND id IN #{id} - + diff --git a/data-agent-backend/src/test/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyServiceTest.java b/data-agent-backend/src/test/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyServiceTest.java new file mode 100644 index 00000000..de226198 --- /dev/null +++ b/data-agent-backend/src/test/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyServiceTest.java @@ -0,0 +1,75 @@ +/* + * Copyright (C) 2026 github.com/MaloneTalk + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as + * published by the Free Software Foundation, either version 3 of the + * License, or any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + * limitations under the License. + */ +package io.github.malonetalk.service.semantic.sync; + +import static org.mockito.ArgumentMatchers.anyList; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import io.github.malonetalk.entity.TableInfo; +import io.github.malonetalk.mapper.ColumnSemanticInfoMapper; +import io.github.malonetalk.mapper.TableInfoMapper; +import io.github.malonetalk.service.semantic.sync.SemanticSyncApplyService.ColumnSyncSource; +import io.github.malonetalk.service.semantic.sync.SemanticSyncApplyService.TableSyncSource; +import java.util.List; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +@ExtendWith(MockitoExtension.class) +class SemanticSyncApplyServiceTest { + + @Mock private TableInfoMapper tableInfoMapper; + @Mock private ColumnSemanticInfoMapper columnSemanticInfoMapper; + @InjectMocks private SemanticSyncApplyService service; + + @Test + void associatesNewColumnWithPersistedTableId() { + TableInfo persistedTable = new TableInfo(); + persistedTable.setId(42); + persistedTable.setDatasourceId(7); + persistedTable.setTableName("orders"); + + when(tableInfoMapper.selectByDatasourceIdAndTableNames(eq(7), anyList())) + .thenReturn(List.of(), List.of(persistedTable)); + when(columnSemanticInfoMapper.selectByDatasourceIdAndTableNames(eq(7), anyList())) + .thenReturn(List.of()); + + service.applyTableSync( + 7, + List.of( + new TableSyncSource( + "orders", + "orders table", + List.of( + new ColumnSyncSource( + "id", "order id", "BIGINT", true, "PRIMARY")))), + List.of()); + + verify(columnSemanticInfoMapper) + .batchUpsertPhysicalCache( + org.mockito.ArgumentMatchers.argThat( + columns -> + columns.size() == 1 + && Integer.valueOf(42) + .equals(columns.get(0).getTableId()))); + } +} diff --git a/sql/data_source.sql b/sql/data_source.sql index 63e4e9eb..1a2cbfa0 100644 --- a/sql/data_source.sql +++ b/sql/data_source.sql @@ -40,7 +40,7 @@ CREATE TABLE IF NOT EXISTS `table_info` ( CREATE TABLE IF NOT EXISTS `column_info` ( `id` INT NOT NULL AUTO_INCREMENT COMMENT '主键ID', `datasource_id` INT NOT NULL COMMENT '关联数据源ID', - `table_name` VARCHAR(255) NOT NULL COMMENT '表名', + `table_id` INT NOT NULL COMMENT '关联表信息ID', `column_name` VARCHAR(255) NOT NULL COMMENT '列名', `physical_column_description` VARCHAR(500) DEFAULT NULL COMMENT '物理列原始描述', `type_name` VARCHAR(255) DEFAULT NULL COMMENT '物理列类型', @@ -53,19 +53,21 @@ CREATE TABLE IF NOT EXISTS `column_info` ( `update_time` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间', `index_info` TEXT COMMENT 'physical index info', PRIMARY KEY (`id`), - UNIQUE KEY `uk_datasource_table_column` (`datasource_id`, `table_name`, `column_name`), - KEY `idx_datasource_table_visible` (`datasource_id`, `table_name`, `is_visible`), - KEY `idx_datasource_table_visible_column` - (`datasource_id`, `table_name`, `is_visible`, `column_name`) + UNIQUE KEY `uk_table_column` (`table_id`, `column_name`), + KEY `idx_datasource_id` (`datasource_id`), + KEY `idx_table_visible` (`table_id`, `is_visible`), + KEY `idx_table_visible_column` (`table_id`, `is_visible`, `column_name`), + CONSTRAINT `fk_column_info_table` + FOREIGN KEY (`table_id`) REFERENCES `table_info` (`id`) ON DELETE RESTRICT ON UPDATE RESTRICT ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='列信息表'; CREATE TABLE IF NOT EXISTS `logical_table_relation` ( `id` INT NOT NULL AUTO_INCREMENT COMMENT '主键ID', `datasource_id` INT NOT NULL COMMENT '关联数据源ID', - `source_table_name` VARCHAR(255) NOT NULL COMMENT '源表名', + `source_table_id` INT NOT NULL COMMENT '源表信息ID', `source_column_names_json` TEXT NOT NULL COMMENT '源列名JSON', `source_column_signature` VARCHAR(500) NOT NULL COMMENT '源列签名', - `target_table_name` VARCHAR(255) NOT NULL COMMENT '目标表名', + `target_table_id` INT NOT NULL COMMENT '目标表信息ID', `target_column_names_json` TEXT NOT NULL COMMENT '目标列名JSON', `target_column_signature` VARCHAR(500) NOT NULL COMMENT '目标列签名', `relation_type` VARCHAR(64) NOT NULL COMMENT '关系类型', @@ -75,12 +77,16 @@ CREATE TABLE IF NOT EXISTS `logical_table_relation` ( `update_time` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间', PRIMARY KEY (`id`), KEY `idx_relation_source_signature` - (`datasource_id`, `source_table_name`, `source_column_signature`), - KEY `idx_relation_source_table` (`datasource_id`, `source_table_name`), - KEY `idx_relation_source_enabled` (`datasource_id`, `source_table_name`, `is_enabled`), - KEY `idx_relation_source_enabled_id` (`datasource_id`, `source_table_name`, `is_enabled`, `id`), + (`datasource_id`, `source_table_id`, `source_column_signature`), + KEY `idx_relation_source_table` (`datasource_id`, `source_table_id`), + KEY `idx_relation_source_enabled` (`datasource_id`, `source_table_id`, `is_enabled`), + KEY `idx_relation_source_enabled_id` (`datasource_id`, `source_table_id`, `is_enabled`, `id`), KEY `idx_relation_source_target_id` - (`datasource_id`, `source_table_name`, `target_table_name`, `id`) + (`datasource_id`, `source_table_id`, `target_table_id`, `id`), + CONSTRAINT `fk_relation_source_table` + FOREIGN KEY (`source_table_id`) REFERENCES `table_info` (`id`) ON DELETE RESTRICT ON UPDATE RESTRICT, + CONSTRAINT `fk_relation_target_table` + FOREIGN KEY (`target_table_id`) REFERENCES `table_info` (`id`) ON DELETE RESTRICT ON UPDATE RESTRICT ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='逻辑表关系表'; CREATE TABLE IF NOT EXISTS `domain_info` ( diff --git a/sql/migration_primary_key_relations.sql b/sql/migration_primary_key_relations.sql new file mode 100644 index 00000000..fabec1fb --- /dev/null +++ b/sql/migration_primary_key_relations.sql @@ -0,0 +1,67 @@ +-- One-time MySQL 5.7 migration for primary-key-based semantic table references. +-- Back up the database first. The NOT NULL steps intentionally fail when legacy names +-- cannot be resolved, so inconsistent metadata is not silently discarded. + +ALTER TABLE `column_info` + ADD COLUMN `table_id` INT NULL COMMENT '关联表信息ID' AFTER `datasource_id`; + +UPDATE `column_info` c +INNER JOIN `table_info` t + ON t.`datasource_id` = c.`datasource_id` + AND LOWER(t.`table_name`) = LOWER(c.`table_name`) +SET c.`table_id` = t.`id`; + +ALTER TABLE `column_info` + MODIFY COLUMN `table_id` INT NOT NULL COMMENT '关联表信息ID', + DROP INDEX `uk_datasource_table_column`, + DROP INDEX `idx_datasource_table_visible`, + DROP INDEX `idx_datasource_table_visible_column`, + ADD UNIQUE KEY `uk_table_column` (`table_id`, `column_name`), + ADD KEY `idx_datasource_id` (`datasource_id`), + ADD KEY `idx_table_visible` (`table_id`, `is_visible`), + ADD KEY `idx_table_visible_column` (`table_id`, `is_visible`, `column_name`), + ADD CONSTRAINT `fk_column_info_table` + FOREIGN KEY (`table_id`) REFERENCES `table_info` (`id`) + ON DELETE RESTRICT ON UPDATE RESTRICT, + DROP COLUMN `table_name`; + +ALTER TABLE `logical_table_relation` + ADD COLUMN `source_table_id` INT NULL COMMENT '源表信息ID' AFTER `datasource_id`, + ADD COLUMN `target_table_id` INT NULL COMMENT '目标表信息ID' + AFTER `source_column_signature`; + +UPDATE `logical_table_relation` relation_meta +INNER JOIN `table_info` source_table + ON source_table.`datasource_id` = relation_meta.`datasource_id` + AND LOWER(source_table.`table_name`) = LOWER(relation_meta.`source_table_name`) +INNER JOIN `table_info` target_table + ON target_table.`datasource_id` = relation_meta.`datasource_id` + AND LOWER(target_table.`table_name`) = LOWER(relation_meta.`target_table_name`) +SET relation_meta.`source_table_id` = source_table.`id`, + relation_meta.`target_table_id` = target_table.`id`; + +ALTER TABLE `logical_table_relation` + MODIFY COLUMN `source_table_id` INT NOT NULL COMMENT '源表信息ID', + MODIFY COLUMN `target_table_id` INT NOT NULL COMMENT '目标表信息ID', + DROP INDEX `idx_relation_source_signature`, + DROP INDEX `idx_relation_source_table`, + DROP INDEX `idx_relation_source_enabled`, + DROP INDEX `idx_relation_source_enabled_id`, + DROP INDEX `idx_relation_source_target_id`, + ADD KEY `idx_relation_source_signature` + (`datasource_id`, `source_table_id`, `source_column_signature`), + ADD KEY `idx_relation_source_table` (`datasource_id`, `source_table_id`), + ADD KEY `idx_relation_source_enabled` + (`datasource_id`, `source_table_id`, `is_enabled`), + ADD KEY `idx_relation_source_enabled_id` + (`datasource_id`, `source_table_id`, `is_enabled`, `id`), + ADD KEY `idx_relation_source_target_id` + (`datasource_id`, `source_table_id`, `target_table_id`, `id`), + ADD CONSTRAINT `fk_relation_source_table` + FOREIGN KEY (`source_table_id`) REFERENCES `table_info` (`id`) + ON DELETE RESTRICT ON UPDATE RESTRICT, + ADD CONSTRAINT `fk_relation_target_table` + FOREIGN KEY (`target_table_id`) REFERENCES `table_info` (`id`) + ON DELETE RESTRICT ON UPDATE RESTRICT, + DROP COLUMN `source_table_name`, + DROP COLUMN `target_table_name`; From 9c571ec4cc3baa67d95361374452cfe91fba071b Mon Sep 17 00:00:00 2001 From: mengnankkkk Date: Thu, 24 Sep 2026 20:50:10 +0800 Subject: [PATCH 2/7] fix --- README.md | 8 ++++- .../mapper/ColumnSemanticInfoMapper.java | 3 ++ .../malonetalk/mapper/TableInfoMapper.java | 3 ++ .../column/ColumnSemanticServiceImpl.java | 33 +++++++++++++++---- .../table/TableSemanticServiceImpl.java | 33 +++++++++++++++---- .../mapper/ColumnSemanticInfoMapper.xml | 30 ++++++++--------- .../main/resources/mapper/TableInfoMapper.xml | 11 +++++++ docs/configuration.md | 3 +- docs/getting-started.md | 9 +++++ sql/data_source.sql | 6 ++-- sql/migration_primary_key_relations.sql | 6 ++-- 11 files changed, 106 insertions(+), 39 deletions(-) diff --git a/README.md b/README.md index 89c347c8..f203e7dc 100644 --- a/README.md +++ b/README.md @@ -91,6 +91,12 @@ cd data-agent-frontend pnpm install && pnpm dev ``` +已有数据库升级到主键关联版本时,先备份数据库并在启动新版后端前执行: + +```bash +mysql -u root -p data_agent < sql/migration_primary_key_relations.sql +``` + 浏览器打开 http://localhost:3000 ,使用 `admin` / `ADMIN_INIT_PASSWORD` 登录后,先在「数据源管理」接入业务库,在「语义管理」同步表结构并维护表/列/指标口径,再到聊天框用自然语言提问,例如:*"上个月各区域销售额是多少?"* 用户与角色入口在「系统管理」。 ## 📚 文档 @@ -151,4 +157,4 @@ pnpm install && pnpm dev ---- \ No newline at end of file +--- diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java index 53dd4c2e..aac95d58 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java @@ -57,4 +57,7 @@ ColumnInfo selectByDatasourceIdAndTableNameAndColumnName( int resetSemanticFieldsByIds( @Param("datasourceId") Integer datasourceId, @Param("ids") List ids); + + int deletePhysicalMissingByIds( + @Param("datasourceId") Integer datasourceId, @Param("ids") List ids); } diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/TableInfoMapper.java b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/TableInfoMapper.java index a8952855..2d0a481e 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/TableInfoMapper.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/TableInfoMapper.java @@ -36,6 +36,9 @@ public interface TableInfoMapper { int resetSemanticFieldsByIds( @Param("datasourceId") Integer datasourceId, @Param("ids") List ids); + int deletePhysicalMissingByIds( + @Param("datasourceId") Integer datasourceId, @Param("ids") List ids); + List selectByDatasourceId(@Param("datasourceId") Integer datasourceId); List selectPageByDatasourceId( diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java index 6bf1f067..26660fae 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java @@ -36,11 +36,13 @@ import io.github.malonetalk.service.semantic.SemanticMergeService; import io.github.malonetalk.utils.SemanticUtils; import java.time.LocalDateTime; +import java.util.ArrayList; import java.util.List; import java.util.Set; import java.util.stream.Collectors; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; @Service @RequiredArgsConstructor @@ -132,6 +134,7 @@ public List getMergedTableSchema(Integer datasourceId, Str } @Override + @Transactional public void resetColumnSemantic(Integer datasourceId, String tableName, String columnName) { requireDatasource(datasourceId); ColumnInfo existing = @@ -141,10 +144,11 @@ public void resetColumnSemantic(Integer datasourceId, String tableName, String c throw BusinessException.of( ErrorCode.RESOURCE_NOT_FOUND, "Column semantic metadata does not exist."); } - columnSemanticInfoMapper.resetSemanticFieldsByIds(datasourceId, List.of(existing.getId())); + resetColumnRecords(datasourceId, List.of(existing)); } @Override + @Transactional public int resetColumnSemantics( Integer datasourceId, String tableName, List columnNames) { requireDatasource(datasourceId); @@ -163,7 +167,7 @@ public int resetColumnSemantics( "Missing columnName for batch column semantic" + " reset.")) .collect(Collectors.toCollection(java.util.LinkedHashSet::new)); - List matchedIds = + List matchedColumns = columnSemanticInfoMapper .selectByDatasourceIdAndTableName(datasourceId, normalizedTableName) .stream() @@ -174,20 +178,35 @@ public int resetColumnSemantics( column.getColumnName(), "Missing columnName while matching column" + " semantic reset."))) - .map(ColumnInfo::getId) - .distinct() .toList(); - if (matchedIds.isEmpty()) { + if (matchedColumns.isEmpty()) { return 0; } - if (matchedIds.size() != normalizedColumnNames.size()) { + if (matchedColumns.size() != normalizedColumnNames.size()) { throw BusinessException.of( ErrorCode.RESOURCE_NOT_FOUND, "Some column semantic metadata does not exist for table " + normalizedTableName + "."); } - return columnSemanticInfoMapper.resetSemanticFieldsByIds(datasourceId, matchedIds); + return resetColumnRecords(datasourceId, matchedColumns); + } + + private int resetColumnRecords(Integer datasourceId, List columns) { + List resetIds = new ArrayList<>(); + List deleteIds = new ArrayList<>(); + for (ColumnInfo column : columns) { + (Boolean.FALSE.equals(column.getPhysicalStatus()) ? deleteIds : resetIds) + .add(column.getId()); + } + int affected = 0; + if (!deleteIds.isEmpty()) { + affected += columnSemanticInfoMapper.deletePhysicalMissingByIds(datasourceId, deleteIds); + } + if (!resetIds.isEmpty()) { + affected += columnSemanticInfoMapper.resetSemanticFieldsByIds(datasourceId, resetIds); + } + return affected; } private void requireDatasource(Integer datasourceId) { diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java index 421d3fcf..f50eb83e 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java @@ -34,11 +34,13 @@ import io.github.malonetalk.service.semantic.SemanticMergeService; import io.github.malonetalk.utils.SemanticUtils; import java.time.LocalDateTime; +import java.util.ArrayList; import java.util.List; import java.util.Set; import java.util.stream.Collectors; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; @Service @RequiredArgsConstructor @@ -171,6 +173,7 @@ public void updateTableSemantic(TableSemanticUpdateRequest request) { } @Override + @Transactional public void resetTableSemantic(Integer datasourceId, String tableName) { requireDatasource(datasourceId); String normalizedTableName = @@ -182,10 +185,11 @@ public void resetTableSemantic(Integer datasourceId, String tableName) { throw BusinessException.of( ErrorCode.RESOURCE_NOT_FOUND, "Table semantic metadata does not exist."); } - tableInfoMapper.resetSemanticFieldsByIds(datasourceId, List.of(existing.getId())); + resetTableRecords(datasourceId, List.of(existing)); } @Override + @Transactional public int resetTableSemantics(Integer datasourceId, List tableNames) { requireDatasource(datasourceId); if (tableNames == null || tableNames.isEmpty()) { @@ -200,7 +204,7 @@ public int resetTableSemantics(Integer datasourceId, List tableNames) { "Missing tableName for batch table semantic" + " reset.")) .collect(Collectors.toCollection(java.util.LinkedHashSet::new)); - List matchedIds = + List matchedTables = tableInfoMapper.selectByDatasourceId(datasourceId).stream() .filter( table -> @@ -209,20 +213,35 @@ public int resetTableSemantics(Integer datasourceId, List tableNames) { table.getTableName(), "Missing tableName while matching table" + " semantic reset."))) - .map(TableInfo::getId) - .distinct() .toList(); - if (matchedIds.isEmpty()) { + if (matchedTables.isEmpty()) { return 0; } - if (matchedIds.size() != normalizedNames.size()) { + if (matchedTables.size() != normalizedNames.size()) { throw BusinessException.of( ErrorCode.RESOURCE_NOT_FOUND, "Some table semantic metadata does not exist for datasource " + datasourceId + "."); } - return tableInfoMapper.resetSemanticFieldsByIds(datasourceId, matchedIds); + return resetTableRecords(datasourceId, matchedTables); + } + + private int resetTableRecords(Integer datasourceId, List tables) { + List resetIds = new ArrayList<>(); + List deleteIds = new ArrayList<>(); + for (TableInfo table : tables) { + (Boolean.FALSE.equals(table.getPhysicalStatus()) ? deleteIds : resetIds) + .add(table.getId()); + } + int affected = 0; + if (!deleteIds.isEmpty()) { + affected += tableInfoMapper.deletePhysicalMissingByIds(datasourceId, deleteIds); + } + if (!resetIds.isEmpty()) { + affected += tableInfoMapper.resetSemanticFieldsByIds(datasourceId, resetIds); + } + return affected; } private void requireDatasource(Integer datasourceId) { diff --git a/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml b/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml index 37224fed..dc73d185 100644 --- a/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml @@ -24,23 +24,6 @@ SELECT c.*, table_meta.table_name FROM column_info c INNER JOIN table_info table_meta ON table_meta.id = c.table_id - WHERE c.id = ( - SELECT c2.id - FROM column_info c2 - WHERE c2.table_id = c.table_id - AND LOWER(c2.column_name) = LOWER(c.column_name) - ORDER BY - CASE - WHEN c2.column_description IS NOT NULL AND c2.column_description <> '' THEN 0 - WHEN c2.semantic_type IS NOT NULL AND c2.semantic_type <> '' THEN 0 - WHEN c2.is_visible = 0 THEN 0 - ELSE 1 - END ASC, - c2.update_time DESC, - c2.create_time DESC, - c2.id DESC - LIMIT 1 - ) - SELECT * FROM ( - - ) column_info - WHERE datasource_id = #{datasourceId} - ORDER BY table_name ASC, column_name ASC, id ASC + + WHERE c.datasource_id = #{datasourceId} + ORDER BY table_meta.table_name ASC, c.column_name ASC, c.id ASC diff --git a/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml b/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml index 12c466da..d3c181dc 100644 --- a/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml @@ -8,11 +8,9 @@ - - @@ -78,12 +76,12 @@ INSERT INTO logical_table_relation ( - datasource_id, source_table_id, source_column_names_json, source_column_signature, - target_table_id, target_column_names_json, target_column_signature, relation_type, + datasource_id, source_table_id, source_column_names_json, + target_table_id, target_column_names_json, relation_type, description, is_enabled, create_time, update_time ) VALUES ( - #{datasourceId}, #{sourceTableId}, #{sourceColumnNamesJson}, #{sourceColumnSignature}, - #{targetTableId}, #{targetColumnNamesJson}, #{targetColumnSignature}, #{relationType}, + #{datasourceId}, #{sourceTableId}, #{sourceColumnNamesJson}, + #{targetTableId}, #{targetColumnNamesJson}, #{relationType}, #{description}, #{isEnabled}, #{createTime}, #{updateTime} ) @@ -93,10 +91,8 @@ source_table_id = #{sourceTableId}, source_column_names_json = #{sourceColumnNamesJson}, - source_column_signature = #{sourceColumnSignature}, target_table_id = #{targetTableId}, target_column_names_json = #{targetColumnNamesJson}, - target_column_signature = #{targetColumnSignature}, relation_type = #{relationType}, description = #{description}, is_enabled = #{isEnabled}, diff --git a/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml b/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml index 75416da9..c769b762 100644 --- a/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml @@ -16,39 +16,14 @@ - - SELECT t.* - FROM table_info t - WHERE t.id = ( - SELECT t2.id - FROM table_info t2 - WHERE t2.datasource_id = t.datasource_id - AND LOWER(t2.table_name) = LOWER(t.table_name) - ORDER BY - CASE - WHEN t2.table_description IS NOT NULL AND t2.table_description <> '' THEN 0 - WHEN t2.is_visible = 0 THEN 0 - ELSE 1 - END ASC, - t2.update_time DESC, - t2.create_time DESC, - t2.id DESC - LIMIT 1 - ) - - - SELECT COUNT(1) FROM ( - - ) table_info + SELECT COUNT(1) FROM table_info WHERE LOWER(domain) = LOWER(#{domain}) diff --git a/sql/data_source.sql b/sql/data_source.sql index 6191ddc8..00848675 100644 --- a/sql/data_source.sql +++ b/sql/data_source.sql @@ -55,8 +55,6 @@ CREATE TABLE IF NOT EXISTS `column_info` ( PRIMARY KEY (`id`), UNIQUE KEY `uk_table_column` (`table_id`, `column_name`), KEY `idx_datasource_id` (`datasource_id`), - KEY `idx_table_visible` (`table_id`, `is_visible`), - KEY `idx_table_visible_column` (`table_id`, `is_visible`, `column_name`), CONSTRAINT `fk_column_info_table` FOREIGN KEY (`table_id`) REFERENCES `table_info` (`id`) ON DELETE CASCADE ON UPDATE RESTRICT ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='列信息表'; @@ -66,23 +64,17 @@ CREATE TABLE IF NOT EXISTS `logical_table_relation` ( `datasource_id` INT NOT NULL COMMENT '关联数据源ID', `source_table_id` INT NOT NULL COMMENT '源表信息ID', `source_column_names_json` TEXT NOT NULL COMMENT '源列名JSON', - `source_column_signature` VARCHAR(500) NOT NULL COMMENT '源列签名', `target_table_id` INT NOT NULL COMMENT '目标表信息ID', `target_column_names_json` TEXT NOT NULL COMMENT '目标列名JSON', - `target_column_signature` VARCHAR(500) NOT NULL COMMENT '目标列签名', `relation_type` VARCHAR(64) NOT NULL COMMENT '关系类型', `description` VARCHAR(1000) DEFAULT NULL COMMENT '关系描述', `is_enabled` TINYINT(1) DEFAULT 1 COMMENT '是否启用', `create_time` DATETIME DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间', `update_time` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间', PRIMARY KEY (`id`), - KEY `idx_relation_source_signature` - (`datasource_id`, `source_table_id`, `source_column_signature`), - KEY `idx_relation_source_table` (`datasource_id`, `source_table_id`), - KEY `idx_relation_source_enabled` (`datasource_id`, `source_table_id`, `is_enabled`), + KEY `idx_relation_source_table` (`source_table_id`), + KEY `idx_relation_target_table` (`target_table_id`), KEY `idx_relation_source_enabled_id` (`datasource_id`, `source_table_id`, `is_enabled`, `id`), - KEY `idx_relation_source_target_id` - (`datasource_id`, `source_table_id`, `target_table_id`, `id`), CONSTRAINT `fk_relation_source_table` FOREIGN KEY (`source_table_id`) REFERENCES `table_info` (`id`) ON DELETE CASCADE ON UPDATE RESTRICT, CONSTRAINT `fk_relation_target_table` diff --git a/sql/migration_primary_key_relations.sql b/sql/migration_primary_key_relations.sql index 6828fa80..1daa0598 100644 --- a/sql/migration_primary_key_relations.sql +++ b/sql/migration_primary_key_relations.sql @@ -53,8 +53,6 @@ ALTER TABLE `column_info` DROP INDEX `idx_datasource_table_visible_column`, ADD UNIQUE KEY `uk_table_column` (`table_id`, `column_name`), ADD KEY `idx_datasource_id` (`datasource_id`), - ADD KEY `idx_table_visible` (`table_id`, `is_visible`), - ADD KEY `idx_table_visible_column` (`table_id`, `is_visible`, `column_name`), ADD CONSTRAINT `fk_column_info_table` FOREIGN KEY (`table_id`) REFERENCES `table_info` (`id`) ON DELETE CASCADE ON UPDATE RESTRICT, @@ -63,7 +61,7 @@ ALTER TABLE `column_info` ALTER TABLE `logical_table_relation` ADD COLUMN `source_table_id` INT NULL COMMENT '源表信息ID' AFTER `datasource_id`, ADD COLUMN `target_table_id` INT NULL COMMENT '目标表信息ID' - AFTER `source_column_signature`; + AFTER `source_column_names_json`; UPDATE `logical_table_relation` relation_meta INNER JOIN `table_info` source_table @@ -83,15 +81,10 @@ ALTER TABLE `logical_table_relation` DROP INDEX `idx_relation_source_enabled`, DROP INDEX `idx_relation_source_enabled_id`, DROP INDEX `idx_relation_source_target_id`, - ADD KEY `idx_relation_source_signature` - (`datasource_id`, `source_table_id`, `source_column_signature`), - ADD KEY `idx_relation_source_table` (`datasource_id`, `source_table_id`), - ADD KEY `idx_relation_source_enabled` - (`datasource_id`, `source_table_id`, `is_enabled`), + ADD KEY `idx_relation_source_table` (`source_table_id`), + ADD KEY `idx_relation_target_table` (`target_table_id`), ADD KEY `idx_relation_source_enabled_id` (`datasource_id`, `source_table_id`, `is_enabled`, `id`), - ADD KEY `idx_relation_source_target_id` - (`datasource_id`, `source_table_id`, `target_table_id`, `id`), ADD CONSTRAINT `fk_relation_source_table` FOREIGN KEY (`source_table_id`) REFERENCES `table_info` (`id`) ON DELETE CASCADE ON UPDATE RESTRICT, @@ -99,4 +92,6 @@ ALTER TABLE `logical_table_relation` FOREIGN KEY (`target_table_id`) REFERENCES `table_info` (`id`) ON DELETE CASCADE ON UPDATE RESTRICT, DROP COLUMN `source_table_name`, - DROP COLUMN `target_table_name`; + DROP COLUMN `target_table_name`, + DROP COLUMN `source_column_signature`, + DROP COLUMN `target_column_signature`; From e2c25c1a9fcc01a815bf3510dbe88ac4363ed670 Mon Sep 17 00:00:00 2001 From: mengnankkkk Date: Sat, 26 Sep 2026 11:53:54 +0800 Subject: [PATCH 5/7] fix --- .../semantic/SemanticMergeService.java | 87 +++++++------------ .../column/ColumnSemanticServiceImpl.java | 4 +- .../relation/LogicalTableRelationHelper.java | 37 +++----- .../relation/RelationSemanticServiceImpl.java | 75 +++++----------- .../sync/SemanticSyncApplyService.java | 4 +- .../sync/SemanticSyncServiceImpl.java | 9 +- .../table/TableSemanticServiceImpl.java | 12 +-- 7 files changed, 76 insertions(+), 152 deletions(-) diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/SemanticMergeService.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/SemanticMergeService.java index f78951f0..bace13f9 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/SemanticMergeService.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/SemanticMergeService.java @@ -58,8 +58,9 @@ public class SemanticMergeService { public List listVisibleTablesByDomains( Datasource datasource, List domains) { List normalizedDomains = normalizeDomains(domains); - TableNameIndex tableIndex = - TableNameIndex.of(tableInfoMapper.selectByDatasourceId(datasource.getId())); + // 关系和列按 table_id 对齐;表名只用于最终提示词和关系去重。 + TableIdIndex tableIndex = + TableIdIndex.of(tableInfoMapper.selectByDatasourceId(datasource.getId())); TableColumnIndex columnIndex = TableColumnIndex.of( columnSemanticInfoMapper.selectByDatasourceId(datasource.getId())); @@ -78,7 +79,7 @@ public List listVisibleTablesByDomains( PromptConverter.mapTablePrompt( table, resolveVisibleRelations( - table.getTableName(), + table.getId(), tableIndex, columnIndex, relationIndex))) @@ -122,27 +123,27 @@ public List getTableSchema(Datasource datasource, String t } private List resolveVisibleRelations( - String sourceTableName, - TableNameIndex tableIndex, + Integer sourceTableId, + TableIdIndex tableIndex, TableColumnIndex columnIndex, RelationSourceIndex relationIndex) { List visibleRelations = filterVisibleLogicalRelations( - relationIndex.get(sourceTableName), tableIndex, columnIndex); + relationIndex.get(sourceTableId), tableIndex, columnIndex); return deduplicateRelations(visibleRelations); } private List filterVisibleLogicalRelations( List logicalRelations, - TableNameIndex tableIndex, + TableIdIndex tableIndex, TableColumnIndex columnIndex) { List visibleRelations = new ArrayList<>(); for (LogicalTableRelation relation : logicalRelations) { if (!Boolean.TRUE.equals(relation.getIsEnabled())) { continue; } - if (tableIndex.isUnavailable(relation.getSourceTableName()) - || tableIndex.isUnavailable(relation.getTargetTableName())) { + if (tableIndex.isUnavailable(relation.getSourceTableId()) + || tableIndex.isUnavailable(relation.getTargetTableId())) { continue; } List sourceColumns = parseRelationColumns(relation, true); @@ -150,9 +151,9 @@ private List filterVisibleLogicalRelations( if (sourceColumns == null || targetColumns == null) { continue; } - if (columnIndex.hasUnavailableColumn(relation.getSourceTableName(), sourceColumns) + if (columnIndex.hasUnavailableColumn(relation.getSourceTableId(), sourceColumns) || columnIndex.hasUnavailableColumn( - relation.getTargetTableName(), targetColumns)) { + relation.getTargetTableId(), targetColumns)) { continue; } visibleRelations.add( @@ -233,46 +234,32 @@ private String targetTableName() { } } - private record TableNameIndex(Map index) { + private record TableIdIndex(Map index) { - private static TableNameIndex of(List tables) { - Map map = new LinkedHashMap<>(); + private static TableIdIndex of(List tables) { + Map map = new LinkedHashMap<>(); for (TableInfo table : tables) { - map.put( - SemanticUtils.normalizeObjectName( - table.getTableName(), - "Missing tableName while building semantic table index."), - table); + map.put(table.getId(), table); } - return new TableNameIndex(map); + return new TableIdIndex(map); } private List asList() { return List.copyOf(index.values()); } - private TableInfo get(String tableName) { - return index.get( - SemanticUtils.normalizeObjectName( - tableName, "Missing tableName while reading semantic table index.")); - } - - private boolean isUnavailable(String tableName) { - TableInfo tableInfo = get(tableName); + private boolean isUnavailable(Integer tableId) { + TableInfo tableInfo = index.get(tableId); return tableInfo == null || !SemanticAvailabilityHelper.isTableAvailable(tableInfo); } } - private record TableColumnIndex(Map> index) { + private record TableColumnIndex(Map> index) { private static TableColumnIndex of(List columns) { - Map> map = new HashMap<>(); + Map> map = new HashMap<>(); for (ColumnInfo column : columns) { - map.computeIfAbsent( - SemanticUtils.normalizeObjectName( - column.getTableName(), - "Missing tableName while building semantic column index."), - key -> new HashMap<>()) + map.computeIfAbsent(column.getTableId(), key -> new HashMap<>()) .put( SemanticUtils.normalizeObjectName( column.getColumnName(), @@ -282,12 +269,8 @@ private static TableColumnIndex of(List columns) { return new TableColumnIndex(map); } - private ColumnInfo get(String tableName, String columnName) { - Map columns = - index.get( - SemanticUtils.normalizeObjectName( - tableName, - "Missing tableName while reading semantic column index.")); + private ColumnInfo get(Integer tableId, String columnName) { + Map columns = index.get(tableId); return columns == null ? null : columns.get( @@ -296,9 +279,9 @@ private ColumnInfo get(String tableName, String columnName) { "Missing columnName while reading semantic column index.")); } - private boolean hasUnavailableColumn(String tableName, List columnNames) { + private boolean hasUnavailableColumn(Integer tableId, List columnNames) { for (String columnName : columnNames) { - ColumnInfo columnInfo = get(tableName, columnName); + ColumnInfo columnInfo = get(tableId, columnName); if (columnInfo == null || !SemanticAvailabilityHelper.isColumnAvailable(columnInfo)) { return true; @@ -308,27 +291,19 @@ private boolean hasUnavailableColumn(String tableName, List columnNames) } } - private record RelationSourceIndex(Map> index) { + private record RelationSourceIndex(Map> index) { private static RelationSourceIndex of(List relations) { - Map> map = new HashMap<>(); + Map> map = new HashMap<>(); for (LogicalTableRelation relation : relations) { - map.computeIfAbsent( - SemanticUtils.normalizeObjectName( - relation.getSourceTableName(), - "Missing sourceTableName while building relation index."), - key -> new ArrayList<>()) + map.computeIfAbsent(relation.getSourceTableId(), key -> new ArrayList<>()) .add(relation); } return new RelationSourceIndex(map); } - private List get(String sourceTableName) { - return index.getOrDefault( - SemanticUtils.normalizeObjectName( - sourceTableName, - "Missing sourceTableName while reading relation index."), - Collections.emptyList()); + private List get(Integer sourceTableId) { + return index.getOrDefault(sourceTableId, Collections.emptyList()); } } } diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java index f1174b93..2e06ed48 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/column/ColumnSemanticServiceImpl.java @@ -193,13 +193,15 @@ public int resetColumnSemantics( private int resetColumnRecords(Integer datasourceId, List columns) { List resetIds = new ArrayList<>(); List deleteIds = new ArrayList<>(); + // 仅删除明确标记为物理缺失的列;其他记录保留并重置语义字段。 for (ColumnInfo column : columns) { (Boolean.FALSE.equals(column.getPhysicalStatus()) ? deleteIds : resetIds) .add(column.getId()); } int affected = 0; if (!deleteIds.isEmpty()) { - affected += columnSemanticInfoMapper.deletePhysicalMissingByIds(datasourceId, deleteIds); + affected += + columnSemanticInfoMapper.deletePhysicalMissingByIds(datasourceId, deleteIds); } if (!resetIds.isEmpty()) { affected += columnSemanticInfoMapper.resetSemanticFieldsByIds(datasourceId, resetIds); diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/LogicalTableRelationHelper.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/LogicalTableRelationHelper.java index 4d916ebf..9a367572 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/LogicalTableRelationHelper.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/LogicalTableRelationHelper.java @@ -28,9 +28,12 @@ import io.github.malonetalk.exception.ErrorCode; import io.github.malonetalk.utils.RequestAssert; import io.github.malonetalk.utils.SemanticUtils; -import java.util.LinkedHashSet; +import java.util.ArrayList; +import java.util.HashSet; import java.util.List; +import java.util.Locale; import java.util.Set; +import java.util.stream.Collectors; import org.springframework.stereotype.Component; @Component @@ -44,22 +47,16 @@ public LogicalTableRelationHelper(ObjectMapper objectMapper) { this.objectMapper = objectMapper; } - public String normalizeTableName(String tableName, String missingMessage) { - return SemanticUtils.normalizeObjectName(tableName, missingMessage); - } - public List normalizeColumnNames(List columnNames, String fieldName) { RequestAssert.requireNotEmpty(columnNames, fieldName + " cannot be empty."); - Set uniqueKeys = new LinkedHashSet<>(); - Set normalizedColumns = new LinkedHashSet<>(); + Set uniqueKeys = new HashSet<>(); + List normalizedColumns = new ArrayList<>(columnNames.size()); + // 按不区分大小写的名称判重,序列化时保留列名原有大小写。 for (String columnName : columnNames) { String normalizedColumnName = RequestAssert.requireNotBlank( columnName, fieldName + " contains a blank column name."); - String uniqueKey = - SemanticUtils.normalizeObjectName( - normalizedColumnName, - "Missing columnName while normalizing logical relation columns."); + String uniqueKey = normalizedColumnName.toLowerCase(Locale.ROOT); if (!uniqueKeys.add(uniqueKey)) { throw BusinessException.of( ErrorCode.BAD_REQUEST, @@ -67,18 +64,13 @@ public List normalizeColumnNames(List columnNames, String fieldN } normalizedColumns.add(normalizedColumnName); } - return normalizedColumns.stream().toList(); + return List.copyOf(normalizedColumns); } - public String buildColumnSignature(List columnNames) { + private String buildColumnSignature(List columnNames) { return normalizeColumnNames(columnNames, "columnNames").stream() - .map( - columnName -> - SemanticUtils.normalizeObjectName( - columnName, - "Missing columnName while building column signature.")) - .reduce((left, right) -> left + RELATION_KEY_SEPARATOR + right) - .orElse(""); + .map(columnName -> columnName.toLowerCase(Locale.ROOT)) + .collect(Collectors.joining(RELATION_KEY_SEPARATOR)); } public String buildRelationKey( @@ -97,10 +89,9 @@ public String buildRelationKey( + buildColumnSignature(targetColumnNames); } - public String toJson(List columnNames) { + public String toJson(List columnNames, String fieldName) { try { - return objectMapper.writeValueAsString( - normalizeColumnNames(columnNames, "columnNames")); + return objectMapper.writeValueAsString(normalizeColumnNames(columnNames, fieldName)); } catch (JsonProcessingException e) { throw BusinessException.of( ErrorCode.OPERATION_FAILED, "Failed to serialize relation columns.", e); diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java index 5a3cf5c2..46b74fce 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java @@ -44,7 +44,6 @@ import io.github.malonetalk.utils.RequestAssert; import io.github.malonetalk.utils.SemanticUtils; import java.time.LocalDateTime; -import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Set; @@ -69,7 +68,7 @@ public PageResponse getRelationPage( RelationSemanticPageQuery query) { requireDatasource(query.datasourceId()); String normalizedTableName = - logicalTableRelationHelper.normalizeTableName(query.tableName(), "tableName"); + SemanticUtils.normalizeObjectName(query.tableName(), "tableName"); int pageNumber = PageResponse.resolvePage(query.page()); int pageSize = PageResponse.resolvePageSize(query.pageSize()); boolean sortDescending = SemanticUtils.isDescendingSort(query.sortOrder()); @@ -117,55 +116,32 @@ public RelationWorkspaceResponse getRelationWorkspace(RelationWorkspacePageQuery PageResponse.empty(pageNumber, pageSize), List.of()); } - List tableNames = page.stream().map(TableInfo::getTableName).distinct().toList(); - Set currentPageTableNames = - tableNames.stream() - .map( - tableName -> - SemanticUtils.normalizeObjectName( - tableName, - "Missing tableName while building relation" - + " workspace page.")) - .collect(Collectors.toSet()); - Map> columnsByTableName = + List tableNames = page.stream().map(TableInfo::getTableName).toList(); + Set currentPageTableIds = + page.stream().map(TableInfo::getId).collect(Collectors.toSet()); + // 列记录已持有 table_id,按主键分组可直接关联当前页的表。 + Map> columnsByTableId = columnSemanticInfoMapper .selectByDatasourceIdAndTableNames(query.datasourceId(), tableNames) .stream() - .collect( - Collectors.groupingBy( - column -> - SemanticUtils.normalizeObjectName( - column.getTableName(), - "Missing tableName while grouping relation" - + " workspace columns."), - LinkedHashMap::new, - Collectors.toList())); + .collect(Collectors.groupingBy(ColumnInfo::getTableId)); List nodes = page.stream() .map( table -> semanticConverter.toWorkspaceTable( table, - columnsByTableName.getOrDefault( - SemanticUtils.normalizeObjectName( - table.getTableName(), - "Missing tableName while mapping" - + " relation workspace" - + " table."), - List.of()))) + columnsByTableId.getOrDefault( + table.getId(), List.of()))) .toList(); + // 来源表来自当前页查询,目标表也必须属于当前页。 List relations = logicalTableRelationMapper .selectByDatasourceIdAndSourceTables(query.datasourceId(), tableNames) .stream() .filter( relation -> - currentPageTableNames.contains( - SemanticUtils.normalizeObjectName( - relation.getTargetTableName(), - "Missing targetTableName while filtering" - + " relation workspace" - + " relations."))) + currentPageTableIds.contains(relation.getTargetTableId())) .map(semanticConverter::toResponse) .toList(); @@ -232,8 +208,7 @@ public boolean deleteRelationSemantic( public int deleteRelationSemantics( Integer datasourceId, String tableName, List relationIds) { requireDatasource(datasourceId); - String normalizedTableName = - logicalTableRelationHelper.normalizeTableName(tableName, "tableName"); + String normalizedTableName = SemanticUtils.normalizeObjectName(tableName, "tableName"); if (relationIds == null || relationIds.isEmpty()) { return 0; } @@ -384,20 +359,13 @@ private void populateRelationFields( LogicalTableRelationType relationType, String description, Boolean enabled) { - relation.setSourceTableName( - logicalTableRelationHelper.normalizeTableName(tableName, "tableName")); - List normalizedSourceColumns = - logicalTableRelationHelper.normalizeColumnNames( - sourceColumnNames, "sourceColumnNames"); + relation.setSourceTableName(SemanticUtils.normalizeObjectName(tableName, "tableName")); relation.setSourceColumnNamesJson( - logicalTableRelationHelper.toJson(normalizedSourceColumns)); + logicalTableRelationHelper.toJson(sourceColumnNames, "sourceColumnNames")); relation.setTargetTableName( - logicalTableRelationHelper.normalizeTableName(targetTableName, "targetTableName")); - List normalizedTargetColumns = - logicalTableRelationHelper.normalizeColumnNames( - targetColumnNames, "targetColumnNames"); + SemanticUtils.normalizeObjectName(targetTableName, "targetTableName")); relation.setTargetColumnNamesJson( - logicalTableRelationHelper.toJson(normalizedTargetColumns)); + logicalTableRelationHelper.toJson(targetColumnNames, "targetColumnNames")); LogicalTableRelationType resolvedRelationType = relationType == null ? LogicalTableRelationType.FOREIGN_KEY : relationType; relation.setRelationType(resolvedRelationType.getCode()); @@ -408,15 +376,14 @@ private void populateRelationFields( private LogicalTableRelation requireRelation( Integer datasourceId, String tableName, Integer relationId) { RequestAssert.requireNonNull(relationId, "relationId cannot be null."); + String normalizedTableName = SemanticUtils.normalizeObjectName(tableName, "tableName"); LogicalTableRelation relation = logicalTableRelationMapper.selectById(relationId); - TableInfo sourceTable = - requireTable( - datasourceId, - logicalTableRelationHelper.normalizeTableName(tableName, "tableName"), - "sourceTable"); + // selectById 已关联 table_info,直接用当前来源表名验证归属。 if (relation == null || !datasourceId.equals(relation.getDatasourceId()) - || !sourceTable.getId().equals(relation.getSourceTableId())) { + || !normalizedTableName.equals( + SemanticUtils.normalizeObjectName( + relation.getSourceTableName(), "sourceTableName"))) { throw BusinessException.of( ErrorCode.RESOURCE_NOT_FOUND, "Logical relation does not exist."); } diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyService.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyService.java index a5763eb7..2c247d80 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyService.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncApplyService.java @@ -105,8 +105,7 @@ public List refreshPhysicalStatus( if (semanticTables.isEmpty()) { return List.of(); } - List tableNames = - semanticTables.stream().map(TableInfo::getTableName).distinct().toList(); + List tableNames = semanticTables.stream().map(TableInfo::getTableName).toList(); Map> columnsByTableName = loadSemanticColumnsByTable(datasourceId, tableNames); List tableNamesToMarkMissing = new ArrayList<>(); @@ -182,6 +181,7 @@ private void batchSavePresentSchema( tableInfoMapper.batchUpsertPhysicalCache(newTables); } + // 新表的主键由数据库生成,写入列之前需重新取得 table_id。 List presentTableNames = presentTables.stream().map(TableSyncSource::tableName).toList(); Map persistedTableIndex = diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncServiceImpl.java index ef358619..a36d46d8 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/sync/SemanticSyncServiceImpl.java @@ -110,13 +110,8 @@ public PageResponse getPhysicalTableCandidates( private List loadCandidateTables( Datasource datasource, List semanticTables) { - Map candidates = new LinkedHashMap<>(); - for (PhysicalTableInfo table : schemaReader.getTables(datasource)) { - candidates.putIfAbsent( - SemanticUtils.normalizeObjectName( - table.tableName(), "Missing physical tableName."), - table); - } + Map candidates = loadPhysicalTableIndex(datasource); + // 保留语义缓存中的表,供用户发现已从物理数据源消失的表。 for (TableInfo table : semanticTables) { String tableName = SemanticUtils.requireTrimmed( diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java index 9ea684c8..ba9e5843 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/table/TableSemanticServiceImpl.java @@ -184,15 +184,8 @@ public int resetTableSemantics(Integer datasourceId, List tableNames) { + " reset.")) .collect(Collectors.toCollection(java.util.LinkedHashSet::new)); List matchedTables = - tableInfoMapper.selectByDatasourceId(datasourceId).stream() - .filter( - table -> - normalizedNames.contains( - SemanticUtils.normalizeObjectName( - table.getTableName(), - "Missing tableName while matching table" - + " semantic reset."))) - .toList(); + tableInfoMapper.selectByDatasourceIdAndTableNames( + datasourceId, List.copyOf(normalizedNames)); if (matchedTables.isEmpty()) { return 0; } @@ -209,6 +202,7 @@ public int resetTableSemantics(Integer datasourceId, List tableNames) { private int resetTableRecords(Integer datasourceId, List tables) { List resetIds = new ArrayList<>(); List deleteIds = new ArrayList<>(); + // 仅删除明确标记为物理缺失的表;其他记录保留并重置语义字段。 for (TableInfo table : tables) { (Boolean.FALSE.equals(table.getPhysicalStatus()) ? deleteIds : resetIds) .add(table.getId()); From 6ba947660f36eba3b921ca12736e614db6f504c0 Mon Sep 17 00:00:00 2001 From: mengnankkkk Date: Thu, 1 Oct 2026 20:26:52 +0800 Subject: [PATCH 6/7] fix: remove add --- .../malonetalk/common/SemanticConstants.java | 3 - .../convertor/SemanticConverter.java | 7 - .../LogicalTableRelationResponse.java | 1 - .../mapper/ColumnSemanticInfoMapper.java | 4 + .../mapper/LogicalTableRelationMapper.java | 13 +- .../semantic/SemanticMergeService.java | 212 ++++++------------ .../relation/LogicalTableRelationHelper.java | 28 --- .../relation/RelationSemanticServiceImpl.java | 22 +- .../mapper/ColumnSemanticInfoMapper.xml | 10 + .../mapper/LogicalTableRelationMapper.xml | 27 +-- data-agent-frontend/src/api/semantic.ts | 3 +- .../semantic/components/RelationWorkspace.vue | 18 +- .../components/TableSemanticManage.vue | 10 +- docs/getting-started.md | 7 +- docs/semantic-layer.md | 2 +- sql/migration_primary_key_relations.sql | 22 +- 16 files changed, 135 insertions(+), 254 deletions(-) diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/common/SemanticConstants.java b/data-agent-backend/src/main/java/io/github/malonetalk/common/SemanticConstants.java index 0f9f0c32..5daf87d9 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/common/SemanticConstants.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/common/SemanticConstants.java @@ -21,9 +21,6 @@ public final class SemanticConstants { - public static final String RELATION_KEY_SEPARATOR = "|"; - public static final String RELATION_TABLE_COLUMN_SEPARATOR = ":"; - public static final String RELATION_GROUP_SEPARATOR = "::"; public static final String DEFAULT_DOMAIN = "default"; public static final String RELATION_TYPE_FOREIGN_KEY = LogicalTableRelationType.FOREIGN_KEY.getCode(); diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/convertor/SemanticConverter.java b/data-agent-backend/src/main/java/io/github/malonetalk/convertor/SemanticConverter.java index 1de51873..299cc6d8 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/convertor/SemanticConverter.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/convertor/SemanticConverter.java @@ -85,15 +85,8 @@ public LogicalTableRelationResponse toResponse(LogicalTableRelation relation) { List targetColumns = logicalTableRelationHelper.fromJson( relation.getTargetColumnNamesJson(), "targetColumnNames"); - String relationKey = - logicalTableRelationHelper.buildRelationKey( - relation.getSourceTableName(), - sourceColumns, - relation.getTargetTableName(), - targetColumns); return LogicalTableRelationResponse.builder() .id(relation.getId()) - .relationKey(relationKey) .datasourceId(relation.getDatasourceId()) .source(SemanticConstants.RELATION_SOURCE_LOGICAL) .sourceTableName(relation.getSourceTableName()) diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/dto/semantic/LogicalTableRelationResponse.java b/data-agent-backend/src/main/java/io/github/malonetalk/dto/semantic/LogicalTableRelationResponse.java index dd269844..2268fdd9 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/dto/semantic/LogicalTableRelationResponse.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/dto/semantic/LogicalTableRelationResponse.java @@ -25,7 +25,6 @@ @Builder public record LogicalTableRelationResponse( Integer id, - String relationKey, Integer datasourceId, String source, String sourceTableName, diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java index a83babcb..fa2bcb2d 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/ColumnSemanticInfoMapper.java @@ -21,6 +21,7 @@ import io.github.malonetalk.entity.ColumnInfo; import java.time.LocalDateTime; import java.util.List; +import java.util.Set; import org.apache.ibatis.annotations.Mapper; import org.apache.ibatis.annotations.Param; @@ -36,6 +37,9 @@ List selectByDatasourceIdAndTableNames( @Param("datasourceId") Integer datasourceId, @Param("tableNames") List tableNames); + List selectByDatasourceIdAndTableIds( + @Param("datasourceId") Integer datasourceId, @Param("tableIds") Set tableIds); + List selectPageByDatasourceIdAndTableName( @Param("query") ColumnSemanticPageQuery query, @Param("sortDescending") boolean sortDescending); diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/LogicalTableRelationMapper.java b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/LogicalTableRelationMapper.java index 72c50a83..b15f5f84 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/mapper/LogicalTableRelationMapper.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/mapper/LogicalTableRelationMapper.java @@ -21,6 +21,7 @@ import io.github.malonetalk.entity.LogicalTableRelation; import java.time.LocalDateTime; import java.util.List; +import java.util.Set; import org.apache.ibatis.annotations.Mapper; import org.apache.ibatis.annotations.Param; @@ -31,13 +32,9 @@ public interface LogicalTableRelationMapper { List selectByDatasourceId(@Param("datasourceId") Integer datasourceId); - List selectByDatasourceIdAndSourceTable( + List selectByDatasourceIdAndSourceTableIds( @Param("datasourceId") Integer datasourceId, - @Param("sourceTableName") String sourceTableName); - - List selectByDatasourceIdAndSourceTables( - @Param("datasourceId") Integer datasourceId, - @Param("sourceTableNames") List sourceTableNames); + @Param("sourceTableIds") Set sourceTableIds); List selectPageByDatasourceIdAndSourceTable( @Param("query") RelationSemanticPageQuery query, @@ -63,8 +60,4 @@ int deleteByIdsAndSourceTable( @Param("datasourceId") Integer datasourceId, @Param("sourceTableId") Integer sourceTableId, @Param("ids") List ids); - - int deleteByIds(@Param("datasourceId") Integer datasourceId, @Param("ids") List ids); - - int deleteByDatasourceId(@Param("datasourceId") Integer datasourceId); } diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/SemanticMergeService.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/SemanticMergeService.java index bace13f9..1a73b972 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/SemanticMergeService.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/SemanticMergeService.java @@ -35,12 +35,13 @@ import io.github.malonetalk.service.semantic.relation.LogicalTableRelationHelper; import io.github.malonetalk.utils.SemanticUtils; import java.util.ArrayList; -import java.util.Collections; import java.util.HashMap; -import java.util.LinkedHashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Set; +import java.util.stream.Collectors; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -58,17 +59,28 @@ public class SemanticMergeService { public List listVisibleTablesByDomains( Datasource datasource, List domains) { List normalizedDomains = normalizeDomains(domains); - // 关系和列按 table_id 对齐;表名只用于最终提示词和关系去重。 - TableIdIndex tableIndex = - TableIdIndex.of(tableInfoMapper.selectByDatasourceId(datasource.getId())); - TableColumnIndex columnIndex = - TableColumnIndex.of( - columnSemanticInfoMapper.selectByDatasourceId(datasource.getId())); - RelationSourceIndex relationIndex = - RelationSourceIndex.of( - logicalTableRelationMapper.selectByDatasourceId(datasource.getId())); + // 关系和列按 table_id 对齐;表名只用于最终提示词。 + List tables = tableInfoMapper.selectByDatasourceId(datasource.getId()); + Map tablesById = new HashMap<>(); + for (TableInfo table : tables) { + tablesById.put(table.getId(), table); + } + Map> availableColumnsByTableId = new HashMap<>(); + for (ColumnInfo column : + columnSemanticInfoMapper.selectByDatasourceId(datasource.getId())) { + if (SemanticAvailabilityHelper.isColumnAvailable(column)) { + availableColumnsByTableId + .computeIfAbsent(column.getTableId(), id -> new HashSet<>()) + .add( + SemanticUtils.normalizeObjectName( + column.getColumnName(), "Missing columnName.")); + } + } + Map> relationsBySourceId = + logicalTableRelationMapper.selectByDatasourceId(datasource.getId()).stream() + .collect(Collectors.groupingBy(LogicalTableRelation::getSourceTableId)); - return tableIndex.asList().stream() + return tables.stream() .filter( table -> domainMatches( @@ -78,11 +90,11 @@ public List listVisibleTablesByDomains( table -> PromptConverter.mapTablePrompt( table, - resolveVisibleRelations( - table.getId(), - tableIndex, - columnIndex, - relationIndex))) + filterVisibleLogicalRelations( + relationsBySourceId.getOrDefault( + table.getId(), List.of()), + tablesById, + availableColumnsByTableId))) .filter(Objects::nonNull) .toList(); } @@ -122,28 +134,17 @@ public List getTableSchema(Datasource datasource, String t .toList(); } - private List resolveVisibleRelations( - Integer sourceTableId, - TableIdIndex tableIndex, - TableColumnIndex columnIndex, - RelationSourceIndex relationIndex) { - List visibleRelations = - filterVisibleLogicalRelations( - relationIndex.get(sourceTableId), tableIndex, columnIndex); - return deduplicateRelations(visibleRelations); - } - - private List filterVisibleLogicalRelations( + private List filterVisibleLogicalRelations( List logicalRelations, - TableIdIndex tableIndex, - TableColumnIndex columnIndex) { - List visibleRelations = new ArrayList<>(); + Map tablesById, + Map> availableColumnsByTableId) { + List visibleRelations = new ArrayList<>(); for (LogicalTableRelation relation : logicalRelations) { if (!Boolean.TRUE.equals(relation.getIsEnabled())) { continue; } - if (tableIndex.isUnavailable(relation.getSourceTableId()) - || tableIndex.isUnavailable(relation.getTargetTableId())) { + if (isUnavailableTable(tablesById.get(relation.getSourceTableId())) + || isUnavailableTable(tablesById.get(relation.getTargetTableId()))) { continue; } List sourceColumns = parseRelationColumns(relation, true); @@ -151,17 +152,41 @@ private List filterVisibleLogicalRelations( if (sourceColumns == null || targetColumns == null) { continue; } - if (columnIndex.hasUnavailableColumn(relation.getSourceTableId(), sourceColumns) - || columnIndex.hasUnavailableColumn( - relation.getTargetTableId(), targetColumns)) { + if (hasUnavailableColumn( + availableColumnsByTableId, relation.getSourceTableId(), sourceColumns) + || hasUnavailableColumn( + availableColumnsByTableId, + relation.getTargetTableId(), + targetColumns)) { continue; } visibleRelations.add( - new ResolvedLogicalRelation(relation, sourceColumns, targetColumns)); + new TableRelationPromptResponse( + LogicalTableRelationType.fromCode(relation.getRelationType()), + SemanticConstants.RELATION_SOURCE_LOGICAL, + relation.getSourceTableName(), + sourceColumns, + relation.getTargetTableName(), + targetColumns, + relation.getDescription())); } return visibleRelations; } + private boolean hasUnavailableColumn( + Map> availableColumnsByTableId, + Integer tableId, + List columnNames) { + Set availableColumns = availableColumnsByTableId.getOrDefault(tableId, Set.of()); + return columnNames.stream() + .map(name -> SemanticUtils.normalizeObjectName(name, "Missing columnName.")) + .anyMatch(name -> !availableColumns.contains(name)); + } + + private boolean isUnavailableTable(TableInfo table) { + return table == null || !SemanticAvailabilityHelper.isTableAvailable(table); + } + private List parseRelationColumns(LogicalTableRelation relation, boolean source) { String fieldName = source ? "sourceColumnNames" : "targetColumnNames"; String json = @@ -178,32 +203,6 @@ private List parseRelationColumns(LogicalTableRelation relation, boolean } } - private List deduplicateRelations( - List relations) { - LinkedHashMap merged = new LinkedHashMap<>(); - for (ResolvedLogicalRelation relation : relations) { - String key = - logicalTableRelationHelper.buildRelationKey( - relation.sourceTableName(), - relation.sourceColumns(), - relation.targetTableName(), - relation.targetColumns()); - merged.put(key, toPromptResponse(relation)); - } - return List.copyOf(merged.values()); - } - - private TableRelationPromptResponse toPromptResponse(ResolvedLogicalRelation relation) { - return new TableRelationPromptResponse( - LogicalTableRelationType.fromCode(relation.relation().getRelationType()), - SemanticConstants.RELATION_SOURCE_LOGICAL, - relation.sourceTableName(), - relation.sourceColumns(), - relation.targetTableName(), - relation.targetColumns(), - relation.relation().getDescription()); - } - private List normalizeDomains(List domains) { if (domains == null || domains.isEmpty()) { return List.of(); @@ -221,89 +220,4 @@ private boolean domainMatches(String domain, List domains) { } return domains.stream().anyMatch(d -> d.equalsIgnoreCase(domain)); } - - private record ResolvedLogicalRelation( - LogicalTableRelation relation, List sourceColumns, List targetColumns) { - - private String sourceTableName() { - return relation.getSourceTableName(); - } - - private String targetTableName() { - return relation.getTargetTableName(); - } - } - - private record TableIdIndex(Map index) { - - private static TableIdIndex of(List tables) { - Map map = new LinkedHashMap<>(); - for (TableInfo table : tables) { - map.put(table.getId(), table); - } - return new TableIdIndex(map); - } - - private List asList() { - return List.copyOf(index.values()); - } - - private boolean isUnavailable(Integer tableId) { - TableInfo tableInfo = index.get(tableId); - return tableInfo == null || !SemanticAvailabilityHelper.isTableAvailable(tableInfo); - } - } - - private record TableColumnIndex(Map> index) { - - private static TableColumnIndex of(List columns) { - Map> map = new HashMap<>(); - for (ColumnInfo column : columns) { - map.computeIfAbsent(column.getTableId(), key -> new HashMap<>()) - .put( - SemanticUtils.normalizeObjectName( - column.getColumnName(), - "Missing columnName while building semantic column index."), - column); - } - return new TableColumnIndex(map); - } - - private ColumnInfo get(Integer tableId, String columnName) { - Map columns = index.get(tableId); - return columns == null - ? null - : columns.get( - SemanticUtils.normalizeObjectName( - columnName, - "Missing columnName while reading semantic column index.")); - } - - private boolean hasUnavailableColumn(Integer tableId, List columnNames) { - for (String columnName : columnNames) { - ColumnInfo columnInfo = get(tableId, columnName); - if (columnInfo == null - || !SemanticAvailabilityHelper.isColumnAvailable(columnInfo)) { - return true; - } - } - return false; - } - } - - private record RelationSourceIndex(Map> index) { - - private static RelationSourceIndex of(List relations) { - Map> map = new HashMap<>(); - for (LogicalTableRelation relation : relations) { - map.computeIfAbsent(relation.getSourceTableId(), key -> new ArrayList<>()) - .add(relation); - } - return new RelationSourceIndex(map); - } - - private List get(Integer sourceTableId) { - return index.getOrDefault(sourceTableId, Collections.emptyList()); - } - } } diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/LogicalTableRelationHelper.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/LogicalTableRelationHelper.java index 9a367572..62777ccb 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/LogicalTableRelationHelper.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/LogicalTableRelationHelper.java @@ -17,23 +17,17 @@ */ package io.github.malonetalk.service.semantic.relation; -import static io.github.malonetalk.common.SemanticConstants.RELATION_GROUP_SEPARATOR; -import static io.github.malonetalk.common.SemanticConstants.RELATION_KEY_SEPARATOR; -import static io.github.malonetalk.common.SemanticConstants.RELATION_TABLE_COLUMN_SEPARATOR; - import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; import io.github.malonetalk.exception.BusinessException; import io.github.malonetalk.exception.ErrorCode; import io.github.malonetalk.utils.RequestAssert; -import io.github.malonetalk.utils.SemanticUtils; import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Locale; import java.util.Set; -import java.util.stream.Collectors; import org.springframework.stereotype.Component; @Component @@ -67,28 +61,6 @@ public List normalizeColumnNames(List columnNames, String fieldN return List.copyOf(normalizedColumns); } - private String buildColumnSignature(List columnNames) { - return normalizeColumnNames(columnNames, "columnNames").stream() - .map(columnName -> columnName.toLowerCase(Locale.ROOT)) - .collect(Collectors.joining(RELATION_KEY_SEPARATOR)); - } - - public String buildRelationKey( - String sourceTableName, - List sourceColumnNames, - String targetTableName, - List targetColumnNames) { - return SemanticUtils.normalizeObjectName( - sourceTableName, "Missing sourceTableName for logical relation key.") - + RELATION_TABLE_COLUMN_SEPARATOR - + buildColumnSignature(sourceColumnNames) - + RELATION_GROUP_SEPARATOR - + SemanticUtils.normalizeObjectName( - targetTableName, "Missing targetTableName for logical relation key.") - + RELATION_TABLE_COLUMN_SEPARATOR - + buildColumnSignature(targetColumnNames); - } - public String toJson(List columnNames, String fieldName) { try { return objectMapper.writeValueAsString(normalizeColumnNames(columnNames, fieldName)); diff --git a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java index 46b74fce..d8a18071 100644 --- a/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java +++ b/data-agent-backend/src/main/java/io/github/malonetalk/service/semantic/relation/RelationSemanticServiceImpl.java @@ -116,13 +116,12 @@ public RelationWorkspaceResponse getRelationWorkspace(RelationWorkspacePageQuery PageResponse.empty(pageNumber, pageSize), List.of()); } - List tableNames = page.stream().map(TableInfo::getTableName).toList(); Set currentPageTableIds = page.stream().map(TableInfo::getId).collect(Collectors.toSet()); // 列记录已持有 table_id,按主键分组可直接关联当前页的表。 Map> columnsByTableId = columnSemanticInfoMapper - .selectByDatasourceIdAndTableNames(query.datasourceId(), tableNames) + .selectByDatasourceIdAndTableIds(query.datasourceId(), currentPageTableIds) .stream() .collect(Collectors.groupingBy(ColumnInfo::getTableId)); List nodes = @@ -137,7 +136,8 @@ public RelationWorkspaceResponse getRelationWorkspace(RelationWorkspacePageQuery // 来源表来自当前页查询,目标表也必须属于当前页。 List relations = logicalTableRelationMapper - .selectByDatasourceIdAndSourceTables(query.datasourceId(), tableNames) + .selectByDatasourceIdAndSourceTableIds( + query.datasourceId(), currentPageTableIds) .stream() .filter( relation -> @@ -213,23 +213,17 @@ public int deleteRelationSemantics( return 0; } TableInfo sourceTable = requireTable(datasourceId, normalizedTableName, "sourceTable"); - List matchedIds = - logicalTableRelationMapper - .selectByDatasourceIdAndSourceTable(datasourceId, normalizedTableName) - .stream() - .map(LogicalTableRelation::getId) - .filter(relationIds::contains) - .distinct() - .toList(); - if (matchedIds.size() != relationIds.size()) { + int deleted = + logicalTableRelationMapper.deleteByIdsAndSourceTable( + datasourceId, sourceTable.getId(), relationIds); + if (deleted != relationIds.size()) { throw BusinessException.of( ErrorCode.RESOURCE_NOT_FOUND, "Some logical relations do not exist or do not belong to source table " + normalizedTableName + "."); } - return logicalTableRelationMapper.deleteByIdsAndSourceTable( - datasourceId, sourceTable.getId(), relationIds); + return deleted; } private void requireDatasource(Integer datasourceId) { diff --git a/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml b/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml index 8742b306..db66b517 100644 --- a/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml @@ -49,6 +49,16 @@ ORDER BY table_meta.table_name ASC, c.column_name ASC, c.id ASC + + - WHERE relation.datasource_id = #{datasourceId} - AND LOWER(source_table.table_name) = LOWER(#{sourceTableName}) - ORDER BY relation.id DESC - - - @@ -127,16 +120,4 @@ - - DELETE FROM logical_table_relation - WHERE datasource_id = #{datasourceId} - AND id IN - - #{id} - - - - - DELETE FROM logical_table_relation WHERE datasource_id = #{datasourceId} - diff --git a/data-agent-frontend/src/api/semantic.ts b/data-agent-frontend/src/api/semantic.ts index b78aad3b..00cec54f 100644 --- a/data-agent-frontend/src/api/semantic.ts +++ b/data-agent-frontend/src/api/semantic.ts @@ -92,8 +92,7 @@ export interface SyncTableSemanticsResponse { } export interface LogicalTableRelationResponse { - id: number | null; - relationKey: string; + id: number; datasourceId: number; source: 'physical' | 'logical'; sourceTableName: string; diff --git a/data-agent-frontend/src/views/semantic/components/RelationWorkspace.vue b/data-agent-frontend/src/views/semantic/components/RelationWorkspace.vue index 42e82158..b8952891 100644 --- a/data-agent-frontend/src/views/semantic/components/RelationWorkspace.vue +++ b/data-agent-frontend/src/views/semantic/components/RelationWorkspace.vue @@ -29,7 +29,7 @@ interface RelationEdge { id: string; - relationId: string; + relationId: number; path: string; label: string; labelX: number; @@ -104,7 +104,7 @@ const localNodes = ref([]); const dragRelation = ref(null); const hoveredDropColumn = ref<{ tableName: string; columnName: string } | null>(null); - const selectedRelationId = ref(null); + const selectedRelationId = ref(null); const canvasPan = ref(null); const nodeDrag = ref(null); const viewport = reactive({ @@ -163,8 +163,8 @@ return [ { - id: `relation-${relation.relationKey}`, - relationId: relation.relationKey, + id: `relation-${relation.id}`, + relationId: relation.id, path: buildRelationPath(sourceAnchor.x, sourceAnchor.y, targetAnchor.x, targetAnchor.y), label, labelX: (sourceAnchor.x + targetAnchor.x) / 2, @@ -461,7 +461,7 @@ zoomAtCenter(1 / 1.2); } - function selectRelation(relationId: string) { + function selectRelation(relationId: number) { selectedRelationId.value = relationId; } @@ -647,7 +647,7 @@ }; } - function isSelected(relationId: string) { + function isSelected(relationId: number) { return selectedRelationId.value === relationId; } @@ -845,10 +845,10 @@
diff --git a/data-agent-frontend/src/views/semantic/components/TableSemanticManage.vue b/data-agent-frontend/src/views/semantic/components/TableSemanticManage.vue index 46279509..d86d9a2e 100644 --- a/data-agent-frontend/src/views/semantic/components/TableSemanticManage.vue +++ b/data-agent-frontend/src/views/semantic/components/TableSemanticManage.vue @@ -176,9 +176,13 @@ const handleReset = async (row: TableSemanticInfo) => { try { - await ElMessageBox.confirm(`确认重置表 ${row.tableName} 的语义信息吗?`, '确认重置', { - type: 'warning', - }); + await ElMessageBox.confirm( + `确认重置表 ${row.tableName} 的语义信息吗?若该表当前已标记为物理缺失,还会删除该表及其列的语义记录和所有关联到该表的逻辑关系;重新同步不会恢复这些关系。`, + '确认重置', + { + type: 'warning', + }, + ); const activeDatasourceId = await ensureDatasourceId(); if (activeDatasourceId === null) return; await resetTableSemantic(activeDatasourceId, row.tableName); diff --git a/docs/getting-started.md b/docs/getting-started.md index 18f390c0..e9028a42 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -94,15 +94,16 @@ mysql -u root -p data_agent < sql/data_source.sql > `sql/data_source.sql` 包含全部元数据库初始化表结构,导入这一份即可。 -已有数据库升级到主键关联版本时,不要重复初始化;请先备份数据库,再执行一次迁移脚本: +已有数据库升级到主键关联版本时,不要重复初始化;请先备份数据库,依次执行迁移脚本: ```bash +mysql -u root -p data_agent < sql/migration_compatibility.sql mysql -u root -p data_agent < sql/migration_primary_key_relations.sql ``` 迁移必须在启动新版后端前完成。脚本会将列和逻辑关系中的表名引用转换为 -`table_info.id` 外键;脚本会在修改业务表结构前检查无法匹配的旧数据。若检查失败, -修复数据后重新执行整份脚本即可。 +`table_info.id` 外键;脚本会在改表前检查无法匹配的引用、同一数据源内重复的表名, +以及同一表内重复的列名。若检查失败,修复数据后重新执行整份脚本即可。 ### 3. 配置并启动后端 diff --git a/docs/semantic-layer.md b/docs/semantic-layer.md index 8e93d762..638f6e39 100644 --- a/docs/semantic-layer.md +++ b/docs/semantic-layer.md @@ -44,7 +44,7 @@ Agent 查询前按需调用语义工具: 这条链路依赖语义层同步后的缓存。新增或变更业务库表结构后,请先在「语义管理 / 表语义」里同步物理表,再让 Agent 查询。 -已有元数据库升级时,执行 `sql/migration_compatibility.sql` 补齐数据粒度和列语义类型字段;脚本可重复执行,已存在的字段会自动跳过。全新安装无需额外执行。 +已有元数据库升级时,先执行 `sql/migration_compatibility.sql` 补齐数据粒度和列语义类型字段,再执行 `sql/migration_primary_key_relations.sql` 将逻辑关系迁移到表、列主键。请在启动新版后端前完成这两个迁移,并在迁移前备份元数据库。全新安装无需额外执行。 ## 5. 最佳实践 diff --git a/sql/migration_primary_key_relations.sql b/sql/migration_primary_key_relations.sql index 1daa0598..4ec95490 100644 --- a/sql/migration_primary_key_relations.sql +++ b/sql/migration_primary_key_relations.sql @@ -1,11 +1,31 @@ -- One-time MySQL 5.7 migration for primary-key-based semantic table references. -- Back up the database first. Validate all references before altering application tables, --- so fixing unmatched legacy names and rerunning does not hit an already-added column. +-- so fixing invalid legacy data and rerunning does not hit an already-added column. DROP PROCEDURE IF EXISTS `check_primary_key_relation_migration`; DELIMITER $$ CREATE PROCEDURE `check_primary_key_relation_migration`() BEGIN + IF EXISTS ( + SELECT 1 + FROM `table_info` + GROUP BY `datasource_id`, LOWER(`table_name`) + HAVING COUNT(*) > 1 + ) THEN + SIGNAL SQLSTATE '45000' + SET MESSAGE_TEXT = 'table_info has duplicate datasource/table names'; + END IF; + + IF EXISTS ( + SELECT 1 + FROM `column_info` + GROUP BY `datasource_id`, LOWER(`table_name`), LOWER(`column_name`) + HAVING COUNT(*) > 1 + ) THEN + SIGNAL SQLSTATE '45000' + SET MESSAGE_TEXT = 'column_info has duplicate datasource/table/column names'; + END IF; + IF EXISTS ( SELECT 1 FROM `column_info` c From e0c93bc83bba0692dc0d91761a557f2c2425fd37 Mon Sep 17 00:00:00 2001 From: mengnankkkk Date: Thu, 1 Oct 2026 21:15:55 +0800 Subject: [PATCH 7/7] fix: consolidate metadata migration into one rerunnable script --- docs/configuration.md | 4 +- docs/getting-started.md | 5 +- docs/semantic-layer.md | 2 +- sql/migration_compatibility.sql | 203 -------- sql/migration_primary_key_relations.sql | 607 ++++++++++++++++++++---- 5 files changed, 521 insertions(+), 300 deletions(-) delete mode 100644 sql/migration_compatibility.sql diff --git a/docs/configuration.md b/docs/configuration.md index 482eac99..81531325 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -13,8 +13,8 @@ | `spring.datasource.password` | `root` | 支持 `${DB_PASSWORD:...}` | | `spring.datasource.driver-class-name` | `com.mysql.cj.jdbc.Driver` | MySQL 驱动 | -新库建表脚本见 `sql/data_source.sql`。已有数据库升级到主键关联版本时,需在启动新版后端前执行 -`sql/migration_primary_key_relations.sql`。 +新库建表脚本见 `sql/data_source.sql`。已有数据库升级到主键关联版本时,需在启动新版后端前 +执行 `sql/migration_primary_key_relations.sql`,其中包含旧库兼容和主键关联升级。 ## 2. LLM 提供商 diff --git a/docs/getting-started.md b/docs/getting-started.md index e9028a42..6b359931 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -94,14 +94,13 @@ mysql -u root -p data_agent < sql/data_source.sql > `sql/data_source.sql` 包含全部元数据库初始化表结构,导入这一份即可。 -已有数据库升级到主键关联版本时,不要重复初始化;请先备份数据库,依次执行迁移脚本: +已有数据库升级到主键关联版本时,不要重复初始化;请先备份数据库,再执行迁移脚本: ```bash -mysql -u root -p data_agent < sql/migration_compatibility.sql mysql -u root -p data_agent < sql/migration_primary_key_relations.sql ``` -迁移必须在启动新版后端前完成。脚本会将列和逻辑关系中的表名引用转换为 +迁移必须在启动新版后端前完成。脚本会补齐旧库缺少的字段与索引,并将列和逻辑关系中的表名引用转换为 `table_info.id` 外键;脚本会在改表前检查无法匹配的引用、同一数据源内重复的表名, 以及同一表内重复的列名。若检查失败,修复数据后重新执行整份脚本即可。 diff --git a/docs/semantic-layer.md b/docs/semantic-layer.md index 638f6e39..d341a461 100644 --- a/docs/semantic-layer.md +++ b/docs/semantic-layer.md @@ -44,7 +44,7 @@ Agent 查询前按需调用语义工具: 这条链路依赖语义层同步后的缓存。新增或变更业务库表结构后,请先在「语义管理 / 表语义」里同步物理表,再让 Agent 查询。 -已有元数据库升级时,先执行 `sql/migration_compatibility.sql` 补齐数据粒度和列语义类型字段,再执行 `sql/migration_primary_key_relations.sql` 将逻辑关系迁移到表、列主键。请在启动新版后端前完成这两个迁移,并在迁移前备份元数据库。全新安装无需额外执行。 +已有元数据库升级时,备份数据库后执行 `sql/migration_primary_key_relations.sql`。脚本会补齐旧库字段与索引,并将列和逻辑关系迁移到表主键引用。请在启动新版后端前完成迁移。全新安装无需额外执行。 ## 5. 最佳实践 diff --git a/sql/migration_compatibility.sql b/sql/migration_compatibility.sql deleted file mode 100644 index a70db097..00000000 --- a/sql/migration_compatibility.sql +++ /dev/null @@ -1,203 +0,0 @@ -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'table_info' - AND COLUMN_NAME = 'data_granularity' - ), - 'SELECT 1', - 'ALTER TABLE `table_info` ADD COLUMN `data_granularity` VARCHAR(500) DEFAULT NULL COMMENT ''数据粒度(一行代表的业务实体或事件)'' AFTER `table_description`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'column_info' - AND COLUMN_NAME = 'semantic_type' - ), - 'SELECT 1', - 'ALTER TABLE `column_info` ADD COLUMN `semantic_type` VARCHAR(32) DEFAULT NULL COMMENT ''语义类型:IDENTIFIER/DIMENSION/MEASURE/TIME/LOCATION'' AFTER `column_description`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'sys_user' - AND COLUMN_NAME = 'is_super_admin' - ), - 'SELECT 1', - 'ALTER TABLE `sys_user` ADD COLUMN `is_super_admin` TINYINT(1) NOT NULL DEFAULT 0 COMMENT ''是否超级管理员:0否,1是'' AFTER `role_id`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - --- sys_user 审计字段:creator_id / create_time / updater_id / update_time / is_deleted -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'sys_user' - AND COLUMN_NAME = 'creator_id' - ), - 'SELECT 1', - 'ALTER TABLE `sys_user` ADD COLUMN `creator_id` BIGINT DEFAULT NULL COMMENT ''创建人ID'' AFTER `status`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'sys_user' - AND COLUMN_NAME = 'create_time' - ), - 'SELECT 1', - 'ALTER TABLE `sys_user` ADD COLUMN `create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT ''创建时间'' AFTER `creator_id`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'sys_user' - AND COLUMN_NAME = 'updater_id' - ), - 'SELECT 1', - 'ALTER TABLE `sys_user` ADD COLUMN `updater_id` BIGINT DEFAULT NULL COMMENT ''修改人ID'' AFTER `create_time`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'sys_user' - AND COLUMN_NAME = 'update_time' - ), - 'SELECT 1', - 'ALTER TABLE `sys_user` ADD COLUMN `update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT ''更新时间'' AFTER `updater_id`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'sys_user' - AND COLUMN_NAME = 'is_deleted' - ), - 'SELECT 1', - 'ALTER TABLE `sys_user` ADD COLUMN `is_deleted` TINYINT(1) NOT NULL DEFAULT 0 COMMENT ''逻辑删除标记:0-未删除,1-已删除'' AFTER `update_time`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - --- 逻辑删除后允许同身份源用户重建:唯一键改为普通索引 -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.STATISTICS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'sys_user' - AND INDEX_NAME = 'uk_idp' - ), - 'ALTER TABLE `sys_user` DROP INDEX `uk_idp`, ADD KEY `idx_idp` (`idp_type`, `idp_user_id`)', - IF( - EXISTS( - SELECT 1 - FROM information_schema.STATISTICS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'sys_user' - AND INDEX_NAME = 'idx_idp' - ), - 'SELECT 1', - 'ALTER TABLE `sys_user` ADD KEY `idx_idp` (`idp_type`, `idp_user_id`)' - ) -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - --- metric_info 审计字段:creator_id / updater_id -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'metric_info' - AND COLUMN_NAME = 'creator_id' - ), - 'SELECT 1', - 'ALTER TABLE `metric_info` ADD COLUMN `creator_id` BIGINT DEFAULT NULL COMMENT ''创建人ID'' AFTER `description`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'metric_info' - AND COLUMN_NAME = 'updater_id' - ), - 'SELECT 1', - 'ALTER TABLE `metric_info` ADD COLUMN `updater_id` BIGINT DEFAULT NULL COMMENT ''修改人ID'' AFTER `create_time`' -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; - --- 逻辑删除后保留行以便审计:同 key 软删后允许重新创建,唯一键改为普通索引 -SET @migration_sql = IF( - EXISTS( - SELECT 1 - FROM information_schema.STATISTICS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'metric_info' - AND INDEX_NAME = 'uk_datasource_metric_key' - ), - 'ALTER TABLE `metric_info` DROP INDEX `uk_datasource_metric_key`, ADD KEY `idx_datasource_metric_key` (`datasource_id`, `metric_key`)', - IF( - EXISTS( - SELECT 1 - FROM information_schema.STATISTICS - WHERE TABLE_SCHEMA = DATABASE() - AND TABLE_NAME = 'metric_info' - AND INDEX_NAME = 'idx_datasource_metric_key' - ), - 'SELECT 1', - 'ALTER TABLE `metric_info` ADD KEY `idx_datasource_metric_key` (`datasource_id`, `metric_key`)' - ) -); -PREPARE migration_statement FROM @migration_sql; -EXECUTE migration_statement; -DEALLOCATE PREPARE migration_statement; diff --git a/sql/migration_primary_key_relations.sql b/sql/migration_primary_key_relations.sql index 4ec95490..0b204ba1 100644 --- a/sql/migration_primary_key_relations.sql +++ b/sql/migration_primary_key_relations.sql @@ -1,117 +1,542 @@ --- One-time MySQL 5.7 migration for primary-key-based semantic table references. --- Back up the database first. Validate all references before altering application tables, --- so fixing invalid legacy data and rerunning does not hit an already-added column. +-- MySQL 5.7: upgrade legacy metadata and migrate semantic references to table_info.id. +-- Back up the database and stop writes before running this single migration script. +-- ALTER TABLE commits implicitly. Every schema step is conditional so an interrupted +-- migration can be rerun after the cause of failure has been fixed. -DROP PROCEDURE IF EXISTS `check_primary_key_relation_migration`; +DROP PROCEDURE IF EXISTS `migrate_primary_key_relations`; DELIMITER $$ -CREATE PROCEDURE `check_primary_key_relation_migration`() +CREATE PROCEDURE `migrate_primary_key_relations`() BEGIN + DECLARE old_column_name INT DEFAULT 0; + DECLARE old_source_name INT DEFAULT 0; + DECLARE old_target_name INT DEFAULT 0; + DECLARE has_column_id INT DEFAULT 0; + DECLARE has_source_id INT DEFAULT 0; + DECLARE has_target_id INT DEFAULT 0; + + SELECT COUNT(*) INTO old_column_name FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'column_info' AND COLUMN_NAME = 'table_name'; + SELECT COUNT(*) INTO old_source_name FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND COLUMN_NAME = 'source_table_name'; + SELECT COUNT(*) INTO old_target_name FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND COLUMN_NAME = 'target_table_name'; + SELECT COUNT(*) INTO has_column_id FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'column_info' AND COLUMN_NAME = 'table_id'; + SELECT COUNT(*) INTO has_source_id FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND COLUMN_NAME = 'source_table_id'; + SELECT COUNT(*) INTO has_target_id FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND COLUMN_NAME = 'target_table_id'; + + IF (old_column_name = 0 AND has_column_id = 0) + OR (old_source_name = 0 AND has_source_id = 0) + OR (old_target_name = 0 AND has_target_id = 0) THEN + SIGNAL SQLSTATE '45000' + SET MESSAGE_TEXT = 'Unsupported schema: table name and ID are both absent'; + END IF; + + -- Validate legacy data before the first ALTER TABLE. Never guess which duplicate + -- table should own a column or relation. IF EXISTS ( - SELECT 1 - FROM `table_info` + SELECT 1 FROM `table_info` GROUP BY `datasource_id`, LOWER(`table_name`) HAVING COUNT(*) > 1 ) THEN SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = 'table_info has duplicate datasource/table names'; END IF; + IF old_column_name = 1 THEN + IF EXISTS ( + SELECT 1 FROM `column_info` + GROUP BY `datasource_id`, LOWER(`table_name`), LOWER(`column_name`) + HAVING COUNT(*) > 1 + ) THEN + SIGNAL SQLSTATE '45000' + SET MESSAGE_TEXT = 'column_info has duplicate datasource/table/column names'; + END IF; + IF EXISTS ( + SELECT 1 FROM `column_info` c + LEFT JOIN `table_info` t ON t.`datasource_id` = c.`datasource_id` + AND LOWER(t.`table_name`) = LOWER(c.`table_name`) + WHERE t.`id` IS NULL + ) THEN + SIGNAL SQLSTATE '45000' + SET MESSAGE_TEXT = 'column_info has table names missing from table_info'; + END IF; + END IF; + IF old_source_name = 1 THEN + IF EXISTS ( + SELECT 1 FROM `logical_table_relation` r + LEFT JOIN `table_info` t ON t.`datasource_id` = r.`datasource_id` + AND LOWER(t.`table_name`) = LOWER(r.`source_table_name`) + WHERE t.`id` IS NULL + ) THEN + SIGNAL SQLSTATE '45000' + SET MESSAGE_TEXT = 'logical_table_relation has missing source tables'; + END IF; + END IF; + IF old_target_name = 1 THEN + IF EXISTS ( + SELECT 1 FROM `logical_table_relation` r + LEFT JOIN `table_info` t ON t.`datasource_id` = r.`datasource_id` + AND LOWER(t.`table_name`) = LOWER(r.`target_table_name`) + WHERE t.`id` IS NULL + ) THEN + SIGNAL SQLSTATE '45000' + SET MESSAGE_TEXT = 'logical_table_relation has missing target tables'; + END IF; + END IF; + + IF has_column_id = 0 THEN + ALTER TABLE `column_info` + ADD COLUMN `table_id` INT NULL COMMENT '关联表信息ID' AFTER `datasource_id`; + END IF; + IF has_source_id = 0 THEN + ALTER TABLE `logical_table_relation` + ADD COLUMN `source_table_id` INT NULL COMMENT '源表信息ID' AFTER `datasource_id`; + END IF; + IF has_target_id = 0 THEN + ALTER TABLE `logical_table_relation` + ADD COLUMN `target_table_id` INT NULL COMMENT '目标表信息ID' + AFTER `source_column_names_json`; + END IF; + + IF old_column_name = 1 THEN + UPDATE `column_info` c + INNER JOIN `table_info` t ON t.`datasource_id` = c.`datasource_id` + AND LOWER(t.`table_name`) = LOWER(c.`table_name`) + SET c.`table_id` = t.`id`; + END IF; + IF old_source_name = 1 THEN + UPDATE `logical_table_relation` r + INNER JOIN `table_info` t ON t.`datasource_id` = r.`datasource_id` + AND LOWER(t.`table_name`) = LOWER(r.`source_table_name`) + SET r.`source_table_id` = t.`id`; + END IF; + IF old_target_name = 1 THEN + UPDATE `logical_table_relation` r + INNER JOIN `table_info` t ON t.`datasource_id` = r.`datasource_id` + AND LOWER(t.`table_name`) = LOWER(r.`target_table_name`) + SET r.`target_table_id` = t.`id`; + END IF; + -- Validate IDs even when resuming from a partially migrated schema. IF EXISTS ( - SELECT 1 - FROM `column_info` - GROUP BY `datasource_id`, LOWER(`table_name`), LOWER(`column_name`) + SELECT 1 FROM `column_info` c + LEFT JOIN `table_info` t ON t.`id` = c.`table_id` + AND t.`datasource_id` = c.`datasource_id` + WHERE t.`id` IS NULL + ) THEN + SIGNAL SQLSTATE '45000' + SET MESSAGE_TEXT = 'column_info has invalid table_id references'; + END IF; + IF EXISTS ( + SELECT 1 FROM `column_info` + GROUP BY `table_id`, LOWER(`column_name`) HAVING COUNT(*) > 1 ) THEN SIGNAL SQLSTATE '45000' - SET MESSAGE_TEXT = 'column_info has duplicate datasource/table/column names'; + SET MESSAGE_TEXT = 'column_info has duplicate table_id/column names'; END IF; - IF EXISTS ( - SELECT 1 - FROM `column_info` c - LEFT JOIN `table_info` t - ON t.`datasource_id` = c.`datasource_id` - AND LOWER(t.`table_name`) = LOWER(c.`table_name`) - WHERE t.`id` IS NULL + SELECT 1 FROM `logical_table_relation` r + LEFT JOIN `table_info` s ON s.`id` = r.`source_table_id` + AND s.`datasource_id` = r.`datasource_id` + LEFT JOIN `table_info` t ON t.`id` = r.`target_table_id` + AND t.`datasource_id` = r.`datasource_id` + WHERE s.`id` IS NULL OR t.`id` IS NULL ) THEN SIGNAL SQLSTATE '45000' - SET MESSAGE_TEXT = 'column_info has table names missing from table_info'; + SET MESSAGE_TEXT = 'logical_table_relation has invalid table_id references'; END IF; IF EXISTS ( - SELECT 1 - FROM `logical_table_relation` r - LEFT JOIN `table_info` source_table - ON source_table.`datasource_id` = r.`datasource_id` - AND LOWER(source_table.`table_name`) = LOWER(r.`source_table_name`) - LEFT JOIN `table_info` target_table - ON target_table.`datasource_id` = r.`datasource_id` - AND LOWER(target_table.`table_name`) = LOWER(r.`target_table_name`) - WHERE source_table.`id` IS NULL OR target_table.`id` IS NULL + SELECT 1 FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'column_info' + AND COLUMN_NAME = 'table_id' AND IS_NULLABLE = 'YES' ) THEN - SIGNAL SQLSTATE '45000' - SET MESSAGE_TEXT = 'logical_table_relation has table names missing from table_info'; + ALTER TABLE `column_info` + MODIFY COLUMN `table_id` INT NOT NULL COMMENT '关联表信息ID'; + END IF; + IF NOT EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'column_info' + AND INDEX_NAME = 'uk_table_column' + ) THEN + ALTER TABLE `column_info` + ADD UNIQUE KEY `uk_table_column` (`table_id`, `column_name`); + END IF; + IF NOT EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'column_info' + AND INDEX_NAME = 'idx_datasource_id' + ) THEN + ALTER TABLE `column_info` ADD KEY `idx_datasource_id` (`datasource_id`); + END IF; + IF NOT EXISTS ( + SELECT 1 FROM information_schema.TABLE_CONSTRAINTS + WHERE CONSTRAINT_SCHEMA = DATABASE() AND TABLE_NAME = 'column_info' + AND CONSTRAINT_NAME = 'fk_column_info_table' + ) THEN + ALTER TABLE `column_info` + ADD CONSTRAINT `fk_column_info_table` + FOREIGN KEY (`table_id`) REFERENCES `table_info` (`id`) + ON DELETE CASCADE ON UPDATE RESTRICT; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'column_info' + AND INDEX_NAME = 'uk_datasource_table_column' + ) THEN + ALTER TABLE `column_info` DROP INDEX `uk_datasource_table_column`; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'column_info' + AND INDEX_NAME = 'idx_datasource_table_visible' + ) THEN + ALTER TABLE `column_info` DROP INDEX `idx_datasource_table_visible`; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'column_info' + AND INDEX_NAME = 'idx_datasource_table_visible_column' + ) THEN + ALTER TABLE `column_info` DROP INDEX `idx_datasource_table_visible_column`; + END IF; + IF old_column_name = 1 THEN + ALTER TABLE `column_info` DROP COLUMN `table_name`; + END IF; + + IF EXISTS ( + SELECT 1 FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND COLUMN_NAME = 'source_table_id' AND IS_NULLABLE = 'YES' + ) THEN + ALTER TABLE `logical_table_relation` + MODIFY COLUMN `source_table_id` INT NOT NULL COMMENT '源表信息ID'; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND COLUMN_NAME = 'target_table_id' AND IS_NULLABLE = 'YES' + ) THEN + ALTER TABLE `logical_table_relation` + MODIFY COLUMN `target_table_id` INT NOT NULL COMMENT '目标表信息ID'; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND INDEX_NAME = 'idx_relation_source_signature' + ) THEN + ALTER TABLE `logical_table_relation` DROP INDEX `idx_relation_source_signature`; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND INDEX_NAME = 'idx_relation_source_enabled' + ) THEN + ALTER TABLE `logical_table_relation` DROP INDEX `idx_relation_source_enabled`; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND INDEX_NAME = 'idx_relation_source_target_id' + ) THEN + ALTER TABLE `logical_table_relation` DROP INDEX `idx_relation_source_target_id`; + END IF; + -- These two index names are reused; remove only the legacy definitions. + IF EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND INDEX_NAME = 'idx_relation_source_table' AND COLUMN_NAME = 'source_table_name' + ) THEN + ALTER TABLE `logical_table_relation` DROP INDEX `idx_relation_source_table`; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND INDEX_NAME = 'idx_relation_source_enabled_id' + AND COLUMN_NAME = 'source_table_name' + ) THEN + ALTER TABLE `logical_table_relation` DROP INDEX `idx_relation_source_enabled_id`; + END IF; + IF NOT EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND INDEX_NAME = 'idx_relation_source_table' + ) THEN + ALTER TABLE `logical_table_relation` + ADD KEY `idx_relation_source_table` (`source_table_id`); + END IF; + IF NOT EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND INDEX_NAME = 'idx_relation_target_table' + ) THEN + ALTER TABLE `logical_table_relation` + ADD KEY `idx_relation_target_table` (`target_table_id`); + END IF; + IF NOT EXISTS ( + SELECT 1 FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND INDEX_NAME = 'idx_relation_source_enabled_id' + ) THEN + ALTER TABLE `logical_table_relation` + ADD KEY `idx_relation_source_enabled_id` + (`datasource_id`, `source_table_id`, `is_enabled`, `id`); + END IF; + IF NOT EXISTS ( + SELECT 1 FROM information_schema.TABLE_CONSTRAINTS + WHERE CONSTRAINT_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND CONSTRAINT_NAME = 'fk_relation_source_table' + ) THEN + ALTER TABLE `logical_table_relation` + ADD CONSTRAINT `fk_relation_source_table` + FOREIGN KEY (`source_table_id`) REFERENCES `table_info` (`id`) + ON DELETE CASCADE ON UPDATE RESTRICT; + END IF; + IF NOT EXISTS ( + SELECT 1 FROM information_schema.TABLE_CONSTRAINTS + WHERE CONSTRAINT_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND CONSTRAINT_NAME = 'fk_relation_target_table' + ) THEN + ALTER TABLE `logical_table_relation` + ADD CONSTRAINT `fk_relation_target_table` + FOREIGN KEY (`target_table_id`) REFERENCES `table_info` (`id`) + ON DELETE CASCADE ON UPDATE RESTRICT; + END IF; + IF old_source_name = 1 THEN + ALTER TABLE `logical_table_relation` DROP COLUMN `source_table_name`; + END IF; + IF old_target_name = 1 THEN + ALTER TABLE `logical_table_relation` DROP COLUMN `target_table_name`; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND COLUMN_NAME = 'source_column_signature' + ) THEN + ALTER TABLE `logical_table_relation` DROP COLUMN `source_column_signature`; + END IF; + IF EXISTS ( + SELECT 1 FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'logical_table_relation' + AND COLUMN_NAME = 'target_column_signature' + ) THEN + ALTER TABLE `logical_table_relation` DROP COLUMN `target_column_signature`; END IF; END$$ DELIMITER ; -CALL `check_primary_key_relation_migration`(); -DROP PROCEDURE `check_primary_key_relation_migration`; - -ALTER TABLE `column_info` - ADD COLUMN `table_id` INT NULL COMMENT '关联表信息ID' AFTER `datasource_id`; - -UPDATE `column_info` c -INNER JOIN `table_info` t - ON t.`datasource_id` = c.`datasource_id` - AND LOWER(t.`table_name`) = LOWER(c.`table_name`) -SET c.`table_id` = t.`id`; - -ALTER TABLE `column_info` - MODIFY COLUMN `table_id` INT NOT NULL COMMENT '关联表信息ID', - DROP INDEX `uk_datasource_table_column`, - DROP INDEX `idx_datasource_table_visible`, - DROP INDEX `idx_datasource_table_visible_column`, - ADD UNIQUE KEY `uk_table_column` (`table_id`, `column_name`), - ADD KEY `idx_datasource_id` (`datasource_id`), - ADD CONSTRAINT `fk_column_info_table` - FOREIGN KEY (`table_id`) REFERENCES `table_info` (`id`) - ON DELETE CASCADE ON UPDATE RESTRICT, - DROP COLUMN `table_name`; - -ALTER TABLE `logical_table_relation` - ADD COLUMN `source_table_id` INT NULL COMMENT '源表信息ID' AFTER `datasource_id`, - ADD COLUMN `target_table_id` INT NULL COMMENT '目标表信息ID' - AFTER `source_column_names_json`; - -UPDATE `logical_table_relation` relation_meta -INNER JOIN `table_info` source_table - ON source_table.`datasource_id` = relation_meta.`datasource_id` - AND LOWER(source_table.`table_name`) = LOWER(relation_meta.`source_table_name`) -INNER JOIN `table_info` target_table - ON target_table.`datasource_id` = relation_meta.`datasource_id` - AND LOWER(target_table.`table_name`) = LOWER(relation_meta.`target_table_name`) -SET relation_meta.`source_table_id` = source_table.`id`, - relation_meta.`target_table_id` = target_table.`id`; - -ALTER TABLE `logical_table_relation` - MODIFY COLUMN `source_table_id` INT NOT NULL COMMENT '源表信息ID', - MODIFY COLUMN `target_table_id` INT NOT NULL COMMENT '目标表信息ID', - DROP INDEX `idx_relation_source_signature`, - DROP INDEX `idx_relation_source_table`, - DROP INDEX `idx_relation_source_enabled`, - DROP INDEX `idx_relation_source_enabled_id`, - DROP INDEX `idx_relation_source_target_id`, - ADD KEY `idx_relation_source_table` (`source_table_id`), - ADD KEY `idx_relation_target_table` (`target_table_id`), - ADD KEY `idx_relation_source_enabled_id` - (`datasource_id`, `source_table_id`, `is_enabled`, `id`), - ADD CONSTRAINT `fk_relation_source_table` - FOREIGN KEY (`source_table_id`) REFERENCES `table_info` (`id`) - ON DELETE CASCADE ON UPDATE RESTRICT, - ADD CONSTRAINT `fk_relation_target_table` - FOREIGN KEY (`target_table_id`) REFERENCES `table_info` (`id`) - ON DELETE CASCADE ON UPDATE RESTRICT, - DROP COLUMN `source_table_name`, - DROP COLUMN `target_table_name`, - DROP COLUMN `source_column_signature`, - DROP COLUMN `target_column_signature`; + +CALL `migrate_primary_key_relations`(); +DROP PROCEDURE `migrate_primary_key_relations`; + +-- Compatibility changes for older metadata schemas follow. They are independently +-- conditional, so they can run after the primary-key migration and on reruns. +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'table_info' + AND COLUMN_NAME = 'data_granularity' + ), + 'SELECT 1', + 'ALTER TABLE `table_info` ADD COLUMN `data_granularity` VARCHAR(500) DEFAULT NULL COMMENT ''数据粒度(一行代表的业务实体或事件)'' AFTER `table_description`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'column_info' + AND COLUMN_NAME = 'semantic_type' + ), + 'SELECT 1', + 'ALTER TABLE `column_info` ADD COLUMN `semantic_type` VARCHAR(32) DEFAULT NULL COMMENT ''语义类型:IDENTIFIER/DIMENSION/MEASURE/TIME/LOCATION'' AFTER `column_description`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'sys_user' + AND COLUMN_NAME = 'is_super_admin' + ), + 'SELECT 1', + 'ALTER TABLE `sys_user` ADD COLUMN `is_super_admin` TINYINT(1) NOT NULL DEFAULT 0 COMMENT ''是否超级管理员:0否,1是'' AFTER `role_id`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +-- sys_user 审计字段:creator_id / create_time / updater_id / update_time / is_deleted +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'sys_user' + AND COLUMN_NAME = 'creator_id' + ), + 'SELECT 1', + 'ALTER TABLE `sys_user` ADD COLUMN `creator_id` BIGINT DEFAULT NULL COMMENT ''创建人ID'' AFTER `status`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'sys_user' + AND COLUMN_NAME = 'create_time' + ), + 'SELECT 1', + 'ALTER TABLE `sys_user` ADD COLUMN `create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT ''创建时间'' AFTER `creator_id`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'sys_user' + AND COLUMN_NAME = 'updater_id' + ), + 'SELECT 1', + 'ALTER TABLE `sys_user` ADD COLUMN `updater_id` BIGINT DEFAULT NULL COMMENT ''修改人ID'' AFTER `create_time`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'sys_user' + AND COLUMN_NAME = 'update_time' + ), + 'SELECT 1', + 'ALTER TABLE `sys_user` ADD COLUMN `update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT ''更新时间'' AFTER `updater_id`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'sys_user' + AND COLUMN_NAME = 'is_deleted' + ), + 'SELECT 1', + 'ALTER TABLE `sys_user` ADD COLUMN `is_deleted` TINYINT(1) NOT NULL DEFAULT 0 COMMENT ''逻辑删除标记:0-未删除,1-已删除'' AFTER `update_time`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +-- 逻辑删除后允许同身份源用户重建:唯一键改为普通索引 +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'sys_user' + AND INDEX_NAME = 'uk_idp' + ), + 'ALTER TABLE `sys_user` DROP INDEX `uk_idp`, ADD KEY `idx_idp` (`idp_type`, `idp_user_id`)', + IF( + EXISTS( + SELECT 1 + FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'sys_user' + AND INDEX_NAME = 'idx_idp' + ), + 'SELECT 1', + 'ALTER TABLE `sys_user` ADD KEY `idx_idp` (`idp_type`, `idp_user_id`)' + ) +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +-- metric_info 审计字段:creator_id / updater_id +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'metric_info' + AND COLUMN_NAME = 'creator_id' + ), + 'SELECT 1', + 'ALTER TABLE `metric_info` ADD COLUMN `creator_id` BIGINT DEFAULT NULL COMMENT ''创建人ID'' AFTER `description`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'metric_info' + AND COLUMN_NAME = 'updater_id' + ), + 'SELECT 1', + 'ALTER TABLE `metric_info` ADD COLUMN `updater_id` BIGINT DEFAULT NULL COMMENT ''修改人ID'' AFTER `create_time`' +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement; + +-- 逻辑删除后保留行以便审计:同 key 软删后允许重新创建,唯一键改为普通索引 +SET @migration_sql = IF( + EXISTS( + SELECT 1 + FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'metric_info' + AND INDEX_NAME = 'uk_datasource_metric_key' + ), + 'ALTER TABLE `metric_info` DROP INDEX `uk_datasource_metric_key`, ADD KEY `idx_datasource_metric_key` (`datasource_id`, `metric_key`)', + IF( + EXISTS( + SELECT 1 + FROM information_schema.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'metric_info' + AND INDEX_NAME = 'idx_datasource_metric_key' + ), + 'SELECT 1', + 'ALTER TABLE `metric_info` ADD KEY `idx_datasource_metric_key` (`datasource_id`, `metric_key`)' + ) +); +PREPARE migration_statement FROM @migration_sql; +EXECUTE migration_statement; +DEALLOCATE PREPARE migration_statement;