snuyanzin commented on code in PR #29060:
URL: https://github.com/apache/flink/pull/29060#discussion_r3912668876
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java:
##########
@@ -293,5 +295,532 @@ 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<INT> 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<a INT, b ARRAY<INT>>")
+ .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<INT, INT> 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<a INT, b MAP<INT, INT>>")
+ .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<x INT, y INT> 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<a INT, b ROW<x INT, y
INT>>")
+ .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<ROW>}
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<ROW> column wrapped
in a ROW(...) projection")
+ .setupTableSource(
+ SourceTestStep.newBuilder("source_t")
+ .addSchema(
+ "id INT PRIMARY KEY NOT
ENFORCED",
+ "arr ARRAY<ROW<x INT, y
INT>> 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<a INT, b
ARRAY<ROW<x INT, y INT>>>")
+ .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<a INT, b STRING>")
+ .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<a INT, b INT>")
+ .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<INT> 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<a INT, b ARRAY<INT>>",
+ "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();
+
+ /**
+ * 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<a INT, b ARRAY<INT>>")
+ .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<a INT, b MAP<INT, INT>>")
+ .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 =
Review Comment:
> INSERT_SELECT_DELETE_BY_KEY_WITH_NOT_NULL_CAST
this is required
it highlights another finding
flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/DeletesByKeyPrograms.java:
##########
@@ -293,5 +295,532 @@ 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<INT> 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<a INT, b ARRAY<INT>>")
+ .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<INT, INT> 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<a INT, b MAP<INT, INT>>")
+ .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<x INT, y INT> 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<a INT, b ROW<x INT, y
INT>>")
+ .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<ROW>}
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<ROW> column wrapped
in a ROW(...) projection")
+ .setupTableSource(
+ SourceTestStep.newBuilder("source_t")
+ .addSchema(
+ "id INT PRIMARY KEY NOT
ENFORCED",
+ "arr ARRAY<ROW<x INT, y
INT>> 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<a INT, b
ARRAY<ROW<x INT, y INT>>>")
+ .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<a INT, b STRING>")
+ .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<a INT, b INT>")
+ .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<INT> 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<a INT, b ARRAY<INT>>",
+ "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();
+
+ /**
+ * 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<a INT, b ARRAY<INT>>")
+ .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<a INT, b MAP<INT, INT>>")
+ .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 =
Review Comment:
> INSERT_SELECT_DELETE_BY_KEY_WITH_NOT_NULL_CAST
this is required
it highlights another finding
flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]