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

exceptionfactory pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git


The following commit(s) were added to refs/heads/main by this push:
     new be626038b3e NIFI-16069 Improved PutIcebergRecord support for complex 
types and Timestamps (#11391)
be626038b3e is described below

commit be626038b3e668c8ba69c2ddce98141e0c2793fa
Author: maltesander <[email protected]>
AuthorDate: Tue Sep 8 22:44:34 2026 +0200

    NIFI-16069 Improved PutIcebergRecord support for complex types and 
Timestamps (#11391)
    
    Signed-off-by: David Handermann <[email protected]>
---
 .../iceberg/parquet/ParquetIcebergWriterTest.java  |  57 ++++
 .../processors/iceberg/record/DelegatedRecord.java |   2 +-
 .../processors/iceberg/record/RecordConverter.java | 146 ++++++++--
 .../iceberg/record/RecordConverterTest.java        | 301 +++++++++++++++++++++
 4 files changed, 479 insertions(+), 27 deletions(-)

diff --git 
a/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-parquet-writer/src/test/java/org/apache/nifi/services/iceberg/parquet/ParquetIcebergWriterTest.java
 
b/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-parquet-writer/src/test/java/org/apache/nifi/services/iceberg/parquet/ParquetIcebergWriterTest.java
index 636b0ba091f..1a29da01046 100644
--- 
a/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-parquet-writer/src/test/java/org/apache/nifi/services/iceberg/parquet/ParquetIcebergWriterTest.java
+++ 
b/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-parquet-writer/src/test/java/org/apache/nifi/services/iceberg/parquet/ParquetIcebergWriterTest.java
@@ -47,6 +47,8 @@ import java.io.IOException;
 import java.time.LocalDate;
 import java.time.LocalDateTime;
 import java.time.LocalTime;
+import java.util.List;
+import java.util.Map;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertInstanceOf;
@@ -76,6 +78,24 @@ class ParquetIcebergWriterTest {
             LocalDate.ofEpochDay(1), LocalTime.ofSecondOfDay(0)
     );
 
+    private static final String ID_FIELD_NAME = "id";
+
+    private static final String TAGS_FIELD_NAME = "tags";
+
+    private static final String ADDRESS_FIELD_NAME = "address";
+
+    private static final String CITY_FIELD_NAME = "city";
+
+    private static final String ATTRIBUTES_FIELD_NAME = "attributes";
+
+    private static final String ID_FIELD_VALUE = "row-1";
+
+    private static final String CITY_FIELD_VALUE = "Berlin";
+
+    private static final List<String> TAGS_FIELD_VALUE = List.of("a", "b");
+
+    private static final Map<String, String> ATTRIBUTES_FIELD_VALUE = 
Map.of("k", "v");
+
     private ParquetIcebergWriter parquetIcebergWriter;
 
     private TestRunner runner;
@@ -194,6 +214,43 @@ class ParquetIcebergWriterTest {
         assertEquals(microsecondsExpected, partitionField);
     }
 
+    @Test
+    void testWriteDataFilesComplexTypes() throws IOException {
+        runner.enableControllerService(parquetIcebergWriter);
+
+        final Types.StructType nestedStruct = Types.StructType.of(
+                Types.NestedField.optional(10, CITY_FIELD_NAME, 
Types.StringType.get())
+        );
+        final Schema schema = new Schema(
+                Types.NestedField.required(1, ID_FIELD_NAME, 
Types.StringType.get()),
+                Types.NestedField.optional(2, TAGS_FIELD_NAME,
+                        Types.ListType.ofOptional(3, Types.StringType.get())),
+                Types.NestedField.optional(4, ADDRESS_FIELD_NAME, 
nestedStruct),
+                Types.NestedField.optional(5, ATTRIBUTES_FIELD_NAME,
+                        Types.MapType.ofOptional(6, 7, Types.StringType.get(), 
Types.StringType.get()))
+        );
+        final InMemoryOutputFile outputFile = new InMemoryOutputFile();
+        final PartitionSpec partitionSpec = PartitionSpec.unpartitioned();
+        setTable(schema, partitionSpec, outputFile);
+        
when(locationProvider.newDataLocation(anyString())).thenReturn(LOCATION);
+
+        final IcebergRowWriter rowWriter = 
parquetIcebergWriter.getRowWriter(table);
+
+        final GenericRecord address = GenericRecord.create(nestedStruct);
+        address.setField(CITY_FIELD_NAME, CITY_FIELD_VALUE);
+
+        final GenericRecord row = GenericRecord.create(schema);
+        row.setField(ID_FIELD_NAME, ID_FIELD_VALUE);
+        row.setField(TAGS_FIELD_NAME, TAGS_FIELD_VALUE);
+        row.setField(ADDRESS_FIELD_NAME, address);
+        row.setField(ATTRIBUTES_FIELD_NAME, ATTRIBUTES_FIELD_VALUE);
+        rowWriter.write(row);
+
+        final DataFile[] dataFiles = rowWriter.dataFiles();
+        final byte[] serialized = outputFile.toByteArray();
+        assertDataFilesFound(dataFiles, serialized);
+    }
+
     private void writeRow(final Schema schema, final IcebergRowWriter 
rowWriter) throws IOException {
         final GenericRecord row = GenericRecord.create(schema);
         row.setField(FIRST_FIELD_NAME, FIRST_FIELD_VALUE);
diff --git 
a/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/main/java/org/apache/nifi/processors/iceberg/record/DelegatedRecord.java
 
b/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/main/java/org/apache/nifi/processors/iceberg/record/DelegatedRecord.java
index eff7652be1a..565277cebf7 100644
--- 
a/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/main/java/org/apache/nifi/processors/iceberg/record/DelegatedRecord.java
+++ 
b/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/main/java/org/apache/nifi/processors/iceberg/record/DelegatedRecord.java
@@ -37,8 +37,8 @@ public class DelegatedRecord implements Record {
             final org.apache.nifi.serialization.record.Record record,
             final Types.StructType struct
     ) {
-        this.record = 
RecordConverter.getConvertedRecord(Objects.requireNonNull(record));
         this.struct = Objects.requireNonNull(struct);
+        this.record = 
RecordConverter.getConvertedRecord(Objects.requireNonNull(record), struct);
     }
 
     @Override
diff --git 
a/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/main/java/org/apache/nifi/processors/iceberg/record/RecordConverter.java
 
b/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/main/java/org/apache/nifi/processors/iceberg/record/RecordConverter.java
index b84159d4af0..27dad53acb4 100644
--- 
a/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/main/java/org/apache/nifi/processors/iceberg/record/RecordConverter.java
+++ 
b/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/main/java/org/apache/nifi/processors/iceberg/record/RecordConverter.java
@@ -16,6 +16,8 @@
  */
 package org.apache.nifi.processors.iceberg.record;
 
+import org.apache.iceberg.types.Type;
+import org.apache.iceberg.types.Types;
 import org.apache.nifi.serialization.record.DataType;
 import org.apache.nifi.serialization.record.MapRecord;
 import org.apache.nifi.serialization.record.Record;
@@ -26,6 +28,10 @@ import org.apache.nifi.serialization.record.RecordSchema;
 import java.sql.Date;
 import java.sql.Time;
 import java.sql.Timestamp;
+import java.time.ZoneOffset;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collection;
 import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
@@ -39,22 +45,34 @@ class RecordConverter {
     private static final Set<RecordFieldType> CONVERSION_REQUIRED_FIELD_TYPES 
= Set.of(
             RecordFieldType.TIMESTAMP,
             RecordFieldType.DATE,
-            RecordFieldType.TIME
+            RecordFieldType.TIME,
+            RecordFieldType.ARRAY,
+            RecordFieldType.RECORD,
+            RecordFieldType.MAP,
+            // CHOICE can wrap any of the above, so it must also trigger 
conversion.
+            RecordFieldType.CHOICE
     );
 
     /**
-     * Get Converted Record with conditional handling for field values 
requiring translation
+     * Get Converted Record with recursive, schema-aware handling for field 
values requiring translation
      *
      * @param inputRecord Input Record to be converted
+     * @param struct Iceberg Struct Type describing the target field types 
(may be null for scalar-only conversion)
      * @return Input Record or new Record with converted field values
      */
-    static Record getConvertedRecord(final Record inputRecord) {
+    static Record getConvertedRecord(final Record inputRecord, final 
Types.StructType struct) {
         final Record convertedRecord;
 
         final RecordSchema recordSchema = inputRecord.getSchema();
         if (isConversionRequired(recordSchema)) {
             final Map<String, Object> values = inputRecord.toMap();
-            convertedRecord = getConvertedRecord(recordSchema, values);
+            final Map<String, Object> convertedValues = new 
LinkedHashMap<>(values.size());
+            for (final Map.Entry<String, Object> entry : values.entrySet()) {
+                final String field = entry.getKey();
+                final Type fieldType = fieldType(struct, field);
+                convertedValues.put(field, convertValue(entry.getValue(), 
fieldType));
+            }
+            convertedRecord = new MapRecord(recordSchema, convertedValues);
         } else {
             convertedRecord = inputRecord;
         }
@@ -62,40 +80,116 @@ class RecordConverter {
         return convertedRecord;
     }
 
-    private static Record getConvertedRecord(final RecordSchema recordSchema, 
final Map<String, Object> values) {
-        final Map<String, Object> convertedValues = new LinkedHashMap<>();
+    static Object convertValue(final Object value, final Type icebergType) {
+        return switch (value) {
+            // Convert java.sql types to corresponding java.time types for 
Apache Iceberg
+            case Timestamp timestamp -> convertTimestamp(timestamp, 
icebergType);
+            case Date date -> date.toLocalDate();
+            case Time time -> time.toLocalTime();
+            // Recursively convert complex types against the matching Iceberg 
type
+            case null, default -> convertComplexValue(value, icebergType);
+        };
+    }
 
-        for (final Map.Entry<String, Object> entry : values.entrySet()) {
-            final String field = entry.getKey();
-            final Object value = entry.getValue();
-            final Object converted = getConvertedValue(value);
-            convertedValues.put(field, converted);
+    /**
+     * Convert a Timestamp to the java.time type required by the target 
Iceberg Type. Iceberg Types declaring an
+     * adjustment to UTC require an OffsetDateTime, and other Types require a 
LocalDateTime. A Timestamp identifies
+     * an instant, so the adjusted conversion preserves that instant expressed 
at UTC
+     *
+     * @param timestamp Timestamp to be converted
+     * @param icebergType Iceberg Type describing the target field type (may 
be null when not resolved)
+     * @return OffsetDateTime at UTC for Iceberg Types adjusted to UTC or 
LocalDateTime for other Types
+     */
+    private static Object convertTimestamp(final Timestamp timestamp, final 
Type icebergType) {
+        return shouldAdjustToUtc(icebergType) ? 
timestamp.toInstant().atOffset(ZoneOffset.UTC) : timestamp.toLocalDateTime();
+    }
+
+    /**
+     * Determine whether the Iceberg Type declares an adjustment to UTC, which 
Apache Iceberg requires for the
+     * timestamptz and timestamptz_ns column types
+     *
+     * @param icebergType Iceberg Type describing the target field type (may 
be null when not resolved)
+     * @return Adjustment to UTC required status
+     */
+    private static boolean shouldAdjustToUtc(final Type icebergType) {
+        return switch (icebergType) {
+            case Types.TimestampType timestampType -> 
timestampType.shouldAdjustToUTC();
+            case Types.TimestampNanoType timestampNanoType -> 
timestampNanoType.shouldAdjustToUTC();
+            case null, default -> false;
+        };
+    }
+
+    /**
+     * Recursively convert array, collection, nested record, and map values 
against the matching Iceberg type
+     *
+     * @param value Field value to be converted
+     * @param icebergType Iceberg Type describing the target field type (may 
be null when not resolved)
+     * @return Converted value or the input value when the Iceberg Type is 
unknown or does not describe a complex
+     * type matching the value
+     */
+    private static Object convertComplexValue(final Object value, final Type 
icebergType) {
+        final Object convertedValue;
+
+        if (icebergType == null) {
+            convertedValue = value;
+        } else if (icebergType.isListType()) {
+            convertedValue = convertListValue(value, icebergType.asListType());
+        } else if (icebergType.isStructType() && value instanceof Record 
nestedRecord) {
+            convertedValue = new DelegatedRecord(nestedRecord, 
icebergType.asStructType());
+        } else if (icebergType.isMapType() && value instanceof Map<?, ?> map) {
+            convertedValue = convertMap(map, icebergType.asMapType());
+        } else {
+            convertedValue = value;
         }
 
-        return new MapRecord(recordSchema, convertedValues);
+        return convertedValue;
     }
 
-    private static Object getConvertedValue(final Object value) {
+    /**
+     * Convert an array or collection value to the List required for Apache 
Iceberg with elements converted against
+     * the Iceberg element type
+     *
+     * @param value Field value to be converted
+     * @param listType Iceberg List Type describing the target element type
+     * @return Converted List or the input value when the value is neither an 
array nor a collection
+     */
+    private static Object convertListValue(final Object value, final 
Types.ListType listType) {
+        final Type elementType = listType.elementType();
         return switch (value) {
-            // Convert java.sql types to corresponding java.time types for 
Apache Iceberg
-            case Timestamp timestamp -> timestamp.toLocalDateTime();
-            case Date date -> date.toLocalDate();
-            case Time time -> time.toLocalTime();
+            case Object[] array -> convertList(Arrays.asList(array), 
elementType);
+            case Collection<?> collection -> convertList(collection, 
elementType);
             case null, default -> value;
         };
     }
 
-    private static boolean isConversionRequired(final RecordSchema 
recordSchema) {
-        final List<RecordField> fields = recordSchema.getFields();
+    private static List<Object> convertList(final Collection<?> collection, 
final Type elementType) {
+        final List<Object> converted = new ArrayList<>(collection.size());
+        for (final Object element : collection) {
+            converted.add(convertValue(element, elementType));
+        }
+        return converted;
+    }
 
-        for (final RecordField field : fields) {
-            final DataType dataType = field.getDataType();
-            final RecordFieldType recordFieldType = dataType.getFieldType();
-            if (CONVERSION_REQUIRED_FIELD_TYPES.contains(recordFieldType)) {
-                return true;
-            }
+    private static Map<Object, Object> convertMap(final Map<?, ?> map, final 
Types.MapType mapType) {
+        // Using LinkedHashMap here to keep input ordering for deterministic 
flows.
+        final Map<Object, Object> converted = new LinkedHashMap<>(map.size());
+        for (final Map.Entry<?, ?> entry : map.entrySet()) {
+            final Object key = convertValue(entry.getKey(), mapType.keyType());
+            final Object mappedValue = convertValue(entry.getValue(), 
mapType.valueType());
+            converted.put(key, mappedValue);
         }
+        return converted;
+    }
+
+    private static Type fieldType(final Types.StructType struct, final String 
fieldName) {
+        final Types.NestedField nestedField = struct == null ? null : 
struct.field(fieldName);
+        return nestedField == null ? null : nestedField.type();
+    }
 
-        return false;
+    private static boolean isConversionRequired(final RecordSchema 
recordSchema) {
+        return recordSchema.getFields().stream()
+                .map(RecordField::getDataType)
+                .map(DataType::getFieldType)
+                .anyMatch(CONVERSION_REQUIRED_FIELD_TYPES::contains);
     }
 }
diff --git 
a/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/test/java/org/apache/nifi/processors/iceberg/record/RecordConverterTest.java
 
b/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/test/java/org/apache/nifi/processors/iceberg/record/RecordConverterTest.java
new file mode 100644
index 00000000000..3dcc0fcbf2f
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-iceberg-bundle/nifi-iceberg-processors/src/test/java/org/apache/nifi/processors/iceberg/record/RecordConverterTest.java
@@ -0,0 +1,301 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.processors.iceberg.record;
+
+import org.apache.iceberg.StructLike;
+import org.apache.iceberg.types.Type;
+import org.apache.iceberg.types.Types;
+import org.apache.nifi.serialization.SimpleRecordSchema;
+import org.apache.nifi.serialization.record.MapRecord;
+import org.apache.nifi.serialization.record.Record;
+import org.apache.nifi.serialization.record.RecordField;
+import org.apache.nifi.serialization.record.RecordFieldType;
+import org.apache.nifi.serialization.record.RecordSchema;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.sql.Timestamp;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.OffsetDateTime;
+import java.time.ZoneOffset;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Stream;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+
+class RecordConverterTest {
+
+    private static final String CITY_FIELD_NAME = "city";
+
+    private static final String CREATED_FIELD_NAME = "created";
+
+    private static final String NAME_FIELD_NAME = "name";
+
+    private static final String ITEMS_FIELD_NAME = "items";
+
+    private static final String ADDRESS_FIELD_NAME = "address";
+
+    private static final String ID_FIELD_NAME = "id";
+
+    private static final String CITY_FIELD_VALUE = "Berlin";
+
+    private static final String NAME_FIELD_VALUE = "widget";
+
+    private static final String ID_FIELD_VALUE = "row-1";
+
+    private static final LocalDateTime CREATED_LOCAL_DATE_TIME = 
LocalDateTime.of(2026, 1, 1, 12, 30, 45);
+
+    @Test
+    void testConvertPrimitiveArrayToList() {
+        final Types.ListType listType = Types.ListType.ofOptional(1, 
Types.StringType.get());
+        final Object[] array = new Object[] {"a", "b", "c"};
+
+        final Object converted = RecordConverter.convertValue(array, listType);
+
+        final List<?> list = assertInstanceOf(List.class, converted);
+        assertEquals(List.of("a", "b", "c"), list);
+    }
+
+    @Test
+    void testConvertArrayElementDateTime() {
+        final Types.ListType listType = Types.ListType.ofOptional(1, 
Types.DateType.get());
+        final Object[] array = new Object[] 
{java.sql.Date.valueOf("2026-02-03")};
+
+        final Object converted = RecordConverter.convertValue(array, listType);
+
+        final List<?> list = assertInstanceOf(List.class, converted);
+        assertEquals(List.of(LocalDate.of(2026, 2, 3)), list);
+    }
+
+    @Test
+    void testConvertNestedRecordToStructLike() {
+        final Types.StructType structType = Types.StructType.of(
+                Types.NestedField.optional(1, CITY_FIELD_NAME, 
Types.StringType.get())
+        );
+
+        final RecordSchema nestedSchema = new SimpleRecordSchema(List.of(
+                new RecordField(CITY_FIELD_NAME, 
RecordFieldType.STRING.getDataType())
+        ));
+        final Map<String, Object> nestedValues = new LinkedHashMap<>();
+        nestedValues.put(CITY_FIELD_NAME, CITY_FIELD_VALUE);
+        final Record nestedRecord = new MapRecord(nestedSchema, nestedValues);
+
+        final Object converted = RecordConverter.convertValue(nestedRecord, 
structType);
+
+        final StructLike struct = assertInstanceOf(StructLike.class, 
converted);
+        assertEquals(CITY_FIELD_VALUE, struct.get(0, String.class));
+    }
+
+    @Test
+    void testConvertNestedRecordDateTimeField() {
+        final Types.StructType structType = Types.StructType.of(
+                Types.NestedField.optional(1, CREATED_FIELD_NAME, 
Types.TimestampType.withoutZone())
+        );
+
+        final RecordSchema nestedSchema = new SimpleRecordSchema(List.of(
+                new RecordField(CREATED_FIELD_NAME, 
RecordFieldType.TIMESTAMP.getDataType())
+        ));
+        final Map<String, Object> nestedValues = new LinkedHashMap<>();
+        nestedValues.put(CREATED_FIELD_NAME, Timestamp.valueOf("2026-01-01 
12:30:45"));
+        final Record nestedRecord = new MapRecord(nestedSchema, nestedValues);
+
+        final Object converted = RecordConverter.convertValue(nestedRecord, 
structType);
+
+        final StructLike struct = assertInstanceOf(StructLike.class, 
converted);
+        assertEquals(LocalDateTime.of(2026, 1, 1, 12, 30, 45), struct.get(0, 
LocalDateTime.class));
+    }
+
+    /**
+     * Iceberg Types not adjusted to UTC require a LocalDateTime. The Iceberg 
Type is not resolved for every field,
+     * so an unknown Type must retain the same conversion.
+     */
+    @ParameterizedTest
+    @MethodSource
+    void testConvertTimestampNotAdjustedToUtc(final Type icebergType) {
+        final Timestamp timestamp = Timestamp.valueOf(CREATED_LOCAL_DATE_TIME);
+
+        final Object converted = RecordConverter.convertValue(timestamp, 
icebergType);
+
+        assertEquals(CREATED_LOCAL_DATE_TIME, converted);
+    }
+
+    private static Stream<Arguments> testConvertTimestampNotAdjustedToUtc() {
+        return Stream.of(
+                Arguments.of(Types.TimestampType.withoutZone()),
+                Arguments.of(Types.TimestampNanoType.withoutZone()),
+                Arguments.of((Type) null)
+        );
+    }
+
+    /**
+     * Iceberg timestamptz columns require an OffsetDateTime rather than a 
LocalDateTime. A Timestamp identifies an
+     * instant, so the converted value must describe that same instant 
expressed at UTC.
+     */
+    @ParameterizedTest
+    @MethodSource
+    void testConvertTimestampAdjustedToUtc(final Type icebergType) {
+        final Timestamp timestamp = Timestamp.valueOf(CREATED_LOCAL_DATE_TIME);
+
+        final Object converted = RecordConverter.convertValue(timestamp, 
icebergType);
+
+        final OffsetDateTime offsetDateTime = 
assertInstanceOf(OffsetDateTime.class, converted);
+        assertEquals(ZoneOffset.UTC, offsetDateTime.getOffset());
+        assertEquals(timestamp.toInstant(), offsetDateTime.toInstant());
+    }
+
+    private static Stream<Arguments> testConvertTimestampAdjustedToUtc() {
+        return Stream.of(
+                Arguments.of(Types.TimestampType.withZone()),
+                Arguments.of(Types.TimestampNanoType.withZone())
+        );
+    }
+
+    /**
+     * A timestamptz column nested inside a struct must be converted through 
the recursive path, which requires the
+     * Iceberg Type of the nested field to be resolved and passed down.
+     */
+    @Test
+    void testGetConvertedRecordNestedTimestampWithZone() {
+        final Types.StructType structType = Types.StructType.of(
+                Types.NestedField.optional(1, CREATED_FIELD_NAME, 
Types.TimestampType.withZone())
+        );
+
+        final RecordSchema nestedSchema = new SimpleRecordSchema(List.of(
+                new RecordField(CREATED_FIELD_NAME, 
RecordFieldType.TIMESTAMP.getDataType())
+        ));
+        final Timestamp timestamp = Timestamp.valueOf(CREATED_LOCAL_DATE_TIME);
+        final Map<String, Object> nestedValues = new LinkedHashMap<>();
+        nestedValues.put(CREATED_FIELD_NAME, timestamp);
+        final Record nestedRecord = new MapRecord(nestedSchema, nestedValues);
+
+        final Object converted = RecordConverter.convertValue(nestedRecord, 
structType);
+
+        final StructLike struct = assertInstanceOf(StructLike.class, 
converted);
+        assertEquals(timestamp.toInstant(), struct.get(0, 
OffsetDateTime.class).toInstant());
+    }
+
+    @Test
+    void testConvertMapValues() {
+        final Types.MapType mapType = Types.MapType.ofOptional(
+                1, 2, Types.StringType.get(), Types.StringType.get()
+        );
+        final Map<String, Object> map = new LinkedHashMap<>();
+        map.put(CITY_FIELD_NAME, CITY_FIELD_VALUE);
+        map.put(NAME_FIELD_NAME, NAME_FIELD_VALUE);
+
+        final Object converted = RecordConverter.convertValue(map, mapType);
+
+        final Map<?, ?> resultMap = assertInstanceOf(Map.class, converted);
+        assertEquals(CITY_FIELD_VALUE, resultMap.get(CITY_FIELD_NAME));
+        assertEquals(NAME_FIELD_VALUE, resultMap.get(NAME_FIELD_NAME));
+    }
+
+    @Test
+    void testConvertMapDateTimeValue() {
+        final Types.MapType mapType = Types.MapType.ofOptional(
+                1, 2, Types.StringType.get(), Types.DateType.get()
+        );
+        final Map<String, Object> map = new LinkedHashMap<>();
+        map.put(CREATED_FIELD_NAME, java.sql.Date.valueOf("2026-02-03"));
+
+        final Object converted = RecordConverter.convertValue(map, mapType);
+
+        final Map<?, ?> resultMap = assertInstanceOf(Map.class, converted);
+        assertEquals(LocalDate.of(2026, 2, 3), 
resultMap.get(CREATED_FIELD_NAME));
+    }
+
+    @Test
+    void testGetConvertedRecordArrayOfStructs() {
+        final Types.StructType elementStruct = Types.StructType.of(
+                Types.NestedField.optional(2, NAME_FIELD_NAME, 
Types.StringType.get())
+        );
+        final Types.StructType struct = Types.StructType.of(
+                Types.NestedField.optional(1, ITEMS_FIELD_NAME,
+                        Types.ListType.ofOptional(3, elementStruct))
+        );
+
+        final RecordSchema elementSchema = new SimpleRecordSchema(List.of(
+                new RecordField(NAME_FIELD_NAME, 
RecordFieldType.STRING.getDataType())
+        ));
+        final Map<String, Object> elementValues = new LinkedHashMap<>();
+        elementValues.put(NAME_FIELD_NAME, NAME_FIELD_VALUE);
+        final Record element = new MapRecord(elementSchema, elementValues);
+
+        final RecordSchema schema = new SimpleRecordSchema(List.of(
+                new RecordField(ITEMS_FIELD_NAME,
+                        
RecordFieldType.ARRAY.getArrayDataType(RecordFieldType.RECORD.getRecordDataType(elementSchema)))
+        ));
+        final Map<String, Object> values = new LinkedHashMap<>();
+        values.put(ITEMS_FIELD_NAME, new Object[] {element});
+        final Record record = new MapRecord(schema, values);
+
+        final org.apache.iceberg.data.Record converted = new 
DelegatedRecord(record, struct);
+        final Object items = converted.getField(ITEMS_FIELD_NAME);
+
+        final List<?> list = assertInstanceOf(List.class, items);
+        final StructLike first = assertInstanceOf(StructLike.class, 
list.get(0));
+        assertEquals(NAME_FIELD_VALUE, first.get(0, String.class));
+    }
+
+    /**
+     * A Record field declared as CHOICE, as schema inference produces when a 
field is an object in some Records and a
+     * scalar in others, must still be converted. Conversion is driven by the 
Iceberg type rather than the Record field
+     * type, so the CHOICE needs no dedicated handling, but it must not short 
circuit conversion of the whole Record.
+     */
+    @Test
+    void testGetConvertedRecordChoiceFieldWithScalarSiblings() {
+        final Types.StructType nestedStruct = Types.StructType.of(
+                Types.NestedField.optional(2, CITY_FIELD_NAME, 
Types.StringType.get())
+        );
+        final Types.StructType struct = Types.StructType.of(
+                Types.NestedField.optional(1, ADDRESS_FIELD_NAME, 
nestedStruct),
+                Types.NestedField.optional(3, ID_FIELD_NAME, 
Types.StringType.get())
+        );
+
+        final RecordSchema nestedSchema = new SimpleRecordSchema(List.of(
+                new RecordField(CITY_FIELD_NAME, 
RecordFieldType.STRING.getDataType())
+        ));
+        final Map<String, Object> nestedValues = new LinkedHashMap<>();
+        nestedValues.put(CITY_FIELD_NAME, CITY_FIELD_VALUE);
+        final Record nestedRecord = new MapRecord(nestedSchema, nestedValues);
+
+        // Every field other than the CHOICE is a scalar, so the CHOICE alone 
must require conversion
+        final RecordSchema schema = new SimpleRecordSchema(List.of(
+                new RecordField(ADDRESS_FIELD_NAME, 
RecordFieldType.CHOICE.getChoiceDataType(
+                        RecordFieldType.RECORD.getRecordDataType(nestedSchema),
+                        RecordFieldType.STRING.getDataType())),
+                new RecordField(ID_FIELD_NAME, 
RecordFieldType.STRING.getDataType())
+        ));
+        final Map<String, Object> values = new LinkedHashMap<>();
+        values.put(ADDRESS_FIELD_NAME, nestedRecord);
+        values.put(ID_FIELD_NAME, ID_FIELD_VALUE);
+        final Record record = new MapRecord(schema, values);
+
+        final org.apache.iceberg.data.Record converted = new 
DelegatedRecord(record, struct);
+        final Object address = converted.getField(ADDRESS_FIELD_NAME);
+
+        final StructLike addressStruct = assertInstanceOf(StructLike.class, 
address);
+        assertEquals(CITY_FIELD_VALUE, addressStruct.get(0, String.class));
+        assertEquals(ID_FIELD_VALUE, converted.getField(ID_FIELD_NAME));
+    }
+}

Reply via email to