fix(dameng): deduplicate column metadata

This commit is contained in:
t8y2 2026-06-03 12:32:51 +08:00
parent e67894bea7
commit 9adf2b8d00
4 changed files with 176 additions and 17 deletions

View File

@ -458,9 +458,25 @@ fn filter_table_infos(tables: Vec<db::TableInfo>, 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::<Vec<db::ColumnInfo>>(
"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::<Vec<db::ColumnInfo>>(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<db::ColumnInfo>) -> Vec<db::ColumnInfo> {
let mut result: Vec<db::ColumnInfo> = 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<String>, candidate: Option<String>) {
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,

View File

@ -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)
}

View File

@ -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());

View File

@ -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", """
{