This is an automated email from the ASF dual-hosted git repository.

mimaison pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new c5ac0af3497 KAFKA-14228: Support for nested structures: ValueToKey 
(#22235)
c5ac0af3497 is described below

commit c5ac0af34972b7e3a31153aad43ec438a39d1836
Author: Hector Geraldino <[email protected]>
AuthorDate: Mon Jul 6 08:58:28 2026 -0400

    KAFKA-14228: Support for nested structures: ValueToKey (#22235)
    
    Adds support for nested structures to the `ValueToKey` transformer.
    
    [KIP-821: Connect Transforms support for nested
    
structures](https://cwiki.apache.org/confluence/display/KAFKA/KIP-821%3A+Connect+Transforms+support+for+nested+structures)
    
    Signed-off-by: Hector Geraldino <[email protected]>
    Reviewers: Mickael Maison <[email protected]>
---
 .../kafka/connect/transforms/ValueToKey.java       |  36 ++-
 .../kafka/connect/transforms/ValueToKeyTest.java   | 308 ++++++++++++++++++++-
 2 files changed, 317 insertions(+), 27 deletions(-)

diff --git 
a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/ValueToKey.java
 
b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/ValueToKey.java
index 673742a1eea..1c89d14f0d1 100644
--- 
a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/ValueToKey.java
+++ 
b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/ValueToKey.java
@@ -27,8 +27,11 @@ import org.apache.kafka.connect.data.Schema;
 import org.apache.kafka.connect.data.SchemaBuilder;
 import org.apache.kafka.connect.data.Struct;
 import org.apache.kafka.connect.errors.DataException;
+import org.apache.kafka.connect.transforms.field.FieldSyntaxVersion;
+import org.apache.kafka.connect.transforms.field.SingleFieldPath;
 import org.apache.kafka.connect.transforms.util.SimpleConfig;
 
+import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -43,15 +46,17 @@ public class ValueToKey<R extends ConnectRecord<R>> 
implements Transformation<R>
     public static final String FIELDS_CONFIG = "fields";
     public static final String REPLACE_NULL_WITH_DEFAULT_CONFIG = 
"replace.null.with.default";
 
-    public static final ConfigDef CONFIG_DEF = new ConfigDef()
-            .define(FIELDS_CONFIG, ConfigDef.Type.LIST, 
ConfigDef.NO_DEFAULT_VALUE, ConfigDef.ValidList.anyNonDuplicateValues(false, 
false), ConfigDef.Importance.HIGH,
-                    "Field names on the record value to extract as the record 
key.")
-            .define(REPLACE_NULL_WITH_DEFAULT_CONFIG, ConfigDef.Type.BOOLEAN, 
true, ConfigDef.Importance.MEDIUM,
-                    "Whether to replace fields that have a default value and 
that are null to the default value. When set to true, the default value is 
used, otherwise null is used.");
+    public static final ConfigDef CONFIG_DEF = 
FieldSyntaxVersion.appendConfigTo(
+            new ConfigDef()
+                    .define(FIELDS_CONFIG, ConfigDef.Type.LIST, 
ConfigDef.NO_DEFAULT_VALUE, ConfigDef.ValidList.anyNonDuplicateValues(false, 
false), ConfigDef.Importance.HIGH,
+                            "Field names on the record value to extract as the 
record key.")
+                    .define(REPLACE_NULL_WITH_DEFAULT_CONFIG, 
ConfigDef.Type.BOOLEAN, true, ConfigDef.Importance.MEDIUM,
+                            "Whether to replace fields that have a default 
value and that are null to the default value. When set to true, the default 
value is used, otherwise null is used."));
 
     private static final String PURPOSE = "copying fields from value to key";
 
     private List<String> fields;
+    private List<SingleFieldPath> fieldPaths;
     private boolean replaceNullWithDefault;
 
     private Cache<Schema, Schema> valueToKeySchemaCache;
@@ -65,6 +70,11 @@ public class ValueToKey<R extends ConnectRecord<R>> 
implements Transformation<R>
     public void configure(Map<String, ?> configs) {
         final SimpleConfig config = new SimpleConfig(CONFIG_DEF, configs);
         fields = config.getList(FIELDS_CONFIG);
+        final FieldSyntaxVersion fieldSyntaxVersion = 
FieldSyntaxVersion.fromConfig(config);
+        fieldPaths = new ArrayList<>(fields.size());
+        for (String field : fields) {
+            fieldPaths.add(new SingleFieldPath(field, fieldSyntaxVersion));
+        }
         replaceNullWithDefault = 
config.getBoolean(REPLACE_NULL_WITH_DEFAULT_CONFIG);
         valueToKeySchemaCache = new SynchronizedCache<>(new LRUCache<>(16));
     }
@@ -81,8 +91,8 @@ public class ValueToKey<R extends ConnectRecord<R>> 
implements Transformation<R>
     private R applySchemaless(R record) {
         final Map<String, Object> value = requireMap(record.value(), PURPOSE);
         final Map<String, Object> key = new HashMap<>(fields.size());
-        for (String field : fields) {
-            key.put(field, value.get(field));
+        for (int i = 0; i < fields.size(); i++) {
+            key.put(fields.get(i), fieldPaths.get(i).valueFrom(value));
         }
         return record.newRecord(record.topic(), record.kafkaPartition(), null, 
key, record.valueSchema(), record.value(), record.timestamp());
     }
@@ -93,20 +103,20 @@ public class ValueToKey<R extends ConnectRecord<R>> 
implements Transformation<R>
         Schema keySchema = valueToKeySchemaCache.get(value.schema());
         if (keySchema == null) {
             final SchemaBuilder keySchemaBuilder = SchemaBuilder.struct();
-            for (String field : fields) {
-                final Field fieldFromValue = value.schema().field(field);
+            for (int i = 0; i < fields.size(); i++) {
+                final Field fieldFromValue = 
fieldPaths.get(i).fieldFrom(value.schema());
                 if (fieldFromValue == null) {
-                    throw new DataException("Field does not exist: " + field);
+                    throw new DataException("Field does not exist: " + 
fields.get(i));
                 }
-                keySchemaBuilder.field(field, fieldFromValue.schema());
+                keySchemaBuilder.field(fields.get(i), fieldFromValue.schema());
             }
             keySchema = keySchemaBuilder.build();
             valueToKeySchemaCache.put(value.schema(), keySchema);
         }
 
         final Struct key = new Struct(keySchema);
-        for (String field : fields) {
-            key.put(field, replaceNullWithDefault ? value.get(field) : 
value.getWithoutDefault(field));
+        for (int i = 0; i < fields.size(); i++) {
+            key.put(fields.get(i), fieldPaths.get(i).valueFrom(value, 
replaceNullWithDefault));
         }
 
         return record.newRecord(record.topic(), record.kafkaPartition(), 
keySchema, key, value.schema(), value, record.timestamp());
diff --git 
a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/ValueToKeyTest.java
 
b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/ValueToKeyTest.java
index 345d252e758..d1be289d403 100644
--- 
a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/ValueToKeyTest.java
+++ 
b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/ValueToKeyTest.java
@@ -22,6 +22,7 @@ import org.apache.kafka.connect.data.SchemaBuilder;
 import org.apache.kafka.connect.data.Struct;
 import org.apache.kafka.connect.errors.DataException;
 import org.apache.kafka.connect.sink.SinkRecord;
+import org.apache.kafka.connect.transforms.field.FieldSyntaxVersion;
 
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Test;
@@ -48,12 +49,12 @@ public class ValueToKeyTest {
     }
 
     @AfterEach
-    public void teardown() {
+    void teardown() {
         xform.close();
     }
 
     @Test
-    public void schemaless() {
+    void schemaless() {
         xform.configure(Map.of("fields", "a,b"));
 
         final HashMap<String, Integer> value = new HashMap<>();
@@ -61,8 +62,8 @@ public class ValueToKeyTest {
         value.put("b", 2);
         value.put("c", 3);
 
-        final SinkRecord record = new SinkRecord("", 0, null, null, null, 
value, 0);
-        final SinkRecord transformedRecord = xform.apply(record);
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, null, 
value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
 
         final HashMap<String, Integer> expectedKey = new HashMap<>();
         expectedKey.put("a", 1);
@@ -73,7 +74,7 @@ public class ValueToKeyTest {
     }
 
     @Test
-    public void withSchema() {
+    void withSchema() {
         xform.configure(Map.of("fields", "a,b"));
 
         final Schema valueSchema = SchemaBuilder.struct()
@@ -87,8 +88,8 @@ public class ValueToKeyTest {
         value.put("b", 2);
         value.put("c", 3);
 
-        final SinkRecord record = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
-        final SinkRecord transformedRecord = xform.apply(record);
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
 
         final Schema expectedKeySchema = SchemaBuilder.struct()
                 .field("a", Schema.INT32_SCHEMA)
@@ -104,7 +105,7 @@ public class ValueToKeyTest {
     }
 
     @Test
-    public void nonExistingField() {
+    void nonExistingField() {
         xform.configure(Map.of("fields", "not_exist"));
 
         final Schema valueSchema = SchemaBuilder.struct()
@@ -114,20 +115,20 @@ public class ValueToKeyTest {
         final Struct value = new Struct(valueSchema);
         value.put("a", 1);
 
-        final SinkRecord record = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
 
-        DataException actual = assertThrows(DataException.class, () -> 
xform.apply(record));
+        DataException actual = assertThrows(DataException.class, () -> 
xform.apply(sinkRecord));
         assertEquals("Field does not exist: not_exist", actual.getMessage());
     }
 
     @Test
-    public void testValueToKeyVersionRetrievedFromAppInfoParser() {
+    void testValueToKeyVersionRetrievedFromAppInfoParser() {
         assertEquals(AppInfoParser.getVersion(), xform.version());
     }
 
     @ParameterizedTest
     @MethodSource("data")
-    public void testReplaceNullWithDefaultConfig(boolean 
replaceNullWithDefault, Object expectedValue) {
+    void testReplaceNullWithDefaultConfig(boolean replaceNullWithDefault, 
Object expectedValue) {
         Map<String, Object> config = new HashMap<>();
         config.put("fields", "optional_with_default");
         config.put("replace.null.with.default", replaceNullWithDefault);
@@ -138,9 +139,288 @@ public class ValueToKeyTest {
                 .build();
         final Struct value = new 
Struct(valueSchema).put("optional_with_default", null);
 
-        final SinkRecord record = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
-        final SinkRecord transformedRecord = xform.apply(record);
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
 
         assertEquals(expectedValue, ((Struct) 
transformedRecord.key()).getWithoutDefault("optional_with_default"));
     }
+
+    @Test
+    void schemalessNestedFieldV2() {
+        Map<String, String> configs = new HashMap<>();
+        configs.put("fields", "name,address.city");
+        configs.put(FieldSyntaxVersion.FIELD_SYNTAX_VERSION_CONFIG, 
FieldSyntaxVersion.V2.name());
+        xform.configure(configs);
+
+        final Map<String, Object> address = new HashMap<>();
+        address.put("city", "New York");
+        address.put("state", "NY");
+        final Map<String, Object> value = new HashMap<>();
+        value.put("name", "Franz");
+        value.put("address", address);
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, null, 
value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
+
+        assertNull(transformedRecord.keySchema());
+
+        final HashMap<String, String> expectedKey = new HashMap<>();
+        expectedKey.put("name", "Franz");
+        expectedKey.put("address.city", "New York");
+
+        @SuppressWarnings("unchecked")
+        final Map<String, Object> key = (Map<String, Object>) 
transformedRecord.key();
+        assertEquals(expectedKey, key);
+    }
+
+    @Test
+    void withSchemaNestedFieldV2() {
+        Map<String, String> configs = new HashMap<>();
+        configs.put("fields", "name,address.city");
+        configs.put(FieldSyntaxVersion.FIELD_SYNTAX_VERSION_CONFIG, 
FieldSyntaxVersion.V2.name());
+        xform.configure(configs);
+
+        final Schema addressSchema = SchemaBuilder.struct()
+                .field("city", Schema.STRING_SCHEMA)
+                .field("state", Schema.STRING_SCHEMA)
+                .build();
+        final Schema valueSchema = SchemaBuilder.struct()
+                .field("name", Schema.STRING_SCHEMA)
+                .field("address", addressSchema)
+                .build();
+        final Struct value = new Struct(valueSchema)
+                .put("name", "Franz")
+                .put("address", new Struct(addressSchema)
+                        .put("city", "New York")
+                        .put("state", "NY"));
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
+
+        final Schema expectedKeySchema = SchemaBuilder.struct()
+                .field("name", Schema.STRING_SCHEMA)
+                .field("address.city", Schema.STRING_SCHEMA)
+                .build();
+
+        assertEquals(expectedKeySchema, transformedRecord.keySchema());
+        assertEquals("Franz", ((Struct) transformedRecord.key()).get("name"));
+        assertEquals("New York", ((Struct) 
transformedRecord.key()).get("address.city"));
+    }
+
+    @Test
+    void withSchemaMultipleNestedFieldsV2() {
+        Map<String, String> configs = new HashMap<>();
+        configs.put("fields", "address.city,address.state");
+        configs.put(FieldSyntaxVersion.FIELD_SYNTAX_VERSION_CONFIG, 
FieldSyntaxVersion.V2.name());
+        xform.configure(configs);
+
+        final Schema addressSchema = SchemaBuilder.struct()
+                .field("city", Schema.STRING_SCHEMA)
+                .field("state", Schema.STRING_SCHEMA)
+                .build();
+        final Schema valueSchema = SchemaBuilder.struct()
+                .field("name", Schema.STRING_SCHEMA)
+                .field("address", addressSchema)
+                .build();
+        final Struct value = new Struct(valueSchema)
+                .put("name", "Franz")
+                .put("address", new Struct(addressSchema)
+                        .put("city", "New York")
+                        .put("state", "NY"));
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
+
+        final Schema expectedKeySchema = SchemaBuilder.struct()
+                .field("address.city", Schema.STRING_SCHEMA)
+                .field("address.state", Schema.STRING_SCHEMA)
+                .build();
+
+        assertEquals(expectedKeySchema, transformedRecord.keySchema());
+        final Struct key = (Struct) transformedRecord.key();
+        assertEquals("New York", key.get("address.city"));
+        assertEquals("NY", key.get("address.state"));
+    }
+
+    @Test
+    void nonExistentNestedFieldWithSchemaV2() {
+        Map<String, String> configs = new HashMap<>();
+        configs.put("fields", "address.zip");
+        configs.put(FieldSyntaxVersion.FIELD_SYNTAX_VERSION_CONFIG, 
FieldSyntaxVersion.V2.name());
+        xform.configure(configs);
+
+        final Schema addressSchema = SchemaBuilder.struct()
+                .field("city", Schema.STRING_SCHEMA)
+                .build();
+        final Schema valueSchema = SchemaBuilder.struct()
+                .field("address", addressSchema)
+                .build();
+        final Struct value = new Struct(valueSchema)
+                .put("address", new Struct(addressSchema).put("city", "New 
York"));
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
+
+        DataException actual = assertThrows(DataException.class, () -> 
xform.apply(sinkRecord));
+        assertEquals("Field does not exist: address.zip", actual.getMessage());
+    }
+
+    @Test
+    void schemalessNestedFieldNotFoundV2() {
+        Map<String, String> configs = new HashMap<>();
+        configs.put("fields", "address.zip");
+        configs.put(FieldSyntaxVersion.FIELD_SYNTAX_VERSION_CONFIG, 
FieldSyntaxVersion.V2.name());
+        xform.configure(configs);
+
+        final Map<String, Object> value = new HashMap<>();
+        value.put("address", Map.of("city", "New York"));
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, null, 
value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
+
+        assertNull(transformedRecord.keySchema());
+        @SuppressWarnings("unchecked")
+        final Map<String, Object> key = (Map<String, Object>) 
transformedRecord.key();
+        assertNull(key.get("address.zip"));
+    }
+
+    @Test
+    void schemalessNestedStructToKeyV2() {
+        Map<String, String> configs = new HashMap<>();
+        configs.put("fields", "parent.child");
+        configs.put(FieldSyntaxVersion.FIELD_SYNTAX_VERSION_CONFIG, 
FieldSyntaxVersion.V2.name());
+        xform.configure(configs);
+
+        final Map<String, Object> child = new HashMap<>();
+        child.put("k2", "123");
+        final Map<String, Object> parent = new HashMap<>();
+        parent.put("child", child);
+        final Map<String, Object> value = new HashMap<>();
+        value.put("k1", 123);
+        value.put("parent", parent);
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, null, 
value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
+
+        assertNull(transformedRecord.keySchema());
+
+        final HashMap<String, Object> expectedKey = new HashMap<>();
+        expectedKey.put("parent.child", child);
+
+        @SuppressWarnings("unchecked")
+        final Map<String, Object> key = (Map<String, Object>) 
transformedRecord.key();
+        assertEquals(expectedKey, key);
+    }
+
+    @Test
+    void withSchemaNestedStructToKeyV2() {
+        Map<String, String> configs = new HashMap<>();
+        configs.put("fields", "parent.child");
+        configs.put(FieldSyntaxVersion.FIELD_SYNTAX_VERSION_CONFIG, 
FieldSyntaxVersion.V2.name());
+        xform.configure(configs);
+
+        final Schema childSchema = SchemaBuilder.struct()
+                .field("k2", Schema.STRING_SCHEMA)
+                .build();
+        final Schema parentSchema = SchemaBuilder.struct()
+                .field("child", childSchema)
+                .build();
+        final Schema valueSchema = SchemaBuilder.struct()
+                .field("k1", Schema.INT32_SCHEMA)
+                .field("parent", parentSchema)
+                .build();
+        final Struct value = new Struct(valueSchema)
+                .put("k1", 123)
+                .put("parent", new Struct(parentSchema)
+                        .put("child", new Struct(childSchema)
+                                .put("k2", "123")));
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
+
+        final Schema expectedKeySchema = SchemaBuilder.struct()
+                .field("parent.child", childSchema)
+                .build();
+
+        assertEquals(expectedKeySchema, transformedRecord.keySchema());
+        final Struct key = (Struct) transformedRecord.key();
+        final Struct childStruct = (Struct) key.get("parent.child");
+        assertEquals("123", childStruct.get("k2"));
+    }
+
+    @Test
+    void schemalessBacktickedFieldNameV2() {
+        Map<String, String> configs = new HashMap<>();
+        configs.put("fields", "`parent.child`.k2");
+        configs.put(FieldSyntaxVersion.FIELD_SYNTAX_VERSION_CONFIG, 
FieldSyntaxVersion.V2.name());
+        xform.configure(configs);
+
+        final Map<String, Object> parentChild = new HashMap<>();
+        parentChild.put("k2", "123");
+        final Map<String, Object> value = new HashMap<>();
+        value.put("k1", 123);
+        value.put("parent.child", parentChild);
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, null, 
value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
+
+        assertNull(transformedRecord.keySchema());
+
+        final HashMap<String, String> expectedKey = new HashMap<>();
+        expectedKey.put("`parent.child`.k2", "123");
+
+        @SuppressWarnings("unchecked")
+        final Map<String, Object> key = (Map<String, Object>) 
transformedRecord.key();
+        assertEquals(expectedKey, key);
+    }
+
+    @Test
+    void withSchemaBacktickedFieldNameV2() {
+        Map<String, String> configs = new HashMap<>();
+        configs.put("fields", "`parent.child`.k2");
+        configs.put(FieldSyntaxVersion.FIELD_SYNTAX_VERSION_CONFIG, 
FieldSyntaxVersion.V2.name());
+        xform.configure(configs);
+
+        final Schema parentChildSchema = SchemaBuilder.struct()
+                .field("k2", Schema.STRING_SCHEMA)
+                .build();
+        final Schema valueSchema = SchemaBuilder.struct()
+                .field("k1", Schema.INT32_SCHEMA)
+                .field("parent.child", parentChildSchema)
+                .build();
+        final Struct value = new Struct(valueSchema)
+                .put("k1", 123)
+                .put("parent.child", new Struct(parentChildSchema)
+                        .put("k2", "123"));
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
+
+        final Schema expectedKeySchema = SchemaBuilder.struct()
+                .field("`parent.child`.k2", Schema.STRING_SCHEMA)
+                .build();
+
+        assertEquals(expectedKeySchema, transformedRecord.keySchema());
+        assertEquals("123", ((Struct) 
transformedRecord.key()).get("`parent.child`.k2"));
+    }
+
+    @Test
+    void v1BackwardCompatibilityDottedFieldName() {
+        // With V1 (default), a dotted field name is treated as a literal 
field name
+        xform.configure(Map.of("fields", "address.city"));
+
+        final Schema valueSchema = SchemaBuilder.struct()
+                .field("address.city", Schema.STRING_SCHEMA)
+                .build();
+        final Struct value = new Struct(valueSchema).put("address.city", "New 
York");
+
+        final SinkRecord sinkRecord = new SinkRecord("", 0, null, null, 
valueSchema, value, 0);
+        final SinkRecord transformedRecord = xform.apply(sinkRecord);
+
+        final Schema expectedKeySchema = SchemaBuilder.struct()
+                .field("address.city", Schema.STRING_SCHEMA)
+                .build();
+
+        assertEquals(expectedKeySchema, transformedRecord.keySchema());
+        assertEquals("New York", ((Struct) 
transformedRecord.key()).get("address.city"));
+    }
 }

Reply via email to