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 5c87afeceff [FLINK-36688][table-planner] Sort metadata keys when 
reusing source
5c87afeceff is described below

commit 5c87afeceffc52d2e0a24cf547bae6eda4075cfe
Author: Xuyang <[email protected]>
AuthorDate: Thu Dec 5 06:03:34 2024 +0800

    [FLINK-36688][table-planner] Sort metadata keys when reusing source
---
 .../planner/connectors/DynamicSourceUtils.java     |  2 +-
 .../flink/table/planner/plan/reuse/ScanReuser.java |  4 +-
 .../table/planner/plan/reuse/ScanReuserUtils.java  | 12 ++++++
 .../table/planner/plan/optimize/ScanReuseTest.java | 43 ++++++++++++++++++++++
 .../table/planner/plan/optimize/ScanReuseTest.xml  | 28 ++++++++++++++
 5 files changed, 87 insertions(+), 2 deletions(-)

diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/DynamicSourceUtils.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/DynamicSourceUtils.java
index 7a61a61a191..b575f0916a0 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/DynamicSourceUtils.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/DynamicSourceUtils.java
@@ -408,7 +408,7 @@ public final class DynamicSourceUtils {
         relBuilder.push(scan);
     }
 
-    private static Map<String, DataType> extractMetadataMap(DynamicTableSource 
source) {
+    public static Map<String, DataType> extractMetadataMap(DynamicTableSource 
source) {
         if (source instanceof SupportsReadingMetadata) {
             return ((SupportsReadingMetadata) source).listReadableMetadata();
         }
diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/reuse/ScanReuser.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/reuse/ScanReuser.java
index 9128e110346..8da14bb06f2 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/reuse/ScanReuser.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/reuse/ScanReuser.java
@@ -55,6 +55,7 @@ import java.util.stream.Collectors;
 import static 
org.apache.flink.table.planner.plan.reuse.ScanReuserUtils.abilitySpecsWithoutEscaped;
 import static 
org.apache.flink.table.planner.plan.reuse.ScanReuserUtils.concatProjectedFields;
 import static 
org.apache.flink.table.planner.plan.reuse.ScanReuserUtils.createCalcForScan;
+import static 
org.apache.flink.table.planner.plan.reuse.ScanReuserUtils.enforceMetadataKeyOrder;
 import static 
org.apache.flink.table.planner.plan.reuse.ScanReuserUtils.getAdjustedWatermarkSpec;
 import static 
org.apache.flink.table.planner.plan.reuse.ScanReuserUtils.indexOf;
 import static 
org.apache.flink.table.planner.plan.reuse.ScanReuserUtils.metadataKeys;
@@ -165,7 +166,8 @@ public class ScanReuser {
             }
 
             int[][] allProjectFields = allProjectFieldSet.toArray(new 
int[0][]);
-            List<String> allMetaKeys = new ArrayList<>(allMetaKeySet);
+            List<String> allMetaKeys =
+                    enforceMetadataKeyOrder(allMetaKeySet, 
pickTable.tableSource());
 
             // 2. Create new source.
             List<SourceAbilitySpec> specs = 
abilitySpecsWithoutEscaped(pickTable);
diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/reuse/ScanReuserUtils.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/reuse/ScanReuserUtils.java
index 43a00d720a2..b034a3ffee2 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/reuse/ScanReuserUtils.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/reuse/ScanReuserUtils.java
@@ -21,6 +21,7 @@ package org.apache.flink.table.planner.plan.reuse;
 import org.apache.flink.table.api.TableException;
 import org.apache.flink.table.catalog.Column;
 import org.apache.flink.table.catalog.ResolvedSchema;
+import org.apache.flink.table.connector.source.DynamicTableSource;
 import org.apache.flink.table.planner.connectors.DynamicSourceUtils;
 import org.apache.flink.table.planner.plan.abilities.source.FilterPushDownSpec;
 import 
org.apache.flink.table.planner.plan.abilities.source.ProjectPushDownSpec;
@@ -50,6 +51,7 @@ import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.Optional;
+import java.util.Set;
 import java.util.function.Predicate;
 import java.util.stream.Collectors;
 import java.util.stream.IntStream;
@@ -58,6 +60,7 @@ import java.util.stream.Stream;
 import scala.Option;
 
 import static 
org.apache.flink.table.planner.connectors.DynamicSourceUtils.createRequiredMetadataColumns;
+import static 
org.apache.flink.table.planner.connectors.DynamicSourceUtils.extractMetadataMap;
 
 /** Utils for {@link ScanReuser}. */
 public class ScanReuserUtils {
@@ -270,6 +273,15 @@ public class ScanReuserUtils {
                 : meta.getMetadataKeys();
     }
 
+    public static List<String> enforceMetadataKeyOrder(
+            Set<String> allUsedMetadataKeys, DynamicTableSource source) {
+        Set<String> allOrderedMetadataKeysFromTable = 
extractMetadataMap(source).keySet();
+
+        return allOrderedMetadataKeysFromTable.stream()
+                .filter(allUsedMetadataKeys::contains)
+                .collect(Collectors.toList());
+    }
+
     public static int[][] concatProjectedFields(
             ResolvedSchema schema,
             RowType originType,
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/optimize/ScanReuseTest.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/optimize/ScanReuseTest.java
index 82b9d851ac8..aa9bb818b3e 100644
--- 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/optimize/ScanReuseTest.java
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/optimize/ScanReuseTest.java
@@ -18,6 +18,7 @@
 
 package org.apache.flink.table.planner.plan.optimize;
 
+import org.apache.flink.table.api.StatementSet;
 import org.apache.flink.table.api.TableConfig;
 import org.apache.flink.table.planner.utils.TableTestBase;
 import org.apache.flink.table.planner.utils.TableTestUtil;
@@ -401,4 +402,46 @@ class ScanReuseTest extends TableTestBase {
                         + " WHERE T1.a1 = T2.a1";
         util.verifyExecPlan(sqlQuery);
     }
+
+    @TestTemplate
+    void testReuseWithReadMetadataKeepOrder() {
+        assumeThat(isStreaming).isTrue();
+        util.tableEnv()
+                .executeSql(
+                        "CREATE TEMPORARY TABLE src( "
+                                + "   `origin_ts` TIMESTAMP(3) METADATA 
VIRTUAL, "
+                                + "   `partition` INT METADATA VIRTUAL, "
+                                + "   `offset` BIGINT METADATA VIRTUAL, "
+                                + "   `id` BIGINT,"
+                                + "   PRIMARY KEY (`id`) NOT ENFORCED "
+                                + ") WITH ( "
+                                + "   'connector' = 'values', "
+                                + "   'readable-metadata' = 
'offset:bigint,origin_ts:timestamp(3),partition:int'"
+                                + ")");
+
+        util.tableEnv()
+                .executeSql(
+                        "CREATE TEMPORARY TABLE snk1( "
+                                + "   `origin_ts` TIMESTAMP(3), "
+                                + "   `partition` INT, "
+                                + "   `offset` BIGINT, "
+                                + "   `id` BIGINT"
+                                + ") WITH ( "
+                                + "   'connector' = 'values' "
+                                + ")");
+
+        util.tableEnv()
+                .executeSql(
+                        "CREATE TEMPORARY TABLE snk2( "
+                                + "   `id` BIGINT"
+                                + ") WITH ( "
+                                + "   'connector' = 'values' "
+                                + ")");
+
+        StatementSet stmt = util.tableEnv().createStatementSet();
+        stmt.addInsertSql(
+                "INSERT INTO snk1 select `origin_ts`, `partition`, `offset`, 
`id` from src");
+        stmt.addInsertSql("INSERT INTO snk2 select `id` from src");
+        util.verifyExecPlan(stmt);
+    }
 }
diff --git 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/ScanReuseTest.xml
 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/ScanReuseTest.xml
index 6f4bfa0f6a7..867f4bba401 100644
--- 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/ScanReuseTest.xml
+++ 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/optimize/ScanReuseTest.xml
@@ -1348,6 +1348,34 @@ Calc(select=[a1, a2])
                +- LocalWindowAggregate(groupBy=[a0], 
window=[TUMBLE(time_col=[rowtime], size=[1 s])], select=[a0, MIN(a1) AS min$0, 
MIN(a2) AS min$1, slice_end('w$) AS $slice_end])
                   +- Calc(select=[a0, a1, a2, metadata_0, 
Reinterpret(TO_TIMESTAMP(ts)) AS rowtime])
                      +- Reused(reference_id=[1])
+]]>
+    </Resource>
+  </TestCase>
+  <TestCase name="testReuseWithReadMetadataKeepOrder[isStreaming: true]">
+    <Resource name="ast">
+      <![CDATA[
+LogicalSink(table=[default_catalog.default_database.snk1], fields=[origin_ts, 
partition, offset, id])
++- LogicalProject(origin_ts=[$0], partition=[$1], offset=[$2], id=[$3])
+   +- LogicalProject(origin_ts=[$2], partition=[$3], offset=[$1], id=[$0])
+      +- LogicalTableScan(table=[[default_catalog, default_database, src, 
metadata=[offset, origin_ts, partition]]])
+
+LogicalSink(table=[default_catalog.default_database.snk2], fields=[id])
++- LogicalProject(id=[$3])
+   +- LogicalProject(origin_ts=[$2], partition=[$3], offset=[$1], id=[$0])
+      +- LogicalTableScan(table=[[default_catalog, default_database, src, 
metadata=[offset, origin_ts, partition]]])
+]]>
+    </Resource>
+    <Resource name="optimized exec plan">
+      <![CDATA[
+TableSourceScan(table=[[default_catalog, default_database, src, project=[id, 
offset, origin_ts, partition], metadata=[offset, origin_ts, partition]]], 
fields=[id, offset, origin_ts, partition])(reuse_id=[1])
+
+Sink(table=[default_catalog.default_database.snk1], fields=[origin_ts, 
partition, offset, id])
++- Calc(select=[origin_ts, partition, offset, id])
+   +- Reused(reference_id=[1])
+
+Sink(table=[default_catalog.default_database.snk2], fields=[id])
++- Calc(select=[id])
+   +- Reused(reference_id=[1])
 ]]>
     </Resource>
   </TestCase>

Reply via email to