From a06db4147ad51e1d13515b1e1278867581d05216 Mon Sep 17 00:00:00 2001 From: Gustavo de Morais Date: Wed, 5 Aug 2026 18:24:29 +0200 Subject: [PATCH] [FLINK-40334][table] Insert-only input should stay unmaterialized under FORCE FLINK-38928 dropped the guard that skipped the sink upsert materializer for insert-only input, because DO ERROR and DO NOTHING need the operator to compare inserts that share a primary key. Without it, an insert-only query into a sink with a primary key gets a SinkUpsertMaterializer and a keyed shuffle whenever upsert materialization is forced, even though there are no changes to reconcile. Restore the guard and pass the conflict strategy into the sink transformation so that only DO ERROR and DO NOTHING keep materializing insert-only input. The decision stays at translation time on purpose. Plans compiled before the clause existed persist requireUpsertMaterialize = true for this case, so suppressing the operator in the planner instead would make it appear on restore. --- .../plan/nodes/exec/batch/BatchExecSink.java | 1 + .../nodes/exec/common/CommonExecSink.java | 24 ++- .../nodes/exec/stream/StreamExecSink.java | 14 +- .../planner/plan/stream/sql/TableSinkTest.xml | 173 ++++++++++++++++++ .../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 8a1dd04bcfe87..cb267e24e80d1 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 @@ protected Transformation translateToPlanInternal( 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 31c0bc4074feb..3738840ab2761 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.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.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 @@ protected Transformation createSinkTransformation( 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 @@ protected Transformation createSinkTransformation( Optional 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 sinkTransform = applyConstraintValidations(inputTransform, config, persistedRowType); @@ -199,10 +209,10 @@ protected Transformation createSinkTransformation( primaryKeys, sinkParallelism, inputParallelism, - upsertMaterialize); + needMaterialization); } - if (upsertMaterialize) { + if (needMaterialization) { sinkTransform = applyUpsertMaterialize( sinkTransform, @@ -246,6 +256,14 @@ protected Transformation createSinkTransformation( 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 22dce136a0e24..acc7fa10b2268 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.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 @@ protected Transformation translateToPlanInternal( tableSink, rowtimeFieldIndex, upsertMaterialize, - inputUpsertKey); + inputUpsertKey, + conflictStrategy); } @Override @@ -359,7 +359,7 @@ protected Transformation applyUpsertMaterialize( // This assigns the current watermark as the timestamp to each record, // which is required for the WatermarkCompactingSinkMaterializer to work correctly Transformation transformForMaterializer = inputTransform; - if (isErrorOrNothingConflictStrategy()) { + if (detectsDuplicateKeys(conflictStrategy)) { // Use input parallelism to preserve watermark semantics transformForMaterializer = ExecNodeUtil.createOneInputTransformation( @@ -410,7 +410,7 @@ private OneInputStreamOperator createSumOperator( 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 @@ private OneInputStreamOperator createSumOperator( 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 90d4c28486542..486e4d045ec2c 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]) ]]> + + + + + + + + + +