diff --git a/crates/dbx-core/src/transfer.rs b/crates/dbx-core/src/transfer.rs index 16f04dbb0..eb9393387 100644 --- a/crates/dbx-core/src/transfer.rs +++ b/crates/dbx-core/src/transfer.rs @@ -1344,6 +1344,15 @@ fn writable_transfer_columns( .collect() } +fn transfer_key_columns(columns: &[db::ColumnInfo], db_type: &DatabaseType) -> Vec { + let uses_unique_key_model = matches!(db_type, DatabaseType::Doris | DatabaseType::StarRocks); + columns + .iter() + .filter(|column| column.is_primary_key || (uses_unique_key_model && column.is_unique)) + .map(|column| column.name.clone()) + .collect() +} + fn dameng_identity_insert_statement(table: &str, schema: &str, enabled: bool) -> String { let full_table = qualified_table(table, schema, &DatabaseType::Dameng, None); format!("SET IDENTITY_INSERT {full_table} {}", if enabled { "ON" } else { "OFF" }) @@ -5609,11 +5618,7 @@ where ) .await?; let col_names = columns.iter().map(|column| column.name.clone()).collect::>(); - let primary_key_columns = columns - .iter() - .filter(|column| column.is_primary_key) - .map(|column| column.name.clone()) - .collect::>(); + let primary_key_columns = transfer_key_columns(&columns, source_db_type); let sql = pagination_sql_with_order( &col_names, table, @@ -5863,8 +5868,7 @@ where let col_names: Vec = writable_columns.iter().map(|c| c.name.clone()).collect(); let col_types: Vec> = writable_columns.iter().map(|c| Some(c.data_type.clone())).collect(); - let primary_key_columns: Vec = - writable_columns.iter().filter(|c| c.is_primary_key).map(|c| c.name.clone()).collect(); + let primary_key_columns = transfer_key_columns(&writable_columns, source_db_type); log::info!("[transfer] {} has {} columns, counting rows...", table, columns.len()); // Fetch source table comment @@ -6146,10 +6150,9 @@ where ) .await .unwrap_or_default(); - let pks: Vec = target_columns - .iter() - .filter(|c| c.is_primary_key && col_names.iter().any(|name| name.eq_ignore_ascii_case(&c.name))) - .map(|c| c.name.clone()) + let pks: Vec = transfer_key_columns(&target_columns, target_db_type) + .into_iter() + .filter(|name| col_names.iter().any(|column_name| column_name.eq_ignore_ascii_case(name))) .collect(); if pks.is_empty() { log::warn!("[transfer] table {} has no primary key, falling back to append", table); @@ -8367,6 +8370,35 @@ mod tests { assert_eq!(sql, "SELECT \"id\", \"name\" FROM \"public\".\"users\" ORDER BY \"id\" LIMIT 100 OFFSET 200"); } + #[test] + fn doris_unique_key_columns_drive_transfer_pagination_order() { + let columns = + vec![db::ColumnInfo { is_unique: true, ..test_column("id", "int") }, test_column("payload", "varchar(64)")]; + let key_columns = transfer_key_columns(&columns, &DatabaseType::Doris); + + assert_eq!(key_columns, vec![String::from("id")]); + assert_eq!( + pagination_sql_with_order( + &[String::from("id"), String::from("payload")], + "events", + "analytics", + &DatabaseType::Doris, + 1000, + 1000, + &key_columns, + None, + ), + "SELECT `id`, `payload` FROM `analytics`.`events` ORDER BY `id` LIMIT 1000 OFFSET 1000" + ); + } + + #[test] + fn mysql_unique_columns_do_not_become_transfer_keys() { + let columns = vec![db::ColumnInfo { is_unique: true, ..test_column("email", "varchar(255)") }]; + + assert!(transfer_key_columns(&columns, &DatabaseType::Mysql).is_empty()); + } + #[test] fn questdb_pagination_uses_stable_primary_key_order() { let sql = pagination_sql_with_order(