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 8a1dd04bcfe87f..cb267e24e80d19 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 31c0bc4074febd..3738840ab27612 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 22dce136a0e241..acc7fa10b2268a 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 90d4c28486542e..486e4d045ec2cc 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]) ]]> + + + + + + + + + +