From 589bd9f047038c22193d136a6a8679d93c7a1eb9 Mon Sep 17 00:00:00 2001 From: Sergey Nuyanzin Date: Mon, 31 Aug 2026 10:30:26 +0200 Subject: [PATCH 1/3] [FLINK-40528][table] Make codegen tolerat to partial deletes --- .../codegen/calls/ScalarOperatorGens.scala | 4 +- .../exec/stream/DeletesByKeyPrograms.java | 303 ++++++++++++++++++ .../stream/DeletesByKeySemanticTests.java | 9 +- 3 files changed, 313 insertions(+), 3 deletions(-) diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala index 4e3a79cf7a21a..15b074791841c 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala @@ -1226,7 +1226,7 @@ object ScalarOperatorGens { val tpe = fieldTypes(idx) if (element.literal) { "" - } else if (tpe.isNullable) { + } else if (tpe.isNullable || element.nullTerm != NEVER_NULL) { s""" |${element.code} |if (${element.nullTerm}) { @@ -1275,7 +1275,7 @@ object ScalarOperatorGens { .map { case (element, idx) => val tpe = fieldTypes(idx) - if (tpe.isNullable) { + if (tpe.isNullable || element.nullTerm != NEVER_NULL) { s""" |${element.code} |if (${element.nullTerm}) { 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..e06f82829f854 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 @@ -24,6 +24,8 @@ import org.apache.flink.types.Row; import org.apache.flink.types.RowKind; +import java.util.Map; + /** * Tests for verifying semantic of operations when sources produce deletes by key only and the sink * can accept deletes by key only as well. @@ -293,5 +295,306 @@ public final class DeletesByKeyPrograms { "INSERT INTO sink_t SELECT l.id, r.name, l.`value` FROM left_t l JOIN right_t r ON l.id = r.id") .build(); + /** + * A delete-by-key tombstone carries null for a NOT NULL ARRAY wrapped in a {@code ROW(...)} + * projection. + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ARRAY = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-nested-not-null-array", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL ARRAY column wrapped in a ROW(...) projection; validates" + + " that row construction does not fail on the null value") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "arr ARRAY NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, new Integer[] {1, 2}), + Row.ofKind(RowKind.INSERT, 2, new Integer[] {3}), + // Delete by key: NOT NULL array column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind( + RowKind.UPDATE_AFTER, 2, new Integer[] {3, 4})) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "r ROW>") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, +I[1, [1, 2]]]", + "+I[2, +I[2, [3]]]", + "-D[1, +I[1, null]]", + "+U[2, +I[2, [3, 4]]]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, ROW(id, arr) FROM source_t") + .build(); + + /** + * Same as the ARRAY variant but for a NOT NULL {@code MAP} column wrapped in {@code ROW(...)}. + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_MAP = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-nested-not-null-map", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL MAP column wrapped in a ROW(...) projection") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "m MAP NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, Map.of(1, 10)), + Row.ofKind(RowKind.INSERT, 2, Map.of(2, 20)), + // Delete by key: NOT NULL map column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, Map.of(2, 30))) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "r ROW>") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, +I[1, {1=10}]]", + "+I[2, +I[2, {2=20}]]", + "-D[1, +I[1, null]]", + "+U[2, +I[2, {2=30}]]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, ROW(id, m) FROM source_t") + .build(); + + /** + * Same as the ARRAY variant but for a NOT NULL {@code ROW} column wrapped in {@code ROW(...)}. + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ROW = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-nested-not-null-row", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL ROW column wrapped in a ROW(...) projection") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "nested ROW NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, Row.of(1, 10)), + Row.ofKind(RowKind.INSERT, 2, Row.of(2, 20)), + // Delete by key: NOT NULL row column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, Row.of(2, 30))) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "r ROW>") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, +I[1, +I[1, 10]]]", + "+I[2, +I[2, +I[2, 20]]]", + "-D[1, +I[1, null]]", + "+U[2, +I[2, +I[2, 30]]]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, ROW(id, nested) FROM source_t") + .build(); + + /** Same as the ARRAY variant but for a NOT NULL {@code ARRAY} column. */ + public static final TableTestProgram + INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ARRAY_OF_ROW = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-nested-not-null-array-of-row", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL ARRAY column wrapped in a ROW(...) projection") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "arr ARRAY> NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind( + RowKind.INSERT, + 1, + new Row[] {Row.of(1, 10)}), + Row.ofKind( + RowKind.INSERT, + 2, + new Row[] {Row.of(2, 20)}), + // Delete by key: NOT NULL array column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind( + RowKind.UPDATE_AFTER, + 2, + new Row[] {Row.of(2, 30)})) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "r ROW>>") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, +I[1, [+I[1, 10]]]]", + "+I[2, +I[2, [+I[2, 20]]]]", + "-D[1, +I[1, null]]", + "+U[2, +I[2, [+I[2, 30]]]]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, ROW(id, arr) FROM source_t") + .build(); + + /** + * Same shape as the ARRAY variant but for a NOT NULL {@code STRING} column. The reference-typed + * String path does not fail today, but the case is kept for coverage. + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_STRING = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-nested-not-null-string", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL STRING column wrapped in a ROW(...) projection") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", "s STRING NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, "a"), + Row.ofKind(RowKind.INSERT, 2, "b"), + // Delete by key: NOT NULL string column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, "c")) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "r ROW") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, +I[1, a]]", + "+I[2, +I[2, b]]", + "-D[1, +I[1, null]]", + "+U[2, +I[2, c]]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, ROW(id, s) FROM source_t") + .build(); + + /** + * Same shape as the ARRAY variant but for a NOT NULL primitive {@code INT} column. The + * primitive path is already guarded; kept for coverage. + */ + public static final TableTestProgram + INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_PRIMITIVE = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-nested-not-null-primitive", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL INT column wrapped in a ROW(...) projection") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "v INT NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, 10), + Row.ofKind(RowKind.INSERT, 2, 20), + // Delete by key: NOT NULL int column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "r ROW") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, +I[1, 10]]", + "+I[2, +I[2, 20]]", + "-D[1, +I[1, null]]", + "+U[2, +I[2, 30]]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, ROW(id, v) FROM source_t") + .build(); + + /** + * A LEFT JOIN whose probe (left) side produces a delete-by-key tombstone carrying null for a + * NOT NULL ARRAY column that is wrapped in a ROW(...) projection. + */ + public static final TableTestProgram JOIN_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ARRAY = + TableTestProgram.of( + "join-delete-on-key-with-nested-not-null-array", + "No ChangelogNormalize: probe-side delete-by-key tombstone carries null" + + " for a NOT NULL ARRAY column wrapped in a ROW(...) projection" + + " across a join") + .setupTableSource( + SourceTestStep.newBuilder("left_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "arr ARRAY NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, new Integer[] {1, 2}), + Row.ofKind(RowKind.INSERT, 2, new Integer[] {3}), + Row.ofKind(RowKind.INSERT, 3, new Integer[] {5}), + // Delete by key: NOT NULL array column is null + Row.ofKind(RowKind.DELETE, 1, null), + Row.ofKind(RowKind.UPDATE_AFTER, 3, new Integer[] {6})) + .build()) + .setupTableSource( + SourceTestStep.newBuilder("right_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.INSERT, 2, "Bob"), + Row.ofKind(RowKind.INSERT, 3, "Emily"), + Row.ofKind(RowKind.DELETE, 1, null), + Row.ofKind(RowKind.UPDATE_AFTER, 2, "BOB")) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "r ROW>", + "name STRING") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .testMaterializedData() + .consumedValues( + "+I[2, +I[2, [3]], BOB]", "+I[3, +I[3, [6]], Emily]") + .build()) + .runSql( + "INSERT INTO sink_t SELECT l.id, ROW(l.id, l.arr), r.name" + + " FROM left_t l JOIN right_t r ON l.id = r.id") + .build(); + private DeletesByKeyPrograms() {} } 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..c0571ee093cad 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 @@ -34,6 +34,13 @@ public List programs() { DeletesByKeyPrograms.INSERT_SELECT_FULL_DELETE_FULL_DELETE, DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_DELETE_BY_KEY_WITH_PROJECTION, DeletesByKeyPrograms.JOIN_INTO_FULL_DELETES, - DeletesByKeyPrograms.JOIN_INTO_DELETES_BY_KEY); + DeletesByKeyPrograms.JOIN_INTO_DELETES_BY_KEY, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ARRAY, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_MAP, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ROW, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ARRAY_OF_ROW, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_STRING, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_PRIMITIVE, + DeletesByKeyPrograms.JOIN_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ARRAY); } } From 749d1c42ae1e28e5bbeac5ffe56a68b55829a16c Mon Sep 17 00:00:00 2001 From: Sergey Nuyanzin Date: Tue, 1 Sep 2026 13:04:06 +0200 Subject: [PATCH 2/3] Address feedback --- .../table/planner/codegen/CodeGenUtils.scala | 2 +- .../planner/codegen/GeneratedExpression.scala | 3 + .../planner/codegen/JsonGenerateUtils.scala | 17 +- .../codegen/calls/ScalarOperatorGens.scala | 6 +- .../exec/stream/DeletesByKeyPrograms.java | 228 +++++++++++++++++- .../stream/DeletesByKeySemanticTests.java | 8 +- 6 files changed, 249 insertions(+), 15 deletions(-) diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala index 739858689bf55..14567319488a6 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala @@ -600,7 +600,7 @@ object CodeGenUtils { s"$rowTerm.setNullAt($indexTerm)" } - if (fieldType.isNullable) { + if (fieldType.isNullable || !fieldExpr.isProvenNotNull) { s""" |${fieldExpr.code} |if (${fieldExpr.nullTerm}) { diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/GeneratedExpression.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/GeneratedExpression.scala index 325a7608c2706..6fdf90d5aa7b2 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/GeneratedExpression.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/GeneratedExpression.scala @@ -56,6 +56,9 @@ case class GeneratedExpression( */ def literal: Boolean = literalValue.isDefined + /** Whether this expression is statically proven never to be null at runtime. */ + def isProvenNotNull: Boolean = nullTerm == GeneratedExpression.NEVER_NULL + /** * Copy result term to target term if the reference is changed. Note: We must ensure that the * target can only be copied out, so that its object is definitely a brand new reference, not the diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/JsonGenerateUtils.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/JsonGenerateUtils.scala index a64a97727ccd8..aa9b8f239d5dd 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/JsonGenerateUtils.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/JsonGenerateUtils.scala @@ -142,15 +142,14 @@ object JsonGenerateUtils { val valueNodeTerm = createNodeTerm(ctx, fieldAccessTerm, fieldType) - if (fieldType.isNullable) { - s""" - |$containerTerm.isNullAt($indexTerm) ? - | (${className[JsonNode]}) $nodeFactoryTerm.nullNode() : - | (${className[JsonNode]}) $valueNodeTerm - |""".stripMargin - } else { - valueNodeTerm - } + // The declared type may be NOT NULL while the runtime container still carries null, e.g. a + // delete-by-key tombstone. Always guard on isNullAt so such nulls serialize as JSON null + // instead of reading a primitive default. + s""" + |$containerTerm.isNullAt($indexTerm) ? + | (${className[JsonNode]}) $nodeFactoryTerm.nullNode() : + | (${className[JsonNode]}) $valueNodeTerm + |""".stripMargin } /** diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala index 15b074791841c..c7b63fe5ab065 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala @@ -1226,7 +1226,7 @@ object ScalarOperatorGens { val tpe = fieldTypes(idx) if (element.literal) { "" - } else if (tpe.isNullable || element.nullTerm != NEVER_NULL) { + } else if (tpe.isNullable || !element.isProvenNotNull) { s""" |${element.code} |if (${element.nullTerm}) { @@ -1275,7 +1275,7 @@ object ScalarOperatorGens { .map { case (element, idx) => val tpe = fieldTypes(idx) - if (tpe.isNullable || element.nullTerm != NEVER_NULL) { + if (tpe.isNullable || !element.isProvenNotNull) { s""" |${element.code} |if (${element.nullTerm}) { @@ -1350,7 +1350,7 @@ object ScalarOperatorGens { case (element, idx) => if (element.literal) { "" - } else if (elementType.isNullable) { + } else if (elementType.isNullable || !element.isProvenNotNull) { s""" |${element.code} |if (${element.nullTerm}) { 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 e06f82829f854..369f690c5e0e6 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 @@ -113,7 +113,7 @@ public final class DeletesByKeyPrograms { .consumedValues( "+I[1, Alice, 12]", "+I[2, Bob, 22]", - "-D[1, , -1]", + "-D[1, null, null]", "+U[2, Bob, 32]") .build()) .runSql("INSERT INTO sink_t SELECT id, name, `value` + 2 FROM source_t") @@ -596,5 +596,231 @@ public final class DeletesByKeyPrograms { + " FROM left_t l JOIN right_t r ON l.id = r.id") .build(); + /** + * A delete-by-key tombstone carries null for a NOT NULL {@code INT} column that is used as an + * element of an {@code ARRAY[...]} literal (element type {@code INT NOT NULL}) wrapped in a + * {@code ROW(...)} projection. Validates that array construction sets the element to null + * instead of writing the primitive default. + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_NOT_NULL_ARRAY_LITERAL = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-not-null-array-literal", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL INT column used as an element of an ARRAY[...] literal") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "v INT NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, 10), + Row.ofKind(RowKind.INSERT, 2, 20), + // Delete by key: NOT NULL int column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "r ROW>") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, +I[1, [10, 99]]]", + "+I[2, +I[2, [20, 99]]]", + "-D[1, +I[1, [null, 99]]]", + "+U[2, +I[2, [30, 99]]]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, ROW(id, ARRAY[v, 99]) FROM source_t") + .build(); + + /** + * Same as the ARRAY literal variant but for a {@code MAP[...]} literal whose value comes from a + * NOT NULL {@code INT} column (value type {@code INT NOT NULL}). + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_NOT_NULL_MAP_LITERAL = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-not-null-map-literal", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL INT column used as a value of a MAP[...] literal") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "v INT NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, 10), + Row.ofKind(RowKind.INSERT, 2, 20), + // Delete by key: NOT NULL int column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id INT PRIMARY KEY NOT ENFORCED", + "r ROW>") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, +I[1, {99=10}]]", + "+I[2, +I[2, {99=20}]]", + "-D[1, +I[1, {99=null}]]", + "+U[2, +I[2, {99=30}]]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, ROW(id, MAP[99, v]) FROM source_t") + .build(); + + /** + * A delete-by-key tombstone carries null for a NOT NULL {@code INT} column that is nested in a + * {@code ROW(...)} serialized by {@code JSON_OBJECT}. The nested field is NOT NULL, so the JSON + * row converter must still guard on the runtime null and emit {@code null} instead of reading + * the primitive default. + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_JSON_OBJECT_NESTED_ROW = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-json-object-nested-row", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL INT column nested in a ROW serialized by JSON_OBJECT") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "v INT NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, 10), + Row.ofKind(RowKind.INSERT, 2, 20), + // Delete by key: NOT NULL int column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "j STRING") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, {\"r\":{\"EXPR$0\":10}}]", + "+I[2, {\"r\":{\"EXPR$0\":20}}]", + "-D[1, {\"r\":{\"EXPR$0\":null}}]", + "+U[2, {\"r\":{\"EXPR$0\":30}}]") + .build()) + .runSql( + "INSERT INTO sink_t SELECT id, JSON_OBJECT('r' VALUE ROW(v)) FROM source_t") + .build(); + + /** + * Same as the JSON_OBJECT nested ROW variant but for a NOT NULL {@code INT} element nested in + * an {@code ARRAY[...]} serialized by {@code JSON_OBJECT}. + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_JSON_OBJECT_NESTED_ARRAY = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-json-object-nested-array", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL INT element nested in an ARRAY serialized by JSON_OBJECT") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "v INT NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, 10), + Row.ofKind(RowKind.INSERT, 2, 20), + // Delete by key: NOT NULL int column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "j STRING") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, {\"a\":[10,99]}]", + "+I[2, {\"a\":[20,99]}]", + "-D[1, {\"a\":[null,99]}]", + "+U[2, {\"a\":[30,99]}]") + .build()) + .runSql( + "INSERT INTO sink_t SELECT id, JSON_OBJECT('a' VALUE ARRAY[v, 99]) FROM source_t") + .build(); + + /** + * Same as the JSON_OBJECT nested ROW variant but for a NOT NULL {@code INT} value nested in a + * {@code MAP[...]} serialized by {@code JSON_OBJECT}. + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_JSON_OBJECT_NESTED_MAP = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-json-object-nested-map", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL INT value nested in a MAP serialized by JSON_OBJECT") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "v INT NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, 10), + Row.ofKind(RowKind.INSERT, 2, 20), + // Delete by key: NOT NULL int column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "j STRING") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, {\"m\":{\"k\":10}}]", + "+I[2, {\"m\":{\"k\":20}}]", + "-D[1, {\"m\":{\"k\":null}}]", + "+U[2, {\"m\":{\"k\":30}}]") + .build()) + .runSql( + "INSERT INTO sink_t SELECT id, JSON_OBJECT('m' VALUE MAP['k', v]) FROM source_t") + .build(); + + /** + * A delete-by-key tombstone carries null for a NOT NULL {@code INT} column that is CAST to + * another primitive type. The cast framework skips the runtime null guard when the input type + * is NOT NULL and the target is a primitive Java type, so it reads the primitive default (0) + * instead of producing null. + */ + public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_NOT_NULL_CAST = + TableTestProgram.of( + "select-delete-on-key-to-delete-on-key-with-not-null-cast", + "No ChangelogNormalize: a delete-by-key tombstone carries null for a NOT" + + " NULL INT column that is CAST to BIGINT") + .setupTableSource( + SourceTestStep.newBuilder("source_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "v INT NOT NULL") + .addOption("changelog-mode", "I,UA,D") + .addOption("source.produces-delete-by-key", "true") + .producedValues( + Row.ofKind(RowKind.INSERT, 1, 10), + Row.ofKind(RowKind.INSERT, 2, 20), + // Delete by key: NOT NULL int column is null + Row.ofKind(RowKind.DELETE, 1, null), + // Update after only + Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("id INT PRIMARY KEY NOT ENFORCED", "v BIGINT") + .addOption("changelog-mode", "I,UA,D") + .addOption("sink.supports-delete-by-key", "true") + .consumedValues( + "+I[1, 10]", "+I[2, 20]", "-D[1, null]", "+U[2, 30]") + .build()) + .runSql("INSERT INTO sink_t SELECT id, CAST(v AS BIGINT) FROM source_t") + .build(); + private DeletesByKeyPrograms() {} } 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 c0571ee093cad..b43a28fdff881 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 @@ -41,6 +41,12 @@ public List programs() { DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ARRAY_OF_ROW, DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_STRING, DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_PRIMITIVE, - DeletesByKeyPrograms.JOIN_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ARRAY); + DeletesByKeyPrograms.JOIN_DELETE_BY_KEY_WITH_NESTED_NOT_NULL_ARRAY, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NOT_NULL_ARRAY_LITERAL, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NOT_NULL_MAP_LITERAL, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_JSON_OBJECT_NESTED_ROW, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_JSON_OBJECT_NESTED_ARRAY, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_JSON_OBJECT_NESTED_MAP, + DeletesByKeyPrograms.INSERT_SELECT_DELETE_BY_KEY_WITH_NOT_NULL_CAST); } } From 894a8da60a8e31a05ae249025c61e1f562ab7dad Mon Sep 17 00:00:00 2001 From: Sergey Nuyanzin Date: Wed, 2 Sep 2026 10:59:41 +0200 Subject: [PATCH 3/3] Address feedback --- .../exec/stream/DeletesByKeyPrograms.java | 20 ------------------- 1 file changed, 20 deletions(-) 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 369f690c5e0e6..493c9680cac3c 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 @@ -521,9 +521,7 @@ public final class DeletesByKeyPrograms { .producedValues( Row.ofKind(RowKind.INSERT, 1, 10), Row.ofKind(RowKind.INSERT, 2, 20), - // Delete by key: NOT NULL int column is null Row.ofKind(RowKind.DELETE, 1, null), - // Update after only Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) .build()) .setupTableSink( @@ -615,9 +613,7 @@ public final class DeletesByKeyPrograms { .producedValues( Row.ofKind(RowKind.INSERT, 1, 10), Row.ofKind(RowKind.INSERT, 2, 20), - // Delete by key: NOT NULL int column is null Row.ofKind(RowKind.DELETE, 1, null), - // Update after only Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) .build()) .setupTableSink( @@ -653,9 +649,7 @@ public final class DeletesByKeyPrograms { .producedValues( Row.ofKind(RowKind.INSERT, 1, 10), Row.ofKind(RowKind.INSERT, 2, 20), - // Delete by key: NOT NULL int column is null Row.ofKind(RowKind.DELETE, 1, null), - // Update after only Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) .build()) .setupTableSink( @@ -693,9 +687,7 @@ public final class DeletesByKeyPrograms { .producedValues( Row.ofKind(RowKind.INSERT, 1, 10), Row.ofKind(RowKind.INSERT, 2, 20), - // Delete by key: NOT NULL int column is null Row.ofKind(RowKind.DELETE, 1, null), - // Update after only Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) .build()) .setupTableSink( @@ -730,9 +722,7 @@ public final class DeletesByKeyPrograms { .producedValues( Row.ofKind(RowKind.INSERT, 1, 10), Row.ofKind(RowKind.INSERT, 2, 20), - // Delete by key: NOT NULL int column is null Row.ofKind(RowKind.DELETE, 1, null), - // Update after only Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) .build()) .setupTableSink( @@ -767,9 +757,7 @@ public final class DeletesByKeyPrograms { .producedValues( Row.ofKind(RowKind.INSERT, 1, 10), Row.ofKind(RowKind.INSERT, 2, 20), - // Delete by key: NOT NULL int column is null Row.ofKind(RowKind.DELETE, 1, null), - // Update after only Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) .build()) .setupTableSink( @@ -787,12 +775,6 @@ public final class DeletesByKeyPrograms { "INSERT INTO sink_t SELECT id, JSON_OBJECT('m' VALUE MAP['k', v]) FROM source_t") .build(); - /** - * A delete-by-key tombstone carries null for a NOT NULL {@code INT} column that is CAST to - * another primitive type. The cast framework skips the runtime null guard when the input type - * is NOT NULL and the target is a primitive Java type, so it reads the primitive default (0) - * instead of producing null. - */ public static final TableTestProgram INSERT_SELECT_DELETE_BY_KEY_WITH_NOT_NULL_CAST = TableTestProgram.of( "select-delete-on-key-to-delete-on-key-with-not-null-cast", @@ -806,9 +788,7 @@ public final class DeletesByKeyPrograms { .producedValues( Row.ofKind(RowKind.INSERT, 1, 10), Row.ofKind(RowKind.INSERT, 2, 20), - // Delete by key: NOT NULL int column is null Row.ofKind(RowKind.DELETE, 1, null), - // Update after only Row.ofKind(RowKind.UPDATE_AFTER, 2, 30)) .build()) .setupTableSink(