JingsongLi commented on code in PR #1021:
URL: https://github.com/apache/paimon-rust/pull/1021#discussion_r4178014722
##########
crates/integrations/datafusion/src/sql_context.rs:
##########
@@ -1055,6 +1055,139 @@ 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, whatever its declared
+ // read engine. Use the load path and treat only TableNotExist as
+ // absence: FileSystemCatalog returns Unsupported for an existing
+ // engine-served table, and mistaking that for "missing" would try to
+ // populate a table that already exists instead of being a no-op.
+ if ct.if_not_exists {
+ match catalog.load_table(&identifier).await {
+ Ok(_) => return ok_result(&self.ctx),
+ Err(paimon::Error::TableNotExist { .. }) => {}
+ Err(e) => return Err(to_datafusion_error(e)),
+ }
+ }
+
+ // Resolve the source query in the session's current namespace — the
same
+ // namespace the INSERT below evaluates it in — not the destination's.
+ // Inferring the schema against the target database would read a
+ // differently-shaped `source` there while the session's `source` is
+ // inserted positionally. Expand catalog SQL functions once (so
+ // `plus_one(...)` resolves exactly as a plain SELECT does through
+ // SQLContext) and reuse the expanded query for both inference and
+ // population; planning (not collecting) is enough for inference.
+ let query_sql = query.to_string();
+ let state = self.ctx.state();
+ let session_catalog =
state.config_options().catalog.default_catalog.clone();
+ let session_schema =
state.config_options().catalog.default_schema.clone();
+ let expanded_query = crate::sql_function::expand_sql(
+ &query_sql,
+ &self.catalogs,
+ &session_catalog,
+ &session_schema,
+ )
+ .await?;
+ let logical_plan = state.create_logical_plan(&expanded_query).await?;
+ let arrow_fields = logical_plan
+ .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 expanded query. Quote every
+ // identifier component so a case-sensitive or special target (e.g.
+ // `"Dst"`) is written as created, not a different existing table that
+ // unquoted text would normalize to. `Box::pin` breaks the async
+ // recursion through `sql` (CTAS -> INSERT -> dispatch); `collect`
drives
+ // the insert plan to completion.
+ let insert_sql = format!(
+ "INSERT INTO {}.{}.{} {}",
+ quote_ident(&catalog_name),
+ quote_ident(identifier.database()),
+ quote_ident(identifier.object()),
+ expanded_query
+ );
+ // On failure, drop the table this statement just created (it did not
+ // exist before) and surface the original error, so a corrected retry
is
+ // not silently satisfied by an empty leftover.
+ let populated: DFResult<Vec<RecordBatch>> =
+ async { Box::pin(self.sql(&insert_sql)).await?.collect().await
}.await;
+ if let Err(populate_err) = populated {
+ let _ = catalog.drop_table(&identifier, true).await;
Review Comment:
[P1] Clean up only a CTAS target this statement created
With `IF NOT EXISTS`, the initial `load_table` can report absence, but
another client can create and populate the target before the subsequent
`create_table(..., ct.if_not_exists)` call. That call now succeeds as a no-op,
so this statement does not own the table. If its population fails, this
unconditional drop deletes the other client's table and committed data. A
deterministic probe wrapping the real FileSystemCatalog reproduced this: client
B creates `test_db.race_ctas` and commits Parquet row90 between client A's
absence check and create; A runs `CREATE TABLE IF NOT EXISTS
paimon.test_db.race_ctas AS SELECT CAST('bad' AS INT) AS id`, fails, and B's
table becomes TableNotExist. Java SparkCatalog passes `false` to createTable
and propagates AlreadyExists before population; changing only this CTAS create
to strict creation makes all 7 actual probes pass and preserves row90. After
the initial IF-NOT-EXISTS check, perform strict creation and handle a raced
AlreadyExists before ente
ring population, so cleanup is reached only after this statement successfully
creates the target.
--
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]