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]

Reply via email to