This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new 35a166db feat(datafusion): support database SQL statements (#627)
35a166db is described below

commit 35a166dbfcf6748978a9f7d2db660a03612cea9f
Author: shyjsarah <[email protected]>
AuthorDate: Thu Jul 30 09:11:58 2026 +0800

    feat(datafusion): support database SQL statements (#627)
---
 crates/integrations/datafusion/src/sql_context.rs  | 178 ++++++++++++++++++++-
 .../datafusion/tests/sql_context_tests.rs          | 109 ++++++++++++-
 docs/src/sql.md                                    |  11 +-
 3 files changed, 287 insertions(+), 11 deletions(-)

diff --git a/crates/integrations/datafusion/src/sql_context.rs 
b/crates/integrations/datafusion/src/sql_context.rs
index d5a92f09..1803892b 100644
--- a/crates/integrations/datafusion/src/sql_context.rs
+++ b/crates/integrations/datafusion/src/sql_context.rs
@@ -18,12 +18,15 @@
 //! SQL support for Paimon tables.
 //!
 //! DataFusion does not natively support all SQL statements needed by Paimon.
-//! This module provides [`SQLContext`] which intercepts CREATE TABLE,
-//! ALTER TABLE, MERGE INTO, UPDATE and other SQL, translates them to Paimon
-//! catalog operations, and delegates everything else (SELECT, CREATE/DROP
-//! SCHEMA, DROP TABLE, etc.) to the underlying [`SessionContext`].
+//! This module provides [`SQLContext`] which intercepts database and table
+//! statements needed by Paimon, translates them to catalog operations, and
+//! delegates everything else to the underlying [`SessionContext`].
 //!
 //! Supported DDL:
+//! - `SHOW DATABASES`
+//! - `CREATE DATABASE [IF NOT EXISTS] [catalog.]db`
+//! - `DROP DATABASE [IF EXISTS] [catalog.]db [CASCADE]`
+//! - `USE [catalog.]db`
 //! - `CREATE TABLE db.t (col TYPE, ..., PRIMARY KEY (col, ...)) [PARTITIONED 
BY (col, ...)] [WITH ('key' = 'val')]`
 //! - `ALTER TABLE db.t ADD COLUMN col TYPE`
 //! - `ALTER TABLE db.t DROP COLUMN col`
@@ -61,7 +64,7 @@ use datafusion::sql::sqlparser::ast::{
     ColumnOption, CreateFunction, CreateFunctionBody, CreateTable, 
CreateTableOptions, CreateView,
     Delete, Expr as SqlExpr, FromTable, FunctionBehavior, FunctionReturnType, 
Insert, Merge,
     ObjectName, ObjectType, RenameTableNameKind, Reset, ResetStatement, Set, 
ShowCreateObject,
-    SqlOption, Statement, TableFactor, TableObject, Truncate, Update, Value as 
SqlValue,
+    SqlOption, Statement, TableFactor, TableObject, Truncate, Update, Use, 
Value as SqlValue,
 };
 use datafusion::sql::sqlparser::dialect::GenericDialect;
 use datafusion::sql::sqlparser::keywords::Keyword;
@@ -218,6 +221,10 @@ impl SQLContext {
 
     /// Sets the current catalog for unqualified table references.
     pub async fn set_current_catalog(&mut self, catalog_name: impl 
Into<String>) -> DFResult<()> {
+        self.set_current_catalog_inner(catalog_name).await
+    }
+
+    async fn set_current_catalog_inner(&self, catalog_name: impl Into<String>) 
-> DFResult<()> {
         let catalog_name = catalog_name.into();
         if !self.catalogs.contains_key(&catalog_name) {
             return Err(DataFusionError::Plan(format!(
@@ -357,8 +364,8 @@ impl SQLContext {
         &self.dynamic_options
     }
 
-    /// Execute a SQL statement. ALTER TABLE is handled by Paimon directly;
-    /// everything else is delegated to DataFusion.
+    /// Execute a SQL statement. Paimon database and table extensions are 
handled
+    /// directly; everything else is delegated to DataFusion.
     pub async fn sql(&self, sql: &str) -> DFResult<DataFrame> {
         let is_create_table = looks_like_create_table(sql);
         let (rewritten_sql, partition_keys) = if is_create_table {
@@ -380,6 +387,48 @@ impl SQLContext {
         }
 
         match &statements[0] {
+            Statement::ShowDatabases {
+                terse,
+                history,
+                show_options,
+            } => {
+                if *terse
+                    || *history
+                    || show_options.show_in.is_some()
+                    || show_options.starts_with.is_some()
+                    || show_options.limit.is_some()
+                    || show_options.limit_from.is_some()
+                    || show_options.filter_position.is_some()
+                {
+                    return Err(DataFusionError::Plan(
+                        "SHOW DATABASES options are not supported".to_string(),
+                    ));
+                }
+                self.handle_show_databases().await
+            }
+            Statement::CreateDatabase {
+                db_name,
+                if_not_exists,
+                location,
+                managed_location,
+                clone,
+                default_charset,
+                default_collation,
+                ..
+            } => {
+                if location.is_some()
+                    || managed_location.is_some()
+                    || clone.is_some()
+                    || default_charset.is_some()
+                    || default_collation.is_some()
+                {
+                    return Err(DataFusionError::Plan(
+                        "CREATE DATABASE options are not 
supported".to_string(),
+                    ));
+                }
+                self.handle_create_database(db_name, *if_not_exists).await
+            }
+            Statement::Use(Use::Object(name)) => 
self.handle_use_database(name).await,
             Statement::CreateTable(create_table) => {
                 if create_table.temporary {
                     self.handle_create_temp_table(create_table).await
@@ -471,6 +520,28 @@ impl SQLContext {
                     self.ctx.sql(sql).await
                 }
             }
+            Statement::Drop {
+                object_type: ObjectType::Database,
+                if_exists,
+                names,
+                cascade,
+                restrict,
+                purge,
+                temporary,
+                table,
+            } => {
+                let [name] = names.as_slice() else {
+                    return Err(DataFusionError::Plan(
+                        "DROP DATABASE requires exactly one 
database".to_string(),
+                    ));
+                };
+                if *restrict || *purge || *temporary || table.is_some() {
+                    return Err(DataFusionError::Plan(
+                        "DROP DATABASE options are not supported".to_string(),
+                    ));
+                }
+                self.handle_drop_database(name, *if_exists, *cascade).await
+            }
             Statement::Drop {
                 object_type,
                 if_exists,
@@ -1027,6 +1098,66 @@ impl SQLContext {
         self.ctx.read_batch(batch)
     }
 
+    async fn handle_show_databases(&self) -> DFResult<DataFrame> {
+        let catalog = self.current_catalog()?;
+        let mut databases = catalog
+            .list_databases()
+            .await
+            .map_err(to_datafusion_error)?;
+        databases.sort_unstable();
+
+        let schema = Arc::new(Schema::new(vec![Field::new(
+            "database_name",
+            ArrowDataType::Utf8,
+            false,
+        )]));
+        let batch = RecordBatch::try_new(schema, 
vec![Arc::new(StringArray::from(databases))])?;
+        self.ctx.read_batch(batch)
+    }
+
+    async fn handle_create_database(
+        &self,
+        name: &ObjectName,
+        if_not_exists: bool,
+    ) -> DFResult<DataFrame> {
+        let (catalog, _, database) = self.resolve_catalog_and_database(name)?;
+        catalog
+            .create_database(&database, if_not_exists, Default::default())
+            .await
+            .map_err(to_datafusion_error)?;
+        ok_result(&self.ctx)
+    }
+
+    async fn handle_drop_database(
+        &self,
+        name: &ObjectName,
+        if_exists: bool,
+        cascade: bool,
+    ) -> DFResult<DataFrame> {
+        let (catalog, _, database) = self.resolve_catalog_and_database(name)?;
+        catalog
+            .drop_database(&database, if_exists, cascade)
+            .await
+            .map_err(to_datafusion_error)?;
+        ok_result(&self.ctx)
+    }
+
+    async fn handle_use_database(&self, name: &ObjectName) -> 
DFResult<DataFrame> {
+        let (catalog, catalog_name, database) = 
self.resolve_catalog_and_database(name)?;
+        if database.contains('\'') {
+            return Err(DataFusionError::Plan(
+                "Database name must not contain single quotes".to_string(),
+            ));
+        }
+        catalog
+            .get_database(&database)
+            .await
+            .map_err(to_datafusion_error)?;
+        self.set_current_catalog_inner(catalog_name).await?;
+        self.set_current_database(&database).await?;
+        ok_result(&self.ctx)
+    }
+
     async fn handle_alter_table(
         &self,
         catalog: &Arc<dyn Catalog>,
@@ -1721,6 +1852,39 @@ impl SQLContext {
         })
     }
 
+    fn resolve_catalog_and_database(
+        &self,
+        name: &ObjectName,
+    ) -> DFResult<(Arc<dyn Catalog>, String, String)> {
+        let parts = name
+            .0
+            .iter()
+            .map(|part| {
+                part.as_ident()
+                    .map(|identifier| identifier.value.clone())
+                    .ok_or_else(|| {
+                        DataFusionError::Plan(format!("Invalid database 
reference: {name}"))
+                    })
+            })
+            .collect::<DFResult<Vec<_>>>()?;
+        match parts.as_slice() {
+            [database] => Ok((
+                self.current_catalog()?,
+                self.current_catalog_name(),
+                database.clone(),
+            )),
+            [catalog_name, database] => {
+                let catalog = 
self.catalogs.get(catalog_name).cloned().ok_or_else(|| {
+                    DataFusionError::Plan(format!("Unknown catalog 
'{catalog_name}'"))
+                })?;
+                Ok((catalog, catalog_name.clone(), database.clone()))
+            }
+            _ => Err(DataFusionError::Plan(format!(
+                "Invalid database reference: {name}"
+            ))),
+        }
+    }
+
     /// Check whether a TableReference targets a registered Paimon catalog.
     fn is_paimon_catalog_ref(&self, table_ref: &TableReference) -> bool {
         let catalog_name = match table_ref {
diff --git a/crates/integrations/datafusion/tests/sql_context_tests.rs 
b/crates/integrations/datafusion/tests/sql_context_tests.rs
index f42acbda..308b806c 100644
--- a/crates/integrations/datafusion/tests/sql_context_tests.rs
+++ b/crates/integrations/datafusion/tests/sql_context_tests.rs
@@ -719,7 +719,114 @@ async fn 
test_branch_partitions_system_table_reads_branch_snapshot() {
     assert!(catalog.take_partition_identifiers().is_empty());
 }
 
-// ======================= CREATE / DROP SCHEMA =======================
+// ======================= DATABASE / SCHEMA =======================
+
+#[tokio::test]
+async fn test_database_statements() {
+    let (_tmp, catalog) = create_test_env();
+    let sql_context = create_sql_context(catalog.clone()).await;
+
+    sql_context.sql("CREATE DATABASE analytics").await.unwrap();
+    sql_context
+        .sql("CREATE DATABASE IF NOT EXISTS analytics")
+        .await
+        .unwrap();
+
+    let databases = collect_string_column(&sql_context, "SHOW DATABASES", 
"database_name").await;
+    assert_eq!(databases, vec!["analytics", "default"]);
+
+    sql_context
+        .sql("CREATE TABLE analytics.events (id INT)")
+        .await
+        .unwrap();
+    sql_context
+        .sql("DROP DATABASE analytics CASCADE")
+        .await
+        .unwrap();
+    sql_context
+        .sql("DROP DATABASE IF EXISTS analytics")
+        .await
+        .unwrap();
+
+    assert!(!catalog
+        .list_databases()
+        .await
+        .unwrap()
+        .contains(&"analytics".to_string()));
+}
+
+#[tokio::test]
+async fn test_database_statements_reject_unsupported_options() {
+    let (_tmp, catalog) = create_test_env();
+    let sql_context = create_sql_context(catalog.clone()).await;
+
+    assert_sql_error_contains(
+        &sql_context,
+        "SHOW DATABASES LIKE 'a%'",
+        "SHOW DATABASES options are not supported",
+    )
+    .await;
+    assert_sql_error_contains(
+        &sql_context,
+        "CREATE DATABASE analytics LOCATION 'file:///tmp/analytics'",
+        "CREATE DATABASE options are not supported",
+    )
+    .await;
+
+    catalog
+        .create_database("keep_me", false, Default::default())
+        .await
+        .unwrap();
+    assert_sql_error_contains(
+        &sql_context,
+        "DROP DATABASE keep_me PURGE",
+        "DROP DATABASE options are not supported",
+    )
+    .await;
+
+    assert!(catalog
+        .list_databases()
+        .await
+        .unwrap()
+        .contains(&"keep_me".to_string()));
+}
+
+#[tokio::test]
+async fn test_use_catalog_qualified_database() {
+    let (_tmp1, catalog1) = create_test_env();
+    let (_tmp2, catalog2) = create_test_env();
+    let mut sql_context = SQLContext::new();
+    sql_context
+        .register_catalog("cat1", catalog1.clone())
+        .await
+        .unwrap();
+    sql_context
+        .register_catalog("cat2", catalog2.clone())
+        .await
+        .unwrap();
+
+    sql_context
+        .sql("CREATE DATABASE cat2.analytics")
+        .await
+        .unwrap();
+    sql_context.sql("USE cat2.analytics").await.unwrap();
+    sql_context
+        .sql("CREATE TABLE events (id INT)")
+        .await
+        .unwrap();
+
+    assert!(catalog1.list_tables("analytics").await.is_err());
+    assert_eq!(
+        catalog2.list_tables("analytics").await.unwrap(),
+        vec!["events"]
+    );
+    assert_sql_error_contains(
+        &sql_context,
+        "USE missing_database",
+        "Database missing_database does not exist",
+    )
+    .await;
+}
 
 #[tokio::test]
 async fn test_create_schema() {
diff --git a/docs/src/sql.md b/docs/src/sql.md
index 08b76482..42c78cbb 100644
--- a/docs/src/sql.md
+++ b/docs/src/sql.md
@@ -546,14 +546,19 @@ register_variant_functions(&ctx);
 
 ## DDL
 
-### CREATE DATABASE / CREATE SCHEMA / DROP SCHEMA
+### DATABASE
 
 ```sql
-CREATE SCHEMA paimon.my_db;
+SHOW DATABASES;
 CREATE DATABASE paimon.my_db;
-DROP SCHEMA paimon.my_db CASCADE;
+USE paimon.my_db;
+DROP DATABASE paimon.my_db CASCADE;
 ```
 
+`CREATE DATABASE` supports `IF NOT EXISTS`, and `DROP DATABASE` supports
+`IF EXISTS` and `CASCADE`. `CREATE SCHEMA` and `DROP SCHEMA` remain supported
+as compatibility aliases.
+
 ### CREATE TABLE
 
 ```sql

Reply via email to