From 87c9b3f86b9e29ade9e611243d71db6f19cbe109 Mon Sep 17 00:00:00 2001 From: mengnankkkk Date: Sun, 9 Aug 2026 15:18:42 +0800 Subject: [PATCH 1/2] fix(semantic): deduplicate synced semantic metadata - avoid inserting duplicate table and column semantics during physical sync - deduplicate table, column, and relation reads by canonical semantic keys - prefer manually described or disabled semantics, then latest updated record --- .../sync/SemanticSyncApplyService.java | 49 +++++++++++++------ .../mapper/ColumnSemanticInfoMapper.xml | 43 +++++++++++++--- .../mapper/LogicalTableRelationMapper.xml | 46 ++++++++++++++--- .../main/resources/mapper/TableInfoMapper.xml | 47 +++++++++++++++--- 4 files changed, 152 insertions(+), 33 deletions(-) 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 1a075c0..232e214 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 @@ -61,8 +61,8 @@ public List applyTableSync( Map> columnsByTableName = loadSemanticColumnsByTable(datasourceId, selectedTableNames); - // 表和列分别执行一次批量 upsert,不按表或列逐条写库。 - batchUpsertPresentSchema(datasourceId, presentTables); + batchSavePresentSchema( + datasourceId, presentTables, existingTableIndex, columnsByTableName); // 基于写入前快照生成统计结果,同时收集需要批量标记缺失的列 ID。 List missingTableNames = new ArrayList<>(); @@ -151,25 +151,46 @@ public List refreshPhysicalStatus( return results; } - private void batchUpsertPresentSchema( - Integer datasourceId, List presentTables) { + private void batchSavePresentSchema( + Integer datasourceId, + List presentTables, + Map existingTableIndex, + Map> columnsByTableName) { if (presentTables.isEmpty()) { return; } - tableInfoMapper.batchUpsertPhysicalCache( - presentTables.stream() - .map(table -> buildPhysicalTableInfo(datasourceId, table)) - .toList()); - List physicalColumns = new ArrayList<>(); + List newTables = new ArrayList<>(); + List newColumns = new ArrayList<>(); for (TableSyncSource table : presentTables) { + TableInfo tableInfo = buildPhysicalTableInfo(datasourceId, table); + TableInfo existingTable = existingTableIndex.get(table.tableName()); + if (existingTable == null) { + newTables.add(tableInfo); + } else { + tableInfo.setId(existingTable.getId()); + tableInfoMapper.updatePhysicalCacheFields(tableInfo); + } + + Map existingColumnIndex = + loadColumnIndex(columnsByTableName.getOrDefault(table.tableName(), List.of())); for (ColumnSyncSource column : table.columns()) { - physicalColumns.add( - buildPhysicalColumnInfo(datasourceId, table.tableName(), column)); + ColumnInfo columnInfo = + buildPhysicalColumnInfo(datasourceId, table.tableName(), column); + ColumnInfo existingColumn = existingColumnIndex.get(column.columnName()); + if (existingColumn == null) { + newColumns.add(columnInfo); + } else { + columnInfo.setId(existingColumn.getId()); + columnSemanticInfoMapper.updatePhysicalCacheFields(columnInfo); + } } } - if (!physicalColumns.isEmpty()) { - columnSemanticInfoMapper.batchUpsertPhysicalCache(physicalColumns); + if (!newTables.isEmpty()) { + tableInfoMapper.batchUpsertPhysicalCache(newTables); + } + if (!newColumns.isEmpty()) { + columnSemanticInfoMapper.batchUpsertPhysicalCache(newColumns); } } @@ -287,7 +308,7 @@ private Map> loadSemanticColumnsByTable( private Map loadColumnIndex(List columns) { Map index = new LinkedHashMap<>(); for (ColumnInfo column : columns) { - index.put( + index.putIfAbsent( SemanticUtils.normalizeObjectName( column.getColumnName(), "Missing semantic columnName."), column); diff --git a/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml b/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml index 87a2302..580b667 100644 --- a/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/ColumnSemanticInfoMapper.xml @@ -17,21 +17,49 @@ + + SELECT * + FROM ( + SELECT + column_info.*, + ROW_NUMBER() OVER ( + PARTITION BY datasource_id, LOWER(table_name), LOWER(column_name) + ORDER BY + CASE + WHEN column_description IS NOT NULL AND column_description <> '' THEN 0 + WHEN is_visible = 0 THEN 0 + ELSE 1 + END ASC, + update_time DESC, + create_time DESC, + id DESC + ) AS row_num + FROM column_info + ) deduplicated_columns + WHERE row_num = 1 + + diff --git a/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml b/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml index 21d2efd..8823050 100644 --- a/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/LogicalTableRelationMapper.xml @@ -18,25 +18,56 @@ + + SELECT * + FROM ( + SELECT + logical_table_relation.*, + ROW_NUMBER() OVER ( + PARTITION BY + datasource_id, + LOWER(source_table_name), + source_column_signature + ORDER BY + CASE + WHEN description IS NOT NULL AND description <> '' THEN 0 + WHEN is_enabled = 0 THEN 0 + ELSE 1 + END ASC, + update_time DESC, + create_time DESC, + id DESC + ) AS row_num + FROM logical_table_relation + ) deduplicated_relations + WHERE row_num = 1 + + diff --git a/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml b/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml index eea4b7b..5a1c021 100644 --- a/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml +++ b/data-agent-backend/src/main/resources/mapper/TableInfoMapper.xml @@ -15,14 +15,40 @@ + + SELECT * + FROM ( + SELECT + table_info.*, + ROW_NUMBER() OVER ( + PARTITION BY datasource_id, LOWER(table_name) + ORDER BY + CASE + WHEN table_description IS NOT NULL AND table_description <> '' THEN 0 + WHEN is_visible = 0 THEN 0 + ELSE 1 + END ASC, + update_time DESC, + create_time DESC, + id DESC + ) AS row_num + FROM table_info + ) deduplicated_tables + WHERE row_num = 1 + + - SELECT * FROM table_info + SELECT * FROM ( + + ) table_info WHERE datasource_id = #{datasourceId} AND domain IN @@ -132,7 +163,9 @@ From 83cd1530c9a62457a10e4be6a2ac8857d8aa824c Mon Sep 17 00:00:00 2001 From: mengnankkkk Date: Sun, 9 Aug 2026 15:24:32 +0800 Subject: [PATCH 2/2] fix: ci --- .../service/semantic/sync/SemanticSyncApplyService.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) 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 232e214..0c8b261 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 @@ -61,8 +61,7 @@ public List applyTableSync( Map> columnsByTableName = loadSemanticColumnsByTable(datasourceId, selectedTableNames); - batchSavePresentSchema( - datasourceId, presentTables, existingTableIndex, columnsByTableName); + batchSavePresentSchema(datasourceId, presentTables, existingTableIndex, columnsByTableName); // 基于写入前快照生成统计结果,同时收集需要批量标记缺失的列 ID。 List missingTableNames = new ArrayList<>();