JingsongLi commented on code in PR #1021:
URL: https://github.com/apache/paimon-rust/pull/1021#discussion_r4177569517
##########
crates/integrations/datafusion/src/sql_context.rs:
##########
@@ -1055,6 +1055,117 @@ 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);
+ }
+
+ // Expand catalog SQL functions once (so `plus_one(...)` and other
+ // catalog-registered functions resolve exactly as a plain SELECT does
+ // through SQLContext), then reuse the expanded query for both schema
+ // inference and population. Planning (not collecting) is enough for
+ // inference; the INSERT below writes the rows from the same expansion.
+ let query_sql = query.to_string();
+ let mut state = self.ctx.state();
+ state.config_mut().options_mut().catalog.default_catalog =
catalog_name.clone();
Review Comment:
[P1] Keep the source SELECT in the current session namespace. These
assignments infer the schema using the destination catalog/database, whereas
the generated INSERT still evaluates its SELECT in the session current
namespace. With current database test_db, test_db.source(id INT, value INT)
containing (1,99), and other_db.source(value INT, id INT), CREATE TABLE
paimon.other_db.out AS SELECT * FROM source succeeds but SELECT id,value FROM
out returns (99,1): the destination schema was inferred from a different source
table and the source values were inserted positionally. If other_db.source does
not exist, an otherwise valid SELECT fails at inference instead. The same issue
affects unqualified function expansion in the target namespace. A real Parquet
probe fails on this head and the main-integrated tree, but passes with the
previous SQLContext. Use the source query/session namespace consistently for
expansion, inference and population, or reuse the same analyzed source plan.
--
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]