NIFI-16069 - PutIcebergRecord fails with ClassCastException when writing complex types (arrays, maps, nested records) - #11391
NIFI-16069 - PutIcebergRecord fails with ClassCastException when writing complex types (arrays, maps, nested records)#11391maltesander wants to merge 9 commits into
Conversation
exceptionfactory
left a comment
There was a problem hiding this comment.
Thanks for proposing this improvement @maltesander. The initial version of Iceberg integration did not include supported for nested and complex types, so this is an important area of improvement. The general approach looks good, and the tests are helpful. I plan on taking a closer look at the conversion details, I noted a few minor recommendations for now.
|
I pushed some improvements @exceptionfactory:
And one more thing / question: Iceberg supports TimestampTz which we cannot differ via Without something like (timezones vary, needs to be decided): private static Object convertTimestamp(final Timestamp timestamp, final Type icebergType) {
return shouldAdjustToUtc(icebergType) ? timestamp.toInstant().atOffset(ZoneOffset.UTC) : timestamp.toLocalDateTime();
}
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;
};
}The first 3 tests fail without explicit conversion, so we are still missing (at least) one more fix: @Test
void testConvertTimestampWithoutZone() {
final Timestamp timestamp = Timestamp.valueOf(CREATED_LOCAL_DATE_TIME);
final Object converted = RecordConverter.convertValue(timestamp, Types.TimestampType.withoutZone());
assertEquals(CREATED_LOCAL_DATE_TIME, converted);
}
/**
* 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.
*/
@Test
void testConvertTimestampWithZone() {
final Timestamp timestamp = Timestamp.valueOf(CREATED_LOCAL_DATE_TIME);
final Object converted = RecordConverter.convertValue(timestamp, Types.TimestampType.withZone());
final OffsetDateTime offsetDateTime = assertInstanceOf(OffsetDateTime.class, converted);
assertEquals(ZoneOffset.UTC, offsetDateTime.getOffset());
assertEquals(timestamp.toInstant(), offsetDateTime.toInstant());
}
@Test
void testConvertTimestampNanoWithZone() {
final Timestamp timestamp = Timestamp.valueOf(CREATED_LOCAL_DATE_TIME);
final Object converted = RecordConverter.convertValue(timestamp, Types.TimestampNanoType.withZone());
final OffsetDateTime offsetDateTime = assertInstanceOf(OffsetDateTime.class, converted);
assertEquals(timestamp.toInstant(), offsetDateTime.toInstant());
}
/**
* The Iceberg type is not known for every field, so an unresolved type must retain the LocalDateTime conversion.
*/
@Test
void testConvertTimestampUnknownIcebergType() {
final Timestamp timestamp = Timestamp.valueOf(CREATED_LOCAL_DATE_TIME);
final Object converted = RecordConverter.convertValue(timestamp, null);
assertEquals(CREATED_LOCAL_DATE_TIME, converted);
}
/**
* A timestamptz column nested inside a struct must be converted through the recursive path, matching the
* positional access Iceberg uses when writing.
*/
@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());
}This is not mentioned in the Jira ticket or this PR so i would defer fixing this in this PR? |
exceptionfactory
left a comment
There was a problem hiding this comment.
Thanks for your patience @maltesander. I'm open to addressing the multiple timestamp types in this pull request, but deferring it to a separate issue also works.
I noted a few remaining minor recommendations.
| if (!isConversionRequired(recordSchema)) { | ||
| return inputRecord; |
There was a problem hiding this comment.
I recommend adjusting the approach to have a single return, instead of a short-circuit return
| * Recursively convert array, collection, nested record, and map values against the matching Iceberg type. | ||
| * The value is returned unchanged when the target Iceberg type is unknown or does not describe a complex type | ||
| * matching the value. | ||
| */ |
There was a problem hiding this comment.
When adding a method-level comment, the parameters and return should be included.
| * The value is returned unchanged when the target Iceberg type is unknown or does not describe a complex type | ||
| * matching the value. | ||
| */ | ||
| private static Object convertComplexValue(final Object value, final Type icebergType) { |
There was a problem hiding this comment.
Managing multiple returns can become difficult, I recommend refactoring to a single return
…ing complex types (arrays, maps, nested records)
Collapse the remaining short-circuit returns in RecordConverter for consistency.
…to UTC Iceberg timestamptz and timestamptz_ns columns require an OffsetDateTime, so resolve the conversion from the target Iceberg type instead of always producing a LocalDateTime. A Timestamp identifies an instant, so the adjusted conversion preserves that instant expressed at UTC.
c5ade08 to
270dc35
Compare
No worries. Tried to adapt to the review:
I pushed the timestamp changes here 270dc35 and added it to the PR description (not the jira ticket). |
Summary
NIFI-16069 - PutIcebergRecord fails with ClassCastException when writing complex types (arrays, maps, nested records)
PutIcebergRecord fails to write FlowFiles whose schema contains complex/nested types like Iceberg list, map, or struct columns. RecordConverter only translates top-level scalar values (java.sql timestamp/date/time -> java.time) and passes complex values through unchanged.
As a result, values reach Iceberg's Parquet writer in NiFi's native representation, which is incompatible with what Iceberg expects:
Because conversion is gated on scalar field types only, records consisting solely of complex fields skip conversion entirely.
Edit(follow-up): Timestamp Types
RecordConvertertranslated everyjava.sql.Timestampto aLocalDateTime, but Iceberg types declaring an adjustment to UTC require anOffsetDateTime, so writing to atimestamptzcolumn failed:class java.time.LocalDateTime cannot be cast to class java.time.OffsetDateTimeThe conversion is now resolved from the target Iceberg type rather than from the Record field type, which cannot distinguish the two: types reporting
shouldAdjustToUTC()produce anOffsetDateTime, and all other types keep the existingLocalDateTimeconversion.Tracking
Please complete the following tracking steps prior to pull request creation.
Issue Tracking
Pull Request Tracking
NIFI-00000NIFI-00000VerifiedstatusPull Request Formatting
mainbranchVerification
Please indicate the verification steps performed prior to pull request creation.
Build
./mvnw clean install -P contrib-checkLicensing
LICENSEandNOTICEfilesDocumentation