From f1a03f5ecb177564b469b854dd48d865f5ff91a8 Mon Sep 17 00:00:00 2001 From: Au_Miner <358671982@qq.com> Date: Thu, 30 Jul 2026 17:05:29 +0800 Subject: [PATCH 1/2] [FLINK-40303][table] Skip non-key NOT NULL checks for key-only deletes sink AI-Contributed/Feature: 0/46 AI-Contributed/UT: 0/30 --- .../nodes/exec/common/CommonExecSink.java | 9 ++++-- .../exec/stream/DeletesByKeyPrograms.java | 29 +++++++++++++++++++ .../stream/DeletesByKeySemanticTests.java | 1 + .../ConstraintEnforcerExecutor.java | 23 +++++++++++++-- .../sink/constraint/NotNullConstraint.java | 14 ++++++++- 5 files changed, 69 insertions(+), 7 deletions(-) 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..8162398a02c65 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 @@ -188,7 +188,7 @@ protected Transformation createSinkTransformation( TableLineageUtils.extractLineageDataset(outputObject); Transformation sinkTransform = - applyConstraintValidations(inputTransform, config, persistedRowType); + applyConstraintValidations(inputTransform, config, persistedRowType, primaryKeys); if (hasPk) { sinkTransform = @@ -252,14 +252,17 @@ protected Transformation createSinkTransformation( private Transformation applyConstraintValidations( Transformation inputTransform, ExecNodeConfig config, - RowType physicalRowType) { + RowType physicalRowType, + int[] primaryKeys) { final Optional enforcerExecutor = ConstraintEnforcerExecutor.create( physicalRowType, config.get(ExecutionConfigOptions.TABLE_EXEC_SINK_NOT_NULL_ENFORCER), config.get(ExecutionConfigOptions.TABLE_EXEC_SINK_TYPE_LENGTH_ENFORCER), config.get( - ExecutionConfigOptions.TABLE_EXEC_SINK_NESTED_CONSTRAINT_ENFORCER)); + ExecutionConfigOptions.TABLE_EXEC_SINK_NESTED_CONSTRAINT_ENFORCER), + inputChangelogMode.keyOnlyDeletes(), + primaryKeys); return enforcerExecutor .map( diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java index a66837d1c4542..8dc396a102bcd 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java @@ -75,6 +75,35 @@ public final class DeletesByKeyPrograms { .runSql("INSERT INTO sink_t SELECT id, name, `value` FROM source_t") .build(); + public static final TableTestProgram + INSERT_SELECT_DELETE_BY_KEY_DELETE_BY_KEY_WITH_NOT_NULL_SINK = + TableTestProgram.of( + "select-delete-on-key-to-not-null-sink", + "No ChangelogNormalize: validates that key-only deletes can contain null" + + " non-key fields when writing to a sink with NOT NULL constraints") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "name STRING") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, "Alice"), + Row.ofKind(RowKind.DELETE, 1, null)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "name STRING NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues("+I[1, Alice]", "-D[1, null]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, name FROM source_t") + .build(); + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_DELETE_BY_KEY_WITH_PROJECTION = TableTestProgram.of( "select-delete-on-key-to-delete-on-key-with-projection", diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeySemanticTests.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeySemanticTests.java index eba93c8bba9ce..e3d5db38e061b 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeySemanticTests.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeySemanticTests.java @@ -30,6 +30,7 @@ public class DeletesByKeySemanticTests extends SemanticTestBase { public List programs() { return List.of( DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_DELETE_BY_KEY, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_DELETE_BY_KEY_WITH_NOT_NULL_SINK, DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_FULL_DELETE, DeletesByKeyPrograms.INSERT_SELECT_FULL_DELETE_FULL_DELETE, DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_DELETE_BY_KEY_WITH_PROJECTION, diff --git a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/constraint/ConstraintEnforcerExecutor.java b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/constraint/ConstraintEnforcerExecutor.java index b8377af468e98..80481eb3d10f5 100644 --- a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/constraint/ConstraintEnforcerExecutor.java +++ b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/constraint/ConstraintEnforcerExecutor.java @@ -70,10 +70,16 @@ public static Optional create( final RowType physicalType, final NotNullEnforcer notNullEnforcer, final TypeLengthEnforcer typeLengthEnforcer, - final NestedEnforcer nestedConstraints) { + final NestedEnforcer nestedConstraints, + final boolean keyOnlyDeletes, + final int[] primaryKeyIndices) { final Constraint[] topLevelConstraints = createConstraints( - physicalType, notNullEnforcer, typeLengthEnforcer, nestedConstraints); + physicalType, + notNullEnforcer, + typeLengthEnforcer, + nestedConstraints, + keyOnlyDeletes ? primaryKeyIndices : null); return create(topLevelConstraints); } @@ -83,6 +89,16 @@ private static Constraint[] createConstraints( final NotNullEnforcer notNullEnforcer, final TypeLengthEnforcer typeLengthEnforcer, final NestedEnforcer nestedEnforcer) { + return createConstraints( + physicalType, notNullEnforcer, typeLengthEnforcer, nestedEnforcer, null); + } + + private static Constraint[] createConstraints( + final RowType physicalType, + final NotNullEnforcer notNullEnforcer, + final TypeLengthEnforcer typeLengthEnforcer, + final NestedEnforcer nestedEnforcer, + final @Nullable int[] keyOnlyDeleteFieldIndices) { final String[] fieldNames = physicalType.getFieldNames().toArray(new String[0]); final List constraints = new ArrayList<>(); @@ -98,7 +114,8 @@ private static Constraint[] createConstraints( new NotNullConstraint( NotNullEnforcementStrategy.of(notNullEnforcer), notNullFieldIndices, - notNullFieldNames)); + notNullFieldNames, + keyOnlyDeleteFieldIndices)); } if (typeLengthEnforcer != TypeLengthEnforcer.IGNORE) { diff --git a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/constraint/NotNullConstraint.java b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/constraint/NotNullConstraint.java index f2ae509c5b51f..95e13c4cfe4d3 100644 --- a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/constraint/NotNullConstraint.java +++ b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/sink/constraint/NotNullConstraint.java @@ -21,6 +21,9 @@ import org.apache.flink.annotation.Internal; import org.apache.flink.table.api.config.ExecutionConfigOptions; import org.apache.flink.table.data.RowData; +import org.apache.flink.types.RowKind; + +import org.apache.commons.lang3.ArrayUtils; import javax.annotation.Nullable; @@ -32,14 +35,17 @@ final class NotNullConstraint implements Constraint { private final NotNullEnforcementStrategy enforcementStrategy; private final int[] notNullFieldIndices; private final String[] notNullFieldNames; + private final @Nullable int[] keyOnlyDeleteFieldIndices; NotNullConstraint( NotNullEnforcementStrategy enforcementStrategy, int[] notNullFieldIndices, - String[] notNullFieldNames) { + String[] notNullFieldNames, + @Nullable int[] keyOnlyDeleteFieldIndices) { this.enforcementStrategy = enforcementStrategy; this.notNullFieldIndices = notNullFieldIndices; this.notNullFieldNames = notNullFieldNames; + this.keyOnlyDeleteFieldIndices = keyOnlyDeleteFieldIndices; } @Nullable @@ -47,6 +53,12 @@ final class NotNullConstraint implements Constraint { public RowData enforce(RowData input) { for (int i = 0; i < notNullFieldIndices.length; i++) { final int index = notNullFieldIndices[i]; + // Non-key fields have no value semantics in a key-only DELETE. + if (input.getRowKind() == RowKind.DELETE + && keyOnlyDeleteFieldIndices != null + && !ArrayUtils.contains(keyOnlyDeleteFieldIndices, index)) { + continue; + } if (input.isNullAt(index)) { switch (enforcementStrategy) { case ERROR: From ee4dfb9da12cadf8a0ca31c2b4a1d9feb499a2a8 Mon Sep 17 00:00:00 2001 From: Au_Miner <358671982@qq.com> Date: Tue, 4 Aug 2026 14:57:12 +0800 Subject: [PATCH 2/2] trigger ci AI-Contributed/Feature: 0/0 AI-Contributed/UT: 0/2 --- .../planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java index 8dc396a102bcd..c78cfd1674c7b 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java @@ -79,7 +79,7 @@ public final class DeletesByKeyPrograms { INSERT_SELECT_DELETE_BY_KEY_DELETE_BY_KEY_WITH_NOT_NULL_SINK = TableTestProgram.of( "select-delete-on-key-to-not-null-sink", - "No ChangelogNormalize: validates that key-only deletes can contain null" + "No ChangelogNormalize: validates that key-only delete can contain null" + " non-key fields when writing to a sink with NOT NULL constraints") .setupTableSource( SourceTestStep.newBuilder("source_t")