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 636b0ba091f1..1a29da010460 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.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 TAGS_FIELD_VALUE = List.of("a", "b"); + + private static final Map ATTRIBUTES_FIELD_VALUE = Map.of("k", "v"); + private ParquetIcebergWriter parquetIcebergWriter; private TestRunner runner; @@ -194,6 +214,43 @@ void testWriteDataFilesPartitionedTimestamp() throws IOException { 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 eff7652be1ae..565277cebf77 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 DelegatedRecord( 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 b84159d4af0a..27dad53acb45 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 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 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 values = inputRecord.toMap(); - convertedRecord = getConvertedRecord(recordSchema, values); + final Map convertedValues = new LinkedHashMap<>(values.size()); + for (final Map.Entry 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 @@ static Record getConvertedRecord(final Record inputRecord) { return convertedRecord; } - private static Record getConvertedRecord(final RecordSchema recordSchema, final Map values) { - final Map 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 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 fields = recordSchema.getFields(); + private static List convertList(final Collection collection, final Type elementType) { + final List 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 convertMap(final Map map, final Types.MapType mapType) { + // Using LinkedHashMap here to keep input ordering for deterministic flows. + final Map 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 000000000000..3dcc0fcbf2f2 --- /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 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 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 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 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 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 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 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 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 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 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 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)); + } +}