laskoviymishka commented on code in PR #17653:
URL: https://github.com/apache/iceberg/pull/17653#discussion_r4149606646
##########
kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/RecordConverter.java:
##########
@@ -356,12 +359,32 @@ private void evolveSchemaFromConnectSchema(
}
}
+ private void logIfNullForRequiredColumn(
+ Object value, NestedField tableField, boolean hasSchemaUpdates) {
+ if (value == null && tableField.isRequired() && !hasSchemaUpdates) {
+ LOG.warn(
+ "Explicit null value for required column {}; the write will fail.
Consider setting"
Review Comment:
This warning fires on the default path too, not just the preserve-null one.
With `replace-null-with-default=true` (the default), a required column that's
null with no schema default still returns null from `struct.get()`, so existing
deployments start emitting this WARN after upgrading for a case that was silent
before — and `Explicit null value` is the wrong wording there, since nothing
was explicitly overridden, it's just a missing value. I'd reword to `Null value
for required column …`, and if we want to keep the explicit framing, only use
it on the `!replaceNullWithDefault` branch — plus a line in the migration note
that the default config now emits this. While we're here, nothing exercises
this branch with a `NestedField.required` column receiving null, so a quick
test there would guard the condition against an accidental flip.
##########
kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestRecordConverter.java:
##########
@@ -962,6 +988,102 @@ public void testNoSchemaEvolutionStructWithNullValue() {
assertThat(consumer.empty()).isTrue();
}
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ public void testNoSchemaEvolutionStructWithNullValueOfFieldWithDefault(
+ boolean replaceNullWithDefault) {
+ when(config.replaceNullWithDefault()).thenReturn(replaceNullWithDefault);
+
+ org.apache.iceberg.Schema nestedStructSchema =
+ new org.apache.iceberg.Schema(
+ NestedField.required(1, "id", IntegerType.get()),
+ NestedField.optional(
+ 2, "nested", StructType.of(NestedField.optional(3, "a",
IntegerType.get()))));
+
+ Table table = mock(Table.class);
+ when(table.schema()).thenReturn(nestedStructSchema);
+ RecordConverter converter = new RecordConverter(table, config);
+
+ SchemaBuilder connectNestedSchemaBuilder =
+ SchemaBuilder.struct().optional().field("a",
Schema.OPTIONAL_INT32_SCHEMA);
+ Struct nestedDefault = new Struct(connectNestedSchemaBuilder).put("a", 42);
Review Comment:
`nestedDefault` is built from `connectNestedSchemaBuilder` and then set as
that same builder's `defaultValue` before `build()`, so
`nestedDefault.schema().defaultValue()` ends up pointing back at
`nestedDefault` — a cycle. It's harmless today because no exercised path calls
`schema().defaultValue()` on a nested struct schema, but it's fragile for
exactly the test meant to lock this behavior down (same pattern in the
nested-evolution test just below). I'd build the schema once without the
default, create the default struct from that built schema, then rebuild with
the default.
##########
docs/docs/kafka-connect.md:
##########
@@ -94,6 +95,19 @@ If `iceberg.tables.dynamic-enabled` is `false` (the default)
then you must speci
`iceberg.tables.dynamic-enabled` is `true` then you must specify
`iceberg.tables.route-field` which will
contain the name of the table.
+When `iceberg.tables.replace-null-with-default` is set to `false`, an explicit
null value in a
+record is written to the table as null instead of being replaced by the record
schema default
+value. This applies to value conversion and to route-field extraction: a
record whose route field
+is explicitly null is skipped, like any other record with a null route value.
A required Iceberg column cannot take an explicit null under this setting and
the write
+fails, so set `iceberg.tables.schema-force-optional` to `true` (or alter the
columns to optional)
Review Comment:
Since Debezium CDC is the main audience, I'd add a line on the id-columns
interaction: with `false`, an explicitly-null equality-key column flows through
as a null key into the equality-delete write, which silently misses the target
rows (or throws, depending on the writer). `schema-force-optional` doesn't
cover this — it operates on the Connect struct, not the converted Iceberg
record the delete is built from — so the guidance is to constrain equality-key
fields at the pipeline level. Worth calling out so CDC users don't get
surprised by it.
##########
kafka-connect/kafka-connect-transforms/src/test/java/org/apache/iceberg/connect/transforms/TestCopyValue.java:
##########
@@ -86,4 +86,31 @@ public void testCopyValueWithSchema() {
assertThat(newValue.get("data_copy")).isEqualTo("foobar");
}
}
+
+ @Test
+ public void testCopyValueWithSchemaPreservesNullOfFieldWithDefault() {
+ Map<String, String> props =
+ ImmutableMap.of(
+ "source.field", "data",
+ "target.field", "data_copy");
+
+ Schema schema =
+ SchemaBuilder.struct()
+ .field("id", Schema.INT64_SCHEMA)
+ .field("data",
SchemaBuilder.string().optional().defaultValue("").build());
+
+ Struct value = new Struct(schema).put("id", 123L).put("data", null);
+
+ try (CopyValue<SinkRecord> smt = new CopyValue<>()) {
+ smt.configure(props);
+ SinkRecord record = new SinkRecord("topic", 0, null, null, schema,
value, 0);
+ SinkRecord result = smt.apply(record);
+
+ Struct newValue = (Struct) result.value();
+
+ // the explicit null is preserved for both the copied and the target
field
+ assertThat(newValue.getWithoutDefault("data")).isNull();
+ assertThat(newValue.getWithoutDefault("data_copy")).isNull();
Review Comment:
The other three SMT tests assert both that `getWithoutDefault(...)` is null
and that `get(...)` still returns the default — this one only checks the first
half. Adding `assertThat(newValue.get("data")).isEqualTo("");` (and the same
for `data_copy`) would catch a regression that strips the schema default, which
is the property the copied schema is supposed to preserve.
##########
kafka-connect/kafka-connect-transforms/src/main/java/org/apache/iceberg/connect/transforms/CopyValue.java:
##########
@@ -90,9 +90,10 @@ private R applyWithSchema(R record) {
Struct updatedValue = new Struct(updatedSchema);
for (Field field : value.schema().fields()) {
- updatedValue.put(field.name(), value.get(field));
+ // getWithoutDefault so an explicit null is not replaced by the schema
default value
+ updatedValue.put(field.name(), value.getWithoutDefault(field.name()));
Review Comment:
`Struct.getWithoutDefault` only landed in Kafka 3.5 (KAFKA-8713), and the
four bundled SMTs (here plus `DebeziumTransform`, `KafkaMetadataTransform`,
`MongoDebeziumTransform`) now call it unconditionally for every copied field,
independent of the sink flag. Anyone running these on a 2.5–3.4 worker hits
`NoSuchMethodError` the first time a record flows through, which takes the task
down — and our documented minimum is Kafka 2.5, so this is a break for every
SMT user on an older worker, not just the opt-in path.
Can you confirm the version floor? If it holds, I'd bump the documented
minimum to 3.5 (module README + kafka-connect.md) and call it out as the new
floor, or gate the call behind a version check if we want to keep 2.5. The
`false` sink path in `RecordConverter`/`RecordUtils` carries the same 3.5
prerequisite, so whichever way we go the doc note should cover both. This is
the one item I'd want resolved before merge — everything else here is polish.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]