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

gustavodemorais 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 70aabf68cd3 [FLINK-40334][table] Insert-only input should stay 
unmaterialized under FORCE
70aabf68cd3 is described below

commit 70aabf68cd312f9f8103b6f02b639f344643b46e
Author: Gustavo de Morais <[email protected]>
AuthorDate: Wed Aug 12 17:33:40 2026 +0200

    [FLINK-40334][table] Insert-only input should stay unmaterialized under 
FORCE
    
    This closes #28930.
---
 .../plan/nodes/exec/batch/BatchExecSink.java       |   1 +
 .../plan/nodes/exec/common/CommonExecSink.java     |  24 ++-
 .../plan/nodes/exec/stream/StreamExecSink.java     |  14 +-
 .../planner/plan/stream/sql/TableSinkTest.xml      | 173 +++++++++++++++++++++
 .../planner/plan/stream/sql/TableSinkTest.scala    |  45 ++++++
 5 files changed, 244 insertions(+), 13 deletions(-)

diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/batch/BatchExecSink.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/batch/BatchExecSink.java
index 8a1dd04bcfe..cb267e24e80 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/batch/BatchExecSink.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/batch/BatchExecSink.java
@@ -120,6 +120,7 @@ public class BatchExecSink extends CommonExecSink 
implements BatchExecNode<Objec
                 tableSink,
                 -1,
                 false,
+                null,
                 null);
     }
 
diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/common/CommonExecSink.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/common/CommonExecSink.java
index 31c0bc4074f..3738840ab27 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/common/CommonExecSink.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/common/CommonExecSink.java
@@ -36,6 +36,8 @@ import 
org.apache.flink.streaming.api.transformations.LegacySinkTransformation;
 import org.apache.flink.streaming.api.transformations.PartitionTransformation;
 import 
org.apache.flink.streaming.api.transformations.TransformationWithLineage;
 import 
org.apache.flink.streaming.runtime.partitioner.KeyGroupStreamPartitioner;
+import org.apache.flink.table.api.InsertConflictStrategy;
+import org.apache.flink.table.api.InsertConflictStrategy.ConflictBehavior;
 import org.apache.flink.table.api.TableException;
 import org.apache.flink.table.api.config.ExecutionConfigOptions;
 import org.apache.flink.table.catalog.ResolvedSchema;
@@ -84,6 +86,8 @@ import org.apache.flink.util.TemporaryClassLoaderContext;
 
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
 
+import javax.annotation.Nullable;
+
 import java.util.Arrays;
 import java.util.List;
 import java.util.Objects;
@@ -146,7 +150,8 @@ public abstract class CommonExecSink extends 
ExecNodeBase<Object>
             DynamicTableSink tableSink,
             int rowtimeFieldIndex,
             boolean upsertMaterialize,
-            int[] inputUpsertKey) {
+            int[] inputUpsertKey,
+            @Nullable InsertConflictStrategy conflictStrategy) {
         final ResolvedSchema schema = 
tableSinkSpec.getContextResolvedTable().getResolvedSchema();
         final SinkRuntimeProvider runtimeProvider =
                 tableSink.getSinkRuntimeProvider(
@@ -187,6 +192,11 @@ public abstract class CommonExecSink extends 
ExecNodeBase<Object>
         Optional<LineageVertex> lineageVertexOpt =
                 TableLineageUtils.extractLineageDataset(outputObject);
 
+        // only add materialization if input has changes, unless the conflict 
strategy has to
+        // compare every insert against the row stored under the same primary 
key
+        final boolean needMaterialization =
+                upsertMaterialize && (!inputInsertOnly || 
detectsDuplicateKeys(conflictStrategy));
+
         Transformation<RowData> sinkTransform =
                 applyConstraintValidations(inputTransform, config, 
persistedRowType);
 
@@ -199,10 +209,10 @@ public abstract class CommonExecSink extends 
ExecNodeBase<Object>
                             primaryKeys,
                             sinkParallelism,
                             inputParallelism,
-                            upsertMaterialize);
+                            needMaterialization);
         }
 
