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).