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