JingsongLi commented on code in PR #1021:
URL: https://github.com/apache/paimon-rust/pull/1021#discussion_r4176849429


##########
crates/integrations/datafusion/src/sql_context.rs:
##########
@@ -1055,6 +1055,105 @@ impl SQLContext {
         ok_result(&self.ctx)
     }
 
+    /// Handle `CREATE TABLE ... AS SELECT ...` for persistent Paimon tables.
+    ///
+    /// The schema is inferred from the query's output. `PARTITIONED BY`, a
+    /// `PRIMARY KEY` constraint, and `WITH (...)` options still apply; an
+    /// explicit column list is not accepted (the columns come from the query).
+    /// After the table is created it is populated with the same query.
+    async fn handle_create_table_as_select(
+        &self,
+        catalog: &Arc<dyn Catalog>,
+        ct: &CreateTable,
+        partition_keys: Vec<String>,
+    ) -> DFResult<DataFrame> {
+        let Some(query) = &ct.query else {
+            return Err(DataFusionError::Internal(
+                "handle_create_table_as_select called without a 
query".to_string(),
+            ));
+        };
+        if !ct.columns.is_empty() {
+            return Err(DataFusionError::Plan(
+                "CREATE TABLE AS SELECT does not accept an explicit column 
list; \
+                 the schema is inferred from the query"
+                    .to_string(),
+            ));
+        }
+
+        let (_, catalog_name, identifier) = 
self.resolve_catalog_and_table(&ct.name)?;
+
+        // IF NOT EXISTS on an existing table is a no-op (no rows appended).
+        if ct.if_not_exists && catalog.get_table(&identifier).await.is_ok() {
+            return ok_result(&self.ctx);
+        }
+
+        // Infer the output schema from the query. Planning (not collecting) is
+        // enough here; the rows are written by the INSERT below, which 
resolves
+        // the query the same way.
+        let query_sql = query.to_string();
+        let df = self.ctx.sql(&query_sql).await?;
+        let arrow_fields = df
+            .schema()
+            .as_arrow()
+            .fields()
+            .iter()
+            .map(|field| field.as_ref().clone())
+            .collect::<Vec<_>>();
+        let fields =
+            
paimon::arrow::arrow_fields_to_paimon(&arrow_fields).map_err(to_datafusion_error)?;
+
+        // Build the Paimon schema: inferred columns + PRIMARY KEY / 
PARTITIONED BY / WITH.
+        let enable_ident_normalization = self.ctx.enable_ident_normalization();
+        let mut builder = paimon::spec::Schema::builder();
+        for field in &fields {
+            builder = builder.column(field.name().to_string(), 
field.data_type().clone());
+        }
+        for constraint in &ct.constraints {
+            if let 
datafusion::sql::sqlparser::ast::TableConstraint::PrimaryKey(pk) = constraint {
+                let pk_cols: Vec<String> = pk
+                    .columns
+                    .iter()
+                    .map(|c| primary_key_column_name(&c.column.expr, 
enable_ident_normalization))
+                    .collect();
+                builder = builder.primary_key(pk_cols);
+            }
+        }
+        if !partition_keys.is_empty() {
+            let field_names: Vec<&str> = fields.iter().map(|f| 
f.name()).collect();
+            for pk in &partition_keys {
+                if !field_names.contains(&pk.as_str()) {
+                    return Err(DataFusionError::Plan(format!(
+                        "PARTITIONED BY column '{pk}' is not produced by the 
query"
+                    )));
+                }
+            }
+            builder = builder.partition_keys(partition_keys);
+        }
+        for (k, v) in extract_options(&ct.table_options)? {
+            builder = builder.option(k, v);
+        }
+        let schema = builder.build().map_err(to_datafusion_error)?;
+
+        catalog
+            .create_table(&identifier, schema, ct.if_not_exists)
+            .await
+            .map_err(to_datafusion_error)?;
+
+        // Populate the new table from the same query. `Box::pin` breaks the
+        // async recursion through `sql` (CTAS -> INSERT -> dispatch); 
`collect`
+        // drives the insert plan to completion.
+        let insert_sql = format!(
+            "INSERT INTO {}.{}.{} {}",
+            catalog_name,
+            identifier.database(),
+            identifier.object(),
+            query_sql
+        );
+        Box::pin(self.sql(&insert_sql)).await?.collect().await?;

Review Comment:
   [P2] Clean up the newly created table when CTAS population fails
   
   `CREATE TABLE paimon.test_db.failed_ctas AS SELECT CAST('bad' AS INT) AS id` 
fails while optimizing the INSERT, after `schema/schema-0` has already been 
persisted. The failed CTAS therefore leaves an empty target behind; a corrected 
retry with IF NOT EXISTS returns success without populating it. Paimon's 
[Java/Spark CTAS 
strategy](https://github.com/apache/paimon/blob/master/paimon-spark/paimon-spark-common/src/main/scala/org/apache/spark/sql/execution/shim/PaimonCreateTableAsSelectStrategy.scala)
 uses [Spark's CTAS 
execution](https://github.com/apache/spark/blob/v3.5.8/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/WriteToDataSourceV2Exec.scala#L572),
 whose population-failure callback drops the new table or aborts staged 
changes. Please preserve that behavior: clean up only the target this operation 
created, and retain the original error.



##########
crates/integrations/datafusion/src/sql_context.rs:
##########
@@ -1055,6 +1055,105 @@ impl SQLContext {
         ok_result(&self.ctx)
     }
 
+    /// Handle `CREATE TABLE ... AS SELECT ...` for persistent Paimon tables.
+    ///
+    /// The schema is inferred from the query's output. `PARTITIONED BY`, a
+    /// `PRIMARY KEY` constraint, and `WITH (...)` options still apply; an
+    /// explicit column list is not accepted (the columns come from the query).
+    /// After the table is created it is populated with the same query.
+    async fn handle_create_table_as_select(
+        &self,
+        catalog: &Arc<dyn Catalog>,
+        ct: &CreateTable,
+        partition_keys: Vec<String>,
+    ) -> DFResult<DataFrame> {
+        let Some(query) = &ct.query else {
+            return Err(DataFusionError::Internal(
+                "handle_create_table_as_select called without a 
query".to_string(),
+            ));
+        };
+        if !ct.columns.is_empty() {
+            return Err(DataFusionError::Plan(
+                "CREATE TABLE AS SELECT does not accept an explicit column 
list; \
+                 the schema is inferred from the query"
+                    .to_string(),
+            ));
+        }
+
+        let (_, catalog_name, identifier) = 
self.resolve_catalog_and_table(&ct.name)?;
+
+        // IF NOT EXISTS on an existing table is a no-op (no rows appended).
+        if ct.if_not_exists && catalog.get_table(&identifier).await.is_ok() {

Review Comment:
   [P2] Check target existence independently of its read engine
   
   `get_table(...).is_ok()` treats every lookup failure as a missing target. 
FileSystemCatalog deliberately returns Unsupported for an existing 
engine-served table. In a real filesystem probe, after creating 
`existing_iceberg` with `WITH ('type'='iceberg-table')`, `CREATE TABLE IF NOT 
EXISTS ...existing_iceberg AS SELECT CAST(7 AS INT) AS id` tries to populate it 
and fails because no Iceberg engine is registered, rather than being the 
promised no-op. Use a type-independent existence/load path and only treat 
TableNotExist as absence; propagate unrelated lookup errors.



##########
crates/integrations/datafusion/src/sql_context.rs:
##########
@@ -1055,6 +1055,105 @@ impl SQLContext {
         ok_result(&self.ctx)
     }
 
+    /// Handle `CREATE TABLE ... AS SELECT ...` for persistent Paimon tables.
+    ///
+    /// The schema is inferred from the query's output. `PARTITIONED BY`, a
+    /// `PRIMARY KEY` constraint, and `WITH (...)` options still apply; an
+    /// explicit column list is not accepted (the columns come from the query).
+    /// After the table is created it is populated with the same query.
+    async fn handle_create_table_as_select(
+        &self,
+        catalog: &Arc<dyn Catalog>,
+        ct: &CreateTable,
+        partition_keys: Vec<String>,
+    ) -> DFResult<DataFrame> {
+        let Some(query) = &ct.query else {
+            return Err(DataFusionError::Internal(
+                "handle_create_table_as_select called without a 
query".to_string(),
+            ));
+        };
+        if !ct.columns.is_empty() {
+            return Err(DataFusionError::Plan(
+                "CREATE TABLE AS SELECT does not accept an explicit column 
list; \
+                 the schema is inferred from the query"
+                    .to_string(),
+            ));
+        }
+
+        let (_, catalog_name, identifier) = 
self.resolve_catalog_and_table(&ct.name)?;
+
+        // IF NOT EXISTS on an existing table is a no-op (no rows appended).
+        if ct.if_not_exists && catalog.get_table(&identifier).await.is_ok() {
+            return ok_result(&self.ctx);
+        }
+
+        // Infer the output schema from the query. Planning (not collecting) is
+        // enough here; the rows are written by the INSERT below, which 
resolves
+        // the query the same way.
+        let query_sql = query.to_string();
+        let df = self.ctx.sql(&query_sql).await?;
+        let arrow_fields = df
+            .schema()
+            .as_arrow()
+            .fields()
+            .iter()
+            .map(|field| field.as_ref().clone())
+            .collect::<Vec<_>>();
+        let fields =
+            
paimon::arrow::arrow_fields_to_paimon(&arrow_fields).map_err(to_datafusion_error)?;
+
+        // Build the Paimon schema: inferred columns + PRIMARY KEY / 
PARTITIONED BY / WITH.
+        let enable_ident_normalization = self.ctx.enable_ident_normalization();
+        let mut builder = paimon::spec::Schema::builder();
+        for field in &fields {
+            builder = builder.column(field.name().to_string(), 
field.data_type().clone());
+        }
+        for constraint in &ct.constraints {
+            if let 
datafusion::sql::sqlparser::ast::TableConstraint::PrimaryKey(pk) = constraint {
+                let pk_cols: Vec<String> = pk
+                    .columns
+                    .iter()
+                    .map(|c| primary_key_column_name(&c.column.expr, 
enable_ident_normalization))
+                    .collect();
+                builder = builder.primary_key(pk_cols);
+            }
+        }
+        if !partition_keys.is_empty() {
+            let field_names: Vec<&str> = fields.iter().map(|f| 
f.name()).collect();
+            for pk in &partition_keys {
+                if !field_names.contains(&pk.as_str()) {
+                    return Err(DataFusionError::Plan(format!(
+                        "PARTITIONED BY column '{pk}' is not produced by the 
query"
+                    )));
+                }
+            }
+            builder = builder.partition_keys(partition_keys);
+        }
+        for (k, v) in extract_options(&ct.table_options)? {
+            builder = builder.option(k, v);
+        }
+        let schema = builder.build().map_err(to_datafusion_error)?;
+
+        catalog
+            .create_table(&identifier, schema, ct.if_not_exists)
+            .await
+            .map_err(to_datafusion_error)?;
+
+        // Populate the new table from the same query. `Box::pin` breaks the
+        // async recursion through `sql` (CTAS -> INSERT -> dispatch); 
`collect`
+        // drives the insert plan to completion.
+        let insert_sql = format!(
+            "INSERT INTO {}.{}.{} {}",

Review Comment:
   [P1] Quote the resolved CTAS target before generating INSERT
   
   The target is created with its exact quoted identifier, but this INSERT 
formats catalog/database/table components without quoting or escaping. A real 
memory-backed FileSystemCatalog/Parquet probe with an existing `dst` and 
`CREATE TABLE paimon.test_db."Dst" AS SELECT CAST(7 AS INT) AS id` completed 
successfully while `"Dst"` had 0 rows and the existing `dst` had 2 rows. A 
quoted `"target-table"` also creates its schema and then fails parsing the 
unquoted INSERT. This can silently write to a different existing table. Quote 
and escape every identifier component, or construct the insert from the 
resolved TableReference without reparsing unquoted text.



##########
crates/integrations/datafusion/src/sql_context.rs:
##########
@@ -1055,6 +1055,105 @@ impl SQLContext {
         ok_result(&self.ctx)
     }
 
+    /// Handle `CREATE TABLE ... AS SELECT ...` for persistent Paimon tables.
+    ///
+    /// The schema is inferred from the query's output. `PARTITIONED BY`, a
+    /// `PRIMARY KEY` constraint, and `WITH (...)` options still apply; an
+    /// explicit column list is not accepted (the columns come from the query).
+    /// After the table is created it is populated with the same query.
+    async fn handle_create_table_as_select(
+        &self,
+        catalog: &Arc<dyn Catalog>,
+        ct: &CreateTable,
+        partition_keys: Vec<String>,
+    ) -> DFResult<DataFrame> {
+        let Some(query) = &ct.query else {
+            return Err(DataFusionError::Internal(
+                "handle_create_table_as_select called without a 
query".to_string(),
+            ));
+        };
+        if !ct.columns.is_empty() {
+            return Err(DataFusionError::Plan(
+                "CREATE TABLE AS SELECT does not accept an explicit column 
list; \
+                 the schema is inferred from the query"
+                    .to_string(),
+            ));
+        }
+
+        let (_, catalog_name, identifier) = 
self.resolve_catalog_and_table(&ct.name)?;
+
+        // IF NOT EXISTS on an existing table is a no-op (no rows appended).
+        if ct.if_not_exists && catalog.get_table(&identifier).await.is_ok() {
+            return ok_result(&self.ctx);
+        }
+
+        // Infer the output schema from the query. Planning (not collecting) is
+        // enough here; the rows are written by the INSERT below, which 
resolves
+        // the query the same way.
+        let query_sql = query.to_string();
+        let df = self.ctx.sql(&query_sql).await?;

Review Comment:
   [P2] Reuse SQLContext's supported query expansion for CTAS
   
   This bare SessionContext call bypasses the catalog SQL-function expansion 
used by ordinary SELECT in SQLContext. Using the existing catalog-function 
fixture, `SELECT plus_one(41) AS answer` returns 42, while `CREATE TABLE ... AS 
SELECT plus_one(41) AS answer` fails here with `Invalid function 'plus_one'`. 
Users cannot materialize a query that already works through the public SQL API. 
Expand the query once and reuse the expanded query/plan for both inference and 
population; only changing this inference call would leave the generated INSERT 
using the original unexpanded query.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to