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]