Skip to content

Commit b94a1fa

Browse files
authored
feat(catalog): support column-level alter table (#370)
1 parent 5617a77 commit b94a1fa

15 files changed

Lines changed: 1932 additions & 193 deletions

File tree

crates/integrations/datafusion/src/sql_context.rs

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -3182,7 +3182,7 @@ mod tests {
31823182
assert_eq!(identifier.object(), "t1");
31833183
assert_eq!(changes.len(), 1);
31843184
assert!(
3185-
matches!(&changes[0], SchemaChange::AddColumn { field_name, .. } if field_name == "age")
3185+
matches!(&changes[0], SchemaChange::AddColumn { field_names, .. } if field_names.first().map(String::as_str) == Some("age"))
31863186
);
31873187
} else {
31883188
panic!("expected AlterTable call");
@@ -3206,10 +3206,10 @@ mod tests {
32063206
assert!(matches!(
32073207
&changes[0],
32083208
SchemaChange::AddColumn {
3209-
field_name,
3209+
field_names,
32103210
data_type,
32113211
..
3212-
} if field_name == "payload" && matches!(data_type, PaimonDataType::Blob(_))
3212+
} if field_names.first().map(String::as_str) == Some("payload") && matches!(data_type, PaimonDataType::Blob(_))
32133213
));
32143214
} else {
32153215
panic!("expected AlterTable call");
@@ -3231,7 +3231,7 @@ mod tests {
32313231
if let CatalogCall::AlterTable { changes, .. } = &calls[0] {
32323232
assert_eq!(changes.len(), 1);
32333233
assert!(
3234-
matches!(&changes[0], SchemaChange::DropColumn { field_name } if field_name == "age")
3234+
matches!(&changes[0], SchemaChange::DropColumn { field_names } if field_names.first().map(String::as_str) == Some("age"))
32353235
);
32363236
} else {
32373237
panic!("expected AlterTable call");
@@ -3254,8 +3254,8 @@ mod tests {
32543254
assert_eq!(changes.len(), 1);
32553255
assert!(matches!(
32563256
&changes[0],
3257-
SchemaChange::RenameColumn { field_name, new_name }
3258-
if field_name == "old_name" && new_name == "new_name"
3257+
SchemaChange::RenameColumn { field_names, new_name }
3258+
if field_names.first().map(String::as_str) == Some("old_name") && new_name == "new_name"
32593259
));
32603260
} else {
32613261
panic!("expected AlterTable call");

crates/integrations/datafusion/tests/sql_context_tests.rs

Lines changed: 10 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -480,23 +480,18 @@ async fn test_alter_table_add_column() {
480480
.await
481481
.unwrap();
482482

483-
// ALTER TABLE is not yet implemented in FileSystemCatalog, so we expect an error
484-
let result = sql_context
483+
sql_context
485484
.sql("ALTER TABLE paimon.mydb.alter_test ADD COLUMN age INT")
486-
.await;
485+
.await
486+
.expect("ALTER TABLE ADD COLUMN should succeed");
487487

488-
// FileSystemCatalog does not support AddColumn schema change yet
489-
assert!(
490-
result.is_err(),
491-
"ALTER TABLE ADD COLUMN should fail because AddColumn is not yet supported"
492-
);
493-
let err_msg = result.unwrap_err().to_string();
494-
assert!(
495-
err_msg.contains("not yet implemented")
496-
|| err_msg.contains("Unsupported")
497-
|| err_msg.contains("not yet supported"),
498-
"Error should indicate alter_table is not implemented, got: {err_msg}"
499-
);
488+
// The new column is appended to the table schema.
489+
let table = catalog
490+
.get_table(&Identifier::new("mydb", "alter_test"))
491+
.await
492+
.unwrap();
493+
let names: Vec<&str> = table.schema().fields().iter().map(|f| f.name()).collect();
494+
assert_eq!(names, vec!["id", "name", "age"]);
500495
}
501496

502497
#[tokio::test]

crates/paimon/src/api/api_request.rs

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,10 @@
2222
use serde::{Deserialize, Serialize};
2323
use std::collections::HashMap;
2424

25-
use crate::{catalog::Identifier, spec::Schema};
25+
use crate::{
26+
catalog::Identifier,
27+
spec::{Schema, SchemaChange},
28+
};
2629

2730
/// Request to create a new database.
2831
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
@@ -95,6 +98,23 @@ impl CreateTableRequest {
9598
}
9699
}
97100

101+
/// Request to alter a table's schema.
102+
///
103+
/// Wire-compatible with Java Paimon's `AlterTableRequest` (`{"changes": [...]}`).
104+
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
105+
#[serde(rename_all = "camelCase")]
106+
pub struct AlterTableRequest {
107+
/// The ordered list of schema changes to apply.
108+
pub changes: Vec<SchemaChange>,
109+
}
110+
111+
impl AlterTableRequest {
112+
/// Create a new AlterTableRequest.
113+
pub fn new(changes: Vec<SchemaChange>) -> Self {
114+
Self { changes }
115+
}
116+
}
117+
98118
#[cfg(test)]
99119
mod tests {
100120
use super::*;

crates/paimon/src/api/mod.rs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,8 @@ mod api_response;
3131

3232
// Re-export request types
3333
pub use api_request::{
34-
AlterDatabaseRequest, CreateDatabaseRequest, CreateTableRequest, RenameTableRequest,
34+
AlterDatabaseRequest, AlterTableRequest, CreateDatabaseRequest, CreateTableRequest,
35+
RenameTableRequest,
3536
};
3637

3738
// Re-export response types

crates/paimon/src/api/rest_api.rs

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,11 +25,12 @@ use std::collections::HashMap;
2525
use crate::api::rest_client::HttpClient;
2626
use crate::catalog::Identifier;
2727
use crate::common::{CatalogOptions, Options};
28-
use crate::spec::{Partition, PartitionStatistics, Schema, Snapshot};
28+
use crate::spec::{Partition, PartitionStatistics, Schema, SchemaChange, Snapshot};
2929
use crate::Result;
3030

3131
use super::api_request::{
32-
AlterDatabaseRequest, CreateDatabaseRequest, CreateTableRequest, RenameTableRequest,
32+
AlterDatabaseRequest, AlterTableRequest, CreateDatabaseRequest, CreateTableRequest,
33+
RenameTableRequest,
3334
};
3435
use super::api_response::{
3536
ConfigResponse, GetDatabaseResponse, GetTableResponse, ListDatabasesResponse,
@@ -343,6 +344,21 @@ impl RESTApi {
343344
Ok(())
344345
}
345346

347+
/// Alter a table's schema by applying a list of schema changes.
348+
pub async fn alter_table(
349+
&self,
350+
identifier: &Identifier,
351+
changes: Vec<SchemaChange>,
352+
) -> Result<()> {
353+
let database = identifier.database();
354+
let table = identifier.object();
355+
validate_non_empty_multi(&[(database, "database name"), (table, "table name")])?;
356+
let path = self.resource_paths.table(database, table);
357+
let request = AlterTableRequest::new(changes);
358+
let _resp: serde_json::Value = self.client.post(&path, &request).await?;
359+
Ok(())
360+
}
361+
346362
/// Get table information.
347363
pub async fn get_table(&self, identifier: &Identifier) -> Result<GetTableResponse> {
348364
let database = identifier.database();

0 commit comments

Comments
 (0)