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));
+ }
+}