From 9adf2b8d00cbe19081592531afe543748c787ee4 Mon Sep 17 00:00:00 2001 From: t8y2 <1156263951@qq.com> Date: Wed, 3 Jun 2026 12:32:51 +0800 Subject: [PATCH] fix(dameng): deduplicate column metadata --- crates/dbx-core/src/schema.rs | 98 +++++++++++++++++-- crates/dbx-web/src/routes/connection.rs | 1 + .../main/java/app/dbx/jdbc/DbxJdbcPlugin.java | 38 +++++-- .../java/app/dbx/jdbc/DbxJdbcPluginTest.java | 56 +++++++++++ 4 files changed, 176 insertions(+), 17 deletions(-) diff --git a/crates/dbx-core/src/schema.rs b/crates/dbx-core/src/schema.rs index ef3824875..db6be893f 100644 --- a/crates/dbx-core/src/schema.rs +++ b/crates/dbx-core/src/schema.rs @@ -458,9 +458,25 @@ fn filter_table_infos(tables: Vec, filter: Option<&str>, limit: O #[cfg(test)] mod tests { use super::{ - clickhouse_metadata_database, duckdb_attach_database, duckdb_list_databases, duckdb_query_tables_in_database, + clickhouse_metadata_database, deduplicate_column_infos, duckdb_attach_database, duckdb_list_databases, + duckdb_query_tables_in_database, }; + fn test_column(name: &str, comment: Option<&str>, is_primary_key: bool) -> super::db::ColumnInfo { + super::db::ColumnInfo { + name: name.to_string(), + data_type: "VARCHAR".to_string(), + is_nullable: true, + column_default: None, + is_primary_key, + extra: None, + comment: comment.map(|value| value.to_string()), + numeric_precision: None, + numeric_scale: None, + character_maximum_length: None, + } + } + #[test] fn duckdb_list_databases_includes_attached_database() { let unique = uuid::Uuid::new_v4(); @@ -505,6 +521,23 @@ mod tests { assert_eq!(clickhouse_metadata_database("testdb", ""), "testdb"); assert_eq!(clickhouse_metadata_database("default", "testdb"), "default"); } + + #[test] + fn deduplicates_columns_and_preserves_later_comment() { + let columns = deduplicate_column_infos(vec![ + test_column("ID", None, false), + test_column("ID", Some("源主键"), true), + test_column("TFBH", Some(""), false), + test_column("TFBH", Some("台账编号"), false), + ]); + + assert_eq!(columns.len(), 2); + assert_eq!(columns[0].name, "ID"); + assert_eq!(columns[0].comment.as_deref(), Some("源主键")); + assert!(columns[0].is_primary_key); + assert_eq!(columns[1].name, "TFBH"); + assert_eq!(columns[1].comment.as_deref(), Some("台账编号")); + } } pub async fn list_objects_core( @@ -729,7 +762,7 @@ pub async fn get_columns_core( let config = config.clone(); let session = session.clone(); drop(connections); - return session + let columns = session .invoke::>( "getColumns", serde_json::json!({ @@ -739,7 +772,8 @@ pub async fn get_columns_core( "table": table, }), ) - .await; + .await?; + return Ok(deduplicate_column_infos(columns)); } if let Some(con) = extract_pool!(&connections, &pool_key, DuckDb) { drop(connections); @@ -755,10 +789,16 @@ pub async fn get_columns_core( if let Some(client) = extract_pool!(&connections, &pool_key, ClickHouse) { drop(connections); return db::clickhouse_driver::get_columns(&client, clickhouse_metadata_database(database, schema), table) - .await; + .await + .map(deduplicate_column_infos); } try_sqlserver!(connections, &pool_key, get_columns, schema, table); - try_agent!(connections, &pool_key, get_columns, database, schema, table); + if let Some(client) = extract_pool!(&connections, &pool_key, Agent) { + drop(connections); + let mut client = client.lock().await; + let columns = client.get_columns::>(database, schema, table).await?; + return Ok(deduplicate_column_infos(columns)); + } } let connections = state.connections.read().await; @@ -767,13 +807,57 @@ pub async fn get_columns_core( match pool { PoolKind::Mysql(p, mode) => { dispatch_mysql!(p, mode, db::mysql::get_columns, db::ob_oracle::get_columns, database, table) + .map(deduplicate_column_infos) } - PoolKind::Postgres(p) => db::postgres::get_columns(p, schema, table).await, - PoolKind::Sqlite(p) => db::sqlite::get_columns(p, schema, table).await, + PoolKind::Postgres(p) => db::postgres::get_columns(p, schema, table).await.map(deduplicate_column_infos), + PoolKind::Sqlite(p) => db::sqlite::get_columns(p, schema, table).await.map(deduplicate_column_infos), _ => Ok(vec![]), } } +fn deduplicate_column_infos(columns: Vec) -> Vec { + let mut result: Vec = Vec::with_capacity(columns.len()); + for column in columns { + if let Some(existing) = result.iter_mut().find(|existing| existing.name == column.name) { + existing.is_primary_key |= column.is_primary_key; + existing.is_nullable &= column.is_nullable; + merge_optional_string(&mut existing.column_default, column.column_default); + merge_optional_string(&mut existing.extra, column.extra); + merge_optional_string(&mut existing.comment, column.comment); + if existing.numeric_precision.is_none() { + existing.numeric_precision = column.numeric_precision; + } + if existing.numeric_scale.is_none() { + existing.numeric_scale = column.numeric_scale; + } + if existing.character_maximum_length.is_none() { + existing.character_maximum_length = column.character_maximum_length; + } + if existing.data_type.trim().is_empty() && !column.data_type.trim().is_empty() { + existing.data_type = column.data_type; + } + } else { + result.push(column); + } + } + result +} + +fn merge_optional_string(target: &mut Option, candidate: Option) { + let Some(candidate) = candidate else { + return; + }; + if candidate.trim().is_empty() { + if target.is_none() { + *target = Some(candidate); + } + return; + } + if target.as_ref().map_or(true, |value| value.trim().is_empty()) { + *target = Some(candidate); + } +} + pub async fn list_indexes_core( state: &AppState, connection_id: &str, diff --git a/crates/dbx-web/src/routes/connection.rs b/crates/dbx-web/src/routes/connection.rs index 4e07f88a0..a59a6d5aa 100644 --- a/crates/dbx-web/src/routes/connection.rs +++ b/crates/dbx-web/src/routes/connection.rs @@ -209,6 +209,7 @@ mod tests { sse_channels: RwLock::new(HashMap::new()), sql_file_executions: RwLock::new(HashMap::new()), login_rate_limit: Mutex::new(LoginRateLimit { fail_count: 0, locked_until: None }), + export_files: RwLock::new(HashMap::new()), }); (state, dir) } diff --git a/plugins/jdbc/src/main/java/app/dbx/jdbc/DbxJdbcPlugin.java b/plugins/jdbc/src/main/java/app/dbx/jdbc/DbxJdbcPlugin.java index 95b26cf88..934f37a9f 100644 --- a/plugins/jdbc/src/main/java/app/dbx/jdbc/DbxJdbcPlugin.java +++ b/plugins/jdbc/src/main/java/app/dbx/jdbc/DbxJdbcPlugin.java @@ -554,18 +554,16 @@ public final class DbxJdbcPlugin { try (ResultSet rs = meta.getColumns(catalog, schema, table, "%")) { while (rs.next()) { String name = rs.getString("COLUMN_NAME"); - ObjectNode item = MAPPER.createObjectNode(); - item.put("name", name); + ObjectNode item = columnNode(result, name); item.put("data_type", rs.getString("TYPE_NAME")); item.put("is_nullable", rs.getInt("NULLABLE") != DatabaseMetaData.columnNoNulls); - putNullable(item, "column_default", rs.getString("COLUMN_DEF")); + putNullablePreferValue(item, "column_default", rs.getString("COLUMN_DEF")); item.put("is_primary_key", primaryKeys.contains(name)); item.putNull("extra"); - putNullable(item, "comment", rs.getString("REMARKS")); + putNullablePreferValue(item, "comment", rs.getString("REMARKS")); putNullableInt(item, "numeric_precision", rs.getObject("COLUMN_SIZE")); putNullableInt(item, "numeric_scale", rs.getObject("DECIMAL_DIGITS")); putNullableInt(item, "character_maximum_length", rs.getObject("COLUMN_SIZE")); - result.add(item); } } } @@ -867,18 +865,16 @@ public final class DbxJdbcPlugin { try (ResultSet rs = ps.executeQuery()) { while (rs.next()) { String name = rs.getString("column_name"); - ObjectNode item = MAPPER.createObjectNode(); - item.put("name", name); + ObjectNode item = columnNode(result, name); item.put("data_type", rs.getString("data_type")); item.put("is_nullable", !"N".equals(rs.getString("nullable"))); - putNullable(item, "column_default", rs.getString("data_default")); + putNullablePreferValue(item, "column_default", rs.getString("data_default")); item.put("is_primary_key", pks.contains(name)); item.putNull("extra"); - putNullable(item, "comment", rs.getString("comments")); + putNullablePreferValue(item, "comment", rs.getString("comments")); putNullableInt(item, "numeric_precision", rs.getObject("data_precision")); putNullableInt(item, "numeric_scale", rs.getObject("data_scale")); putNullableInt(item, "character_maximum_length", rs.getObject("char_length")); - result.add(item); } } } @@ -955,6 +951,28 @@ public final class DbxJdbcPlugin { } } + private static ObjectNode columnNode(ArrayNode result, String name) { + for (JsonNode node : result) { + if (name.equals(node.path("name").asText()) && node instanceof ObjectNode objectNode) { + return objectNode; + } + } + ObjectNode item = MAPPER.createObjectNode(); + item.put("name", name); + result.add(item); + return item; + } + + private static void putNullablePreferValue(ObjectNode node, String field, String value) { + if (value == null || value.isBlank()) { + if (!node.has(field)) { + node.putNull(field); + } + return; + } + node.put(field, value); + } + private static void putNullableInt(ObjectNode node, String field, Object value) { if (value instanceof Number number) { node.put(field, number.intValue()); diff --git a/plugins/jdbc/src/test/java/app/dbx/jdbc/DbxJdbcPluginTest.java b/plugins/jdbc/src/test/java/app/dbx/jdbc/DbxJdbcPluginTest.java index 504605b7d..2581133d7 100644 --- a/plugins/jdbc/src/test/java/app/dbx/jdbc/DbxJdbcPluginTest.java +++ b/plugins/jdbc/src/test/java/app/dbx/jdbc/DbxJdbcPluginTest.java @@ -311,6 +311,62 @@ final class DbxJdbcPluginTest { } } + @Test + void oracleGetColumnsMergesDuplicateMetadataRowsAndKeepsComments() throws Exception { + Method method = DbxJdbcPlugin.class.getDeclaredMethod("oracleGetColumns", Connection.class, String.class, String.class); + method.setAccessible(true); + + try (Connection conn = DriverManager.getConnection("jdbc:h2:mem:dbx_oracle_duplicate_columns;DB_CLOSE_DELAY=-1", "sa", "")) { + conn.createStatement().execute( + "CREATE TABLE all_tab_comments (owner VARCHAR(64), table_name VARCHAR(64), table_type VARCHAR(16))" + ); + conn.createStatement().execute( + "CREATE TABLE all_tab_columns (" + + "owner VARCHAR(64), table_name VARCHAR(64), column_name VARCHAR(64), data_type VARCHAR(32), " + + "nullable VARCHAR(1), data_default VARCHAR(64), data_precision INT, data_scale INT, char_length INT, column_id INT)" + ); + conn.createStatement().execute( + "CREATE TABLE all_col_comments (owner VARCHAR(64), table_name VARCHAR(64), column_name VARCHAR(64), comments VARCHAR(128))" + ); + conn.createStatement().execute( + "CREATE TABLE all_constraints (owner VARCHAR(64), table_name VARCHAR(64), constraint_name VARCHAR(64), constraint_type VARCHAR(1))" + ); + conn.createStatement().execute( + "CREATE TABLE all_cons_columns (owner VARCHAR(64), table_name VARCHAR(64), constraint_name VARCHAR(64), column_name VARCHAR(64))" + ); + conn.createStatement().execute( + "INSERT INTO all_tab_comments(owner, table_name, table_type) VALUES ('SYSDBA', 'F02_TFBH', 'TABLE')" + ); + conn.createStatement().execute( + "INSERT INTO all_tab_columns(owner, table_name, column_name, data_type, nullable, data_default, data_precision, data_scale, char_length, column_id) " + + "VALUES ('SYSDBA', 'F02_TFBH', 'ID', 'INT', 'N', NULL, 10, 0, NULL, 1), " + + "('SYSDBA', 'F02_TFBH', 'TFBH', 'VARCHAR', 'Y', NULL, NULL, NULL, 8, 2)" + ); + conn.createStatement().execute( + "INSERT INTO all_col_comments(owner, table_name, column_name, comments) VALUES " + + "('SYSDBA', 'F02_TFBH', 'ID', NULL), " + + "('SYSDBA', 'F02_TFBH', 'ID', '源主键'), " + + "('SYSDBA', 'F02_TFBH', 'TFBH', NULL), " + + "('SYSDBA', 'F02_TFBH', 'TFBH', '台账编号')" + ); + conn.createStatement().execute( + "INSERT INTO all_constraints(owner, table_name, constraint_name, constraint_type) VALUES ('SYSDBA', 'F02_TFBH', 'PK_F02_TFBH', 'P')" + ); + conn.createStatement().execute( + "INSERT INTO all_cons_columns(owner, table_name, constraint_name, column_name) VALUES ('SYSDBA', 'F02_TFBH', 'PK_F02_TFBH', 'ID')" + ); + + JsonNode columns = MAPPER.valueToTree(method.invoke(null, conn, "SYSDBA", "F02_TFBH")); + + assertEquals(2, columns.size()); + assertEquals("ID", columns.path(0).path("name").asText()); + assertEquals("源主键", columns.path(0).path("comment").asText()); + assertEquals(true, columns.path(0).path("is_primary_key").asBoolean()); + assertEquals("TFBH", columns.path(1).path("name").asText()); + assertEquals("台账编号", columns.path(1).path("comment").asText()); + } + } + private static void createPeopleTable() throws Exception { request("executeQuery", """ {