This is an automated email from the ASF dual-hosted git repository.

snuyanzin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git


The following commit(s) were added to refs/heads/master by this push:
     new f5c44b683cd [FLINK-36783][table-planner] Use validated SqlNode in 
`asQuery` while converting `SqlCreateTableAS` to `CreateTableASOperation`
f5c44b683cd is described below

commit f5c44b683cdaa647152fd6f64861f4fd3bd1c7b8
Author: Xuyang <[email protected]>
AuthorDate: Tue Dec 3 16:06:04 2024 +0800

    [FLINK-36783][table-planner] Use validated SqlNode in `asQuery` while 
converting `SqlCreateTableAS` to `CreateTableASOperation`
---
 .../operations/SqlCreateTableConverter.java        | 10 ++--
 .../table/planner/plan/batch/sql/TableSinkTest.xml | 53 ++++++++++++++++------
 .../planner/plan/batch/sql/TableSinkTest.scala     | 14 ++++++
 .../runtime/batch/sql/TableSinkITCase.scala        | 23 ++++++++++
 4 files changed, 82 insertions(+), 18 deletions(-)

diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/operations/SqlCreateTableConverter.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/operations/SqlCreateTableConverter.java
index dbff49f5310..b45e93f7b85 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/operations/SqlCreateTableConverter.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/operations/SqlCreateTableConverter.java
@@ -104,16 +104,18 @@ class SqlCreateTableConverter {
                 UnresolvedIdentifier.of(sqlCreateTableAs.fullTableName());
         ObjectIdentifier identifier = 
catalogManager.qualifyIdentifier(unresolvedIdentifier);
 
+        SqlNode asQuerySqlNode = sqlCreateTableAs.getAsQuery();
+        SqlNode validatedAsQuery = flinkPlanner.validate(asQuerySqlNode);
+
         PlannerQueryOperation query =
                 (PlannerQueryOperation)
                         SqlNodeToOperationConversion.convert(
-                                        flinkPlanner, catalogManager, 
sqlCreateTableAs.getAsQuery())
+                                        flinkPlanner, catalogManager, 
validatedAsQuery)
                                 .orElseThrow(
                                         () ->
                                                 new TableException(
                                                         "CTAS unsupported node 
type "
-                                                                + 
sqlCreateTableAs
-                                                                        
.getAsQuery()
+                                                                + 
validatedAsQuery
                                                                         
.getClass()
                                                                         
.getSimpleName()));
         ResolvedCatalogTable tableWithResolvedSchema =
@@ -125,7 +127,7 @@ class SqlCreateTableConverter {
                         catalogManager,
                         flinkPlanner,
                         query,
-                        sqlCreateTableAs.getAsQuery(),
+                        validatedAsQuery,
                         tableWithResolvedSchema);
 
         CreateTableOperation createTableOperation =
diff --git 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/batch/sql/TableSinkTest.xml
 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/batch/sql/TableSinkTest.xml
index dcc55d73b48..b7152ec7d27 100644
--- 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/batch/sql/TableSinkTest.xml
+++ 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/batch/sql/TableSinkTest.xml
@@ -16,16 +16,32 @@ See the License for the specific language governing 
permissions and
 limitations under the License.
 -->
 <Root>
-  <TestCase name="testManagedTableSinkWithEnableCheckpointing">
-    <Resource name="ast">
-      <![CDATA[
-LogicalSink(table=[default_catalog.default_database.sink], fields=[a, b, c])
-+- LogicalProject(a=[$0], b=[$1], c=[$2])
-   +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]])
+  <TestCase name="testCreateTableAsSelectWithOrderKeyNotProjected">
+    <Resource name="explain">
+      <![CDATA[== Abstract Syntax Tree ==
+LogicalSink(table=[default_catalog.default_database.MyCtasSource], fields=[b, 
c, d])
++- LogicalProject(b=[$0], c=[$1], d=[$2])
+   +- LogicalSort(sort0=[$3], dir0=[ASC-nulls-first])
+      +- LogicalProject(b=[$1], c=[$2], d=[$3], a=[$0])
+         +- LogicalValues(tuples=[[{ 1, 1, 2, _UTF-16LE'd1' }]])
+
+== Optimized Physical Plan ==
+Sink(table=[default_catalog.default_database.MyCtasSource], fields=[b, c, d])
++- Calc(select=[b, c, d])
+   +- Sort(orderBy=[a ASC])
+      +- Exchange(distribution=[single])
+         +- Values(tuples=[[{ 1, 1, 2, _UTF-16LE'd1' }]], values=[a, b, c, d])
+
+== Optimized Execution Plan ==
+Sink(table=[default_catalog.default_database.MyCtasSource], fields=[b, c, d])
++- Calc(select=[b, c, d])
+   +- Sort(orderBy=[a ASC])
+      +- Exchange(distribution=[single])
+         +- Values(tuples=[[{ 1, 1, 2, _UTF-16LE'd1' }]], values=[a, b, c, d])
 ]]>
     </Resource>
   </TestCase>
-  <TestCase name="testDynamicPartWithOrderBy">
+  <TestCase name="testDistribution">
     <Resource name="ast">
       <![CDATA[
 LogicalSink(table=[default_catalog.default_database.sink], fields=[a, b])
@@ -44,24 +60,33 @@ Sink(table=[default_catalog.default_database.sink], 
fields=[a, b])
 ]]>
     </Resource>
   </TestCase>
-  <TestCase name="testDistribution">
-       <Resource name="ast">
-               <![CDATA[
+  <TestCase name="testManagedTableSinkWithEnableCheckpointing">
+    <Resource name="ast">
+      <![CDATA[
+LogicalSink(table=[default_catalog.default_database.sink], fields=[a, b, c])
++- LogicalProject(a=[$0], b=[$1], c=[$2])
+   +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]])
+]]>
+    </Resource>
+  </TestCase>
+  <TestCase name="testDynamicPartWithOrderBy">
+    <Resource name="ast">
+      <![CDATA[
 LogicalSink(table=[default_catalog.default_database.sink], fields=[a, b])
 +- LogicalSort(sort0=[$0], dir0=[ASC-nulls-first])
    +- LogicalProject(a=[$0], b=[$1])
       +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]])
 ]]>
-       </Resource>
-       <Resource name="optimized exec plan">
-               <![CDATA[
+    </Resource>
+    <Resource name="optimized exec plan">
+      <![CDATA[
 Sink(table=[default_catalog.default_database.sink], fields=[a, b])
 +- Sort(orderBy=[a ASC])
    +- Exchange(distribution=[single])
       +- Calc(select=[a, b])
          +- BoundedStreamScan(table=[[default_catalog, default_database, 
MyTable]], fields=[a, b, c])
 ]]>
-       </Resource>
+    </Resource>
   </TestCase>
   <TestCase name="testManagedTableSinkWithDisableCheckpointing">
     <Resource name="ast">
diff --git 
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/batch/sql/TableSinkTest.scala
 
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/batch/sql/TableSinkTest.scala
index 2dc3faa96db..4d4cdda2af9 100644
--- 
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/batch/sql/TableSinkTest.scala
+++ 
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/batch/sql/TableSinkTest.scala
@@ -174,4 +174,18 @@ class TableSinkTest extends TableTestBase {
 
     util.verifyAstPlan(stmtSet)
   }
+
+  @Test
+  def testCreateTableAsSelectWithOrderKeyNotProjected(): Unit = {
+    util.verifyExplainInsert(s"""
+                                |create table MyCtasSource
+                                |WITH (
+                                |   'connector' = 'values'
+                                |) as select b, c, d from
+                                |  (values
+                                |    (1, 1, 2, 'd1')
+                                |  ) as V(a, b, c, d)
+                                |  order by a
+                                |""".stripMargin)
+  }
 }
diff --git 
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/runtime/batch/sql/TableSinkITCase.scala
 
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/runtime/batch/sql/TableSinkITCase.scala
index a32e7374c06..623f63cf13c 100644
--- 
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/runtime/batch/sql/TableSinkITCase.scala
+++ 
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/runtime/batch/sql/TableSinkITCase.scala
@@ -30,6 +30,8 @@ import org.apache.flink.table.planner.utils.TableTestUtil
 import org.assertj.core.api.Assertions.{assertThat, assertThatThrownBy}
 import org.junit.jupiter.api.{BeforeEach, Test}
 
+import scala.collection.convert.ImplicitConversions._
+
 class TableSinkITCase extends BatchTestBase {
 
   @BeforeEach
@@ -160,4 +162,25 @@ class TableSinkITCase extends BatchTestBase {
           .await())
       .hasRootCauseMessage("\nExpecting actual not to be null")
   }
+
+  @Test
+  def testCreateTableAsSelectWithOrderKeyNotProjected(): Unit = {
+    tEnv
+      .executeSql(s"""
+                     |create table MyCtasTable
+                     |WITH (
+                     |   'connector' = 'values'
+                     |) as select b, c, d from
+                     |  (values
+                     |    (1, 1, 2, 'd1'),
+                     |    (2, 2, 4, 'd2')
+                     |  ) as V(a, b, c, d)
+                     |  order by a
+                     |""".stripMargin)
+      .await()
+
+    val expected = List("+I[1, 2, d1]", "+I[2, 4, d2]")
+    assertThat(TestValuesTableFactory.getResultsAsStrings("MyCtasTable").toSeq)
+      .isEqualTo(expected)
+  }
 }

Reply via email to