-        if (upsertMaterialize) {
+        if (needMaterialization) {
             sinkTransform =
                     applyUpsertMaterialize(
                             sinkTransform,
@@ -246,6 +256,14 @@ public abstract class CommonExecSink extends 
ExecNodeBase<Object>
         return transformation;
     }
 
+    /** Whether the conflict strategy has to detect rows that share a primary 
key. */
+    protected static boolean detectsDuplicateKeys(
+            @Nullable InsertConflictStrategy conflictStrategy) {
+        return conflictStrategy != null
+                && (conflictStrategy.getBehavior() == ConflictBehavior.ERROR
+                        || conflictStrategy.getBehavior() == 
ConflictBehavior.NOTHING);
+    }
+
     /**
      * Apply an operator to filter or report error to process not-null values 
for not-null fields.
      */
diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecSink.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecSink.java
index 22dce136a0e..acc7fa10b22 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecSink.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecSink.java
@@ -26,7 +26,6 @@ import org.apache.flink.streaming.api.TimeDomain;
 import org.apache.flink.streaming.api.operators.OneInputStreamOperator;
 import org.apache.flink.streaming.api.transformations.OneInputTransformation;
 import org.apache.flink.table.api.InsertConflictStrategy;
-import org.apache.flink.table.api.InsertConflictStrategy.ConflictBehavior;
 import org.apache.flink.table.api.TableException;
 import org.apache.flink.table.api.config.ExecutionConfigOptions;
 import 
org.apache.flink.table.api.config.ExecutionConfigOptions.RowtimeInserter;
@@ -291,7 +290,8 @@ public class StreamExecSink extends CommonExecSink 
implements StreamExecNode<Obj
                 tableSink,
                 rowtimeFieldIndex,
                 upsertMaterialize,
-                inputUpsertKey);
+                inputUpsertKey,
+                conflictStrategy);
     }
 
     @Override
@@ -359,7 +359,7 @@ public class StreamExecSink extends CommonExecSink 
implements StreamExecNode<Obj
         // This assigns the current watermark as the timestamp to each record,
         // which is required for the WatermarkCompactingSinkMaterializer to 
work correctly
         Transformation<RowData> transformForMaterializer = inputTransform;
-        if (isErrorOrNothingConflictStrategy()) {
+        if (detectsDuplicateKeys(conflictStrategy)) {
             // Use input parallelism to preserve watermark semantics
             transformForMaterializer =
                     ExecNodeUtil.createOneInputTransformation(
@@ -410,7 +410,7 @@ public class StreamExecSink extends CommonExecSink 
implements StreamExecNode<Obj
             GeneratedHashFunction rowHashFunction) {
 
         // Check if we should use the watermark-compacting materializer for 
ERROR/NOTHING strategies
-        if (isErrorOrNothingConflictStrategy()) {
+        if (detectsDuplicateKeys(conflictStrategy)) {
             RowType keyType = RowTypeUtils.projectRowType(physicalRowType, 
primaryKeys);
 
             return WatermarkCompactingSinkMaterializer.create(
@@ -450,12 +450,6 @@ public class StreamExecSink extends CommonExecSink 
implements StreamExecNode<Obj
                                 config));
     }
 
-    private boolean isErrorOrNothingConflictStrategy() {
-        return conflictStrategy != null
-                && (conflictStrategy.getBehavior() == ConflictBehavior.ERROR
-                        || conflictStrategy.getBehavior() == 
ConflictBehavior.NOTHING);
-    }
-
     private static SequencedMultiSetStateConfig createStateConfig(
             SinkUpsertMaterializeStrategy strategy,
             TimeDomain ttlTimeDomain,
diff --git 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/TableSinkTest.xml
 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/TableSinkTest.xml
index 90d4c284865..486e4d045ec 100644
--- 
a/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/TableSinkTest.xml
+++ 
b/flink-table/flink-table-planner/src/test/resources/org/apache/flink/table/planner/plan/stream/sql/TableSinkTest.xml
@@ -504,6 +504,179 @@ Sink(table=[default_catalog.default_database.sink], 
fields=[a, b])
 ]]>
     </Resource>
   </TestCase>
+  <TestCase name="testForcedMaterializeWithAppendOnlyInput">
+    <Resource name="explain">
+      <![CDATA[== Abstract Syntax Tree ==
+LogicalSink(table=[default_catalog.default_database.forcedSink], fields=[a, b])
++- LogicalProject(a=[$0], b=[$1])
+   +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]])
+
+== Optimized Physical Plan ==
+Sink(table=[default_catalog.default_database.forcedSink], fields=[a, b], 
upsertMaterialize=[true])
++- Calc(select=[a, b])
+   +- DataStreamScan(table=[[default_catalog, default_database, MyTable]], 
fields=[a, b, c])
+
+== Optimized Execution Plan ==
+Sink(table=[default_catalog.default_database.forcedSink], fields=[a, b], 
upsertMaterialize=[true])
++- Calc(select=[a, b])
+   +- DataStreamScan(table=[[default_catalog, default_database, MyTable]], 
fields=[a, b, c])
+
+== Physical Execution Plan ==
+{
+  "nodes" : [ {
+    "id" : ,
+    "type" : "Source: Collection Source",
+    "pact" : "Data Source",
+    "contents" : "Source: Collection Source",
+    "parallelism" : 1
+  }, {
+    "id" : ,
+    "type" : "SourceConversion[]",
+    "pact" : "Operator",
+    "contents" : 
"[]:SourceConversion(table=[default_catalog.default_database.MyTable], 
fields=[a, b, c])",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "FORWARD",
+      "side" : "second"
+    } ]
+  }, {
+    "id" : ,
+    "type" : "Calc[]",
+    "pact" : "Operator",
+    "contents" : "[]:Calc(select=[a, b])",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "FORWARD",
+      "side" : "second"
+    } ]
+  }, {
+    "id" : ,
+    "type" : "ConstraintEnforcer[]",
+    "pact" : "Operator",
+    "contents" : "[]:ConstraintEnforcer[NotNullEnforcer(fields=[a])]",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "FORWARD",
+      "side" : "second"
+    } ]
+  }, {
+    "id" : ,
+    "type" : "Sink: forcedSink[]",
+    "pact" : "Data Sink",
+    "contents" : "[]:Sink(table=[default_catalog.default_database.forcedSink], 
fields=[a, b], upsertMaterialize=[true])",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "FORWARD",
+      "side" : "second"
+    } ]
+  } ]
+}]]>
+    </Resource>
+  </TestCase>
+  <TestCase name="testForcedMaterializeWithUpdatingInput">
+    <Resource name="explain">
+      <![CDATA[== Abstract Syntax Tree ==
+LogicalSink(table=[default_catalog.default_database.forcedSinkWithCount], 
fields=[c, EXPR$1])
++- LogicalAggregate(group=[{0}], EXPR$1=[COUNT()])
+   +- LogicalProject(c=[$2])
+      +- LogicalTableScan(table=[[default_catalog, default_database, MyTable]])
+
+== Optimized Physical Plan ==
+Sink(table=[default_catalog.default_database.forcedSinkWithCount], fields=[c, 
EXPR$1], upsertMaterialize=[true])
++- GroupAggregate(groupBy=[c], select=[c, COUNT(*) AS EXPR$1])
+   +- Exchange(distribution=[hash[c]])
+      +- Calc(select=[c])
+         +- DataStreamScan(table=[[default_catalog, default_database, 
MyTable]], fields=[a, b, c])
+
+== Optimized Execution Plan ==
+Sink(table=[default_catalog.default_database.forcedSinkWithCount], fields=[c, 
EXPR$1], upsertMaterialize=[true])
++- GroupAggregate(groupBy=[c], select=[c, COUNT(*) AS EXPR$1])
+   +- Exchange(distribution=[hash[c]])
+      +- Calc(select=[c])
+         +- DataStreamScan(table=[[default_catalog, default_database, 
MyTable]], fields=[a, b, c])
+
+== Physical Execution Plan ==
+{
+  "nodes" : [ {
+    "id" : ,
+    "type" : "Source: Collection Source",
+    "pact" : "Data Source",
+    "contents" : "Source: Collection Source",
+    "parallelism" : 1
+  }, {
+    "id" : ,
+    "type" : "SourceConversion[]",
+    "pact" : "Operator",
+    "contents" : 
"[]:SourceConversion(table=[default_catalog.default_database.MyTable], 
fields=[a, b, c])",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "FORWARD",
+      "side" : "second"
+    } ]
+  }, {
+    "id" : ,
+    "type" : "Calc[]",
+    "pact" : "Operator",
+    "contents" : "[]:Calc(select=[c])",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "FORWARD",
+      "side" : "second"
+    } ]
+  }, {
+    "id" : ,
+    "type" : "GroupAggregate[]",
+    "pact" : "Operator",
+    "contents" : "[]:GroupAggregate(groupBy=[c], select=[c, COUNT(*) AS 
EXPR$1])",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "HASH",
+      "side" : "second"
+    } ]
+  }, {
+    "id" : ,
+    "type" : "ConstraintEnforcer[]",
+    "pact" : "Operator",
+    "contents" : "[]:ConstraintEnforcer[NotNullEnforcer(fields=[c])]",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "FORWARD",
+      "side" : "second"
+    } ]
+  }, {
+    "id" : ,
+    "type" : "SinkMaterializer[]",
+    "pact" : "Operator",
+    "contents" : "[]:SinkMaterializer(pk=[c])",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "HASH",
+      "side" : "second"
+    } ]
+  }, {
+    "id" : ,
+    "type" : "Sink: forcedSinkWithCount[]",
+    "pact" : "Data Sink",
+    "contents" : 
"[]:Sink(table=[default_catalog.default_database.forcedSinkWithCount], 
fields=[c, EXPR$1], upsertMaterialize=[true])",
+    "parallelism" : 1,
+    "predecessors" : [ {
+      "id" : ,
+      "ship_strategy" : "FORWARD",
+      "side" : "second"
+    } ]
+  } ]
+}]]>
+    </Resource>
+  </TestCase>
   <TestCase name="testInjectiveCastPreservesUpsertKey">
     <Resource name="ast">
       <![CDATA[
diff --git 
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/TableSinkTest.scala
 
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/TableSinkTest.scala
index d3e46cfc7bd..a2bad91048f 100644
--- 
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/TableSinkTest.scala
+++ 
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/plan/stream/sql/TableSinkTest.scala
@@ -905,6 +905,51 @@ class TableSinkTest extends TableTestBase {
     util.verifyRelPlan(stmtSet, ExplainDetail.CHANGELOG_MODE)
   }
 
+  @Test
+  def testForcedMaterializeWithAppendOnlyInput(): Unit = {
+    util.getStreamEnv.setParallelism(1)
+    util.tableEnv.getConfig.set(
+      ExecutionConfigOptions.TABLE_EXEC_SINK_UPSERT_MATERIALIZE,
+      ExecutionConfigOptions.UpsertMaterialize.FORCE)
+    util.addTable(s"""
+                     |CREATE TABLE forcedSink (
+                     |  `a` INT,
+                     |  `b` BIGINT,
+                     |  PRIMARY KEY (a) NOT ENFORCED
+                     |) WITH (
+                     |  'connector' = 'values',
+                     |  'sink-insert-only' = 'false'
+                     |)
+                     |""".stripMargin)
+    val stmtSet = util.tableEnv.createStatementSet()
+    stmtSet.addInsertSql("INSERT INTO forcedSink SELECT a, b FROM MyTable")
+    // There is nothing to materialize, so the plan must not contain a 
SinkMaterializer.
+    util.verifyExplain(stmtSet, ExplainDetail.JSON_EXECUTION_PLAN)
+  }
+
+  @Test
+  def testForcedMaterializeWithUpdatingInput(): Unit = {
+    util.getStreamEnv.setParallelism(1)
+    util.tableEnv.getConfig.set(
+      ExecutionConfigOptions.TABLE_EXEC_SINK_UPSERT_MATERIALIZE,
+      ExecutionConfigOptions.UpsertMaterialize.FORCE)
+    util.addTable(s"""
+                     |CREATE TABLE forcedSinkWithCount (
+                     |  `c` STRING,
+                     |  `cnt` BIGINT,
+                     |  PRIMARY KEY (c) NOT ENFORCED
+                     |) WITH (
+                     |  'connector' = 'values',
+                     |  'sink-insert-only' = 'false'
+                     |)
+                     |""".stripMargin)
+    val stmtSet = util.tableEnv.createStatementSet()
+    stmtSet.addInsertSql(
+      "INSERT INTO forcedSinkWithCount SELECT c, COUNT(*) FROM MyTable GROUP 
BY c")
+    // The upsert key already matches the primary key, so only FORCE asks for 
a SinkMaterializer.
+    util.verifyExplain(stmtSet, ExplainDetail.JSON_EXECUTION_PLAN)
+  }
+
   @Test
   def testInjectiveCastPreservesUpsertKey(): Unit = {
     // Aggregation produces upsert stream with key (a).

Reply via email to