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)
+ }
}