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]

Reply via email to