voonhous commented on code in PR #20049:
URL: https://github.com/apache/hudi/pull/20049#discussion_r4166270379
##########
hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieStreamerUtils.java:
##########
@@ -238,4 +250,145 @@ void testCombinePropertiesWithPropsOverride() {
// propsOverride takes precedence
assertEquals("overrideValue", result.getString("hoodie.overridden.key"));
}
+
+ /**
+ * A null value in the ordering field must be quarantined as a
record-creation failure rather than
+ * flowing into the write.
+ */
+ @Test
+ void testCreateHoodieRecordsWithNullOrderingValue() {
+ HoodieSchema schema = HoodieSchema.parse(NULLABLE_ORDERING_SCHEMA_STRING);
+ JavaRDD<GenericRecord> recordRdd =
jsc.parallelize(Collections.singletonList(1)).map(i -> {
+ GenericRecord genericRecord = new
GenericData.Record(schema.toAvroSchema());
+ genericRecord.put(0, i * 1000L);
+ genericRecord.put(1, "key" + i);
+ genericRecord.put(2, "path" + i);
+ genericRecord.put(3, "rider1");
+ genericRecord.put(4, "driver1");
+ genericRecord.put(5, null);
+ return genericRecord;
+ });
+ HoodieStreamer.Config cfg = new HoodieStreamer.Config();
+ cfg.payloadClassName = DefaultHoodieRecordPayload.class.getName();
+ cfg.sourceOrderingFields = ORDERING_FIELD;
+ TypedProperties props = new TypedProperties();
+ props.put(KeyGeneratorOptions.PARTITIONPATH_FIELD_NAME.key(),
"partition_path");
+ props.put(KeyGeneratorOptions.RECORDKEY_FIELD_NAME.key(), "_row_key");
+ BaseErrorTableWriter errorTableWriter =
Mockito.mock(BaseErrorTableWriter.class);
+ ArgumentCaptor<JavaRDD<?>> errorEventCaptor =
ArgumentCaptor.forClass(JavaRDD.class);
+
doNothing().when(errorTableWriter).addErrorEvents(errorEventCaptor.capture());
+
+ Option<JavaRDD<HoodieRecord>> recordOpt =
HoodieStreamerUtils.createHoodieRecords(
+ cfg, props, Option.of(recordRdd), new SimpleSchemaProvider(jsc,
schema, props),
+ HoodieRecordType.AVRO, false, "000", Option.of(errorTableWriter), new
HoodieTableConfig());
+
+ assertTrue(recordOpt.isPresent());
+ assertEquals(Collections.emptyList(), recordOpt.get().collect());
+ List<ErrorEvent<String>> errorEvents = (List<ErrorEvent<String>>)
errorEventCaptor.getValue().collect();
+ // The record is schema-valid, so it is serialized by the Avro
JsonEncoder, which wraps
+ // nullable union values.
+ ErrorEvent<String> expectedErrorEvent = new ErrorEvent<>(
+
"{\"timestamp\":1000,\"_row_key\":\"key1\",\"partition_path\":{\"string\":\"path1\"},"
+ + "\"rider\":\"rider1\",\"driver\":\"driver1\",\"" +
ORDERING_FIELD + "\":null}",
+ ErrorEvent.ErrorReason.RECORD_CREATION);
+ assertEquals(Collections.singletonList(expectedErrorEvent), errorEvents);
+ }
+
+ /**
+ * With several ordering fields, a null in any one of them is still a
missing ordering value:
+ * OrderingValues.create returns a non-null ArrayComparable holding the
null, so it survives a
+ * plain null check and fails later when the merger compares it.
+ */
+ @Test
+ void testCreateHoodieRecordsWithNullInOneOfSeveralOrderingFields() {
+ HoodieSchema schema = HoodieSchema.parse(NULLABLE_ORDERING_SCHEMA_STRING);
+ JavaRDD<GenericRecord> recordRdd = nullOrderingRecords(schema, false);
+ HoodieStreamer.Config cfg = nullOrderingConfig("timestamp," +
ORDERING_FIELD);
+
+ List<ErrorEvent<String>> errorEvents = quarantinedEvents(cfg, schema,
recordRdd);
+
+ assertEquals(1, errorEvents.size());
+ assertEquals(ErrorEvent.ErrorReason.RECORD_CREATION,
errorEvents.get(0).getReason());
+ }
+
+ /**
+ * A delete carrying a null ordering value is rejected too. It is not a
commit-time-ordering
+ * delete, since that requires the default ordering value rather than a null
one, so the merger
+ * would compare the null and fail.
+ */
+ @Test
+ void testCreateHoodieRecordsWithNullOrderingValueOnDelete() {
+ HoodieSchema schema =
HoodieSchema.parse(NULLABLE_ORDERING_DELETE_SCHEMA_STRING);
+ JavaRDD<GenericRecord> recordRdd = nullOrderingRecords(schema, true);
+ HoodieStreamer.Config cfg = nullOrderingConfig(ORDERING_FIELD);
+
+ List<ErrorEvent<String>> errorEvents = quarantinedEvents(cfg, schema,
recordRdd);
+
+ assertEquals(1, errorEvents.size());
+ assertEquals(ErrorEvent.ErrorReason.RECORD_CREATION,
errorEvents.get(0).getReason());
+ }
+
+ /**
+ * Without an error table there is nowhere to quarantine the record, so the
batch fails instead of
+ * writing an unusable ordering value.
+ */
+ @Test
+ void testCreateHoodieRecordsWithNullOrderingValueFailsWithoutErrorTable() {
+ HoodieSchema schema = HoodieSchema.parse(NULLABLE_ORDERING_SCHEMA_STRING);
+ JavaRDD<GenericRecord> recordRdd = nullOrderingRecords(schema, false);
+ HoodieStreamer.Config cfg = nullOrderingConfig(ORDERING_FIELD);
+
+ Option<JavaRDD<HoodieRecord>> recordOpt =
HoodieStreamerUtils.createHoodieRecords(
+ cfg, nullOrderingProps(), Option.of(recordRdd), new
SimpleSchemaProvider(jsc, schema, nullOrderingProps()),
+ HoodieRecordType.AVRO, false, "000", Option.empty(), new
HoodieTableConfig());
+
+ assertTrue(recordOpt.isPresent());
+ SparkException sparkException = assertThrows(SparkException.class, () ->
recordOpt.get().collect());
+ assertEquals(HoodieRecordCreationException.class,
sparkException.getCause().getClass());
+ }
+
+ private static TypedProperties nullOrderingProps() {
+ TypedProperties props = new TypedProperties();
+ props.put(KeyGeneratorOptions.PARTITIONPATH_FIELD_NAME.key(),
"partition_path");
+ props.put(KeyGeneratorOptions.RECORDKEY_FIELD_NAME.key(), "_row_key");
+ return props;
+ }
+
+ private static HoodieStreamer.Config nullOrderingConfig(String
orderingFields) {
+ HoodieStreamer.Config cfg = new HoodieStreamer.Config();
+ cfg.payloadClassName = DefaultHoodieRecordPayload.class.getName();
+ cfg.sourceOrderingFields = orderingFields;
+ return cfg;
+ }
Review Comment:
All four new tests use this function, so `requiresOrderingValue` is always
true. No test reaches the false gate. This config causes the
`cfg.recordMergeMode` to be pinned to null, which IIUC, `HoodieStreamer` never
produces, and passes an empty `HoodieTableConfig`.
Could we adde a parameterized case for `COMMIT_TIME_ORDERING` and
`OverwriteWithLatestAvroPayload` asserting the null-ordering record is returned
and nothing is quarantined?
##########
hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerUtils.java:
##########
@@ -97,6 +99,9 @@ public static Option<JavaRDD<HoodieRecord>>
createHoodieRecords(HoodieStreamer.C
String payloadClassName = StringUtils.isNullOrEmpty(cfg.payloadClassName)
? HoodieRecordPayload.getAvroPayloadForMergeMode(cfg.recordMergeMode,
cfg.payloadClassName)
: cfg.payloadClassName;
+ boolean requiresOrderingValue = shouldUseOrderingField
+ && cfg.recordMergeMode != RecordMergeMode.COMMIT_TIME_ORDERING
Review Comment:
+1 on this, please chec.
##########
hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerUtils.java:
##########
@@ -130,6 +135,13 @@ public static Option<JavaRDD<HoodieRecord>>
createHoodieRecords(HoodieStreamer.C
? OrderingValues.create(orderingFieldsStr.split(","),
field -> (Comparable)
HoodieAvroUtils.getNestedFieldVal(gr, field, false,
useConsistentLogicalTimestamp))
: null;
+ if (requiresOrderingValue &&
OrderingValues.isMissing(orderingValue)) {
+ throw new IllegalArgumentException(
+ "Ordering fields '" + orderingFieldsStr + "' resolved
to a null value for record key '"
+ + hoodieKey.getRecordKey() + "'. Please ensure all
records carry non-null values for "
+ + "the ordering fields, or use a merge mode or
payload class that does not order "
+ + "(e.g., COMMIT_TIME_ORDERING or
OverwriteWithLatestAvroPayload).");
Review Comment:
IIUC, this impacts MOR deletes. MOR deletes with a null ordering value,
which this PR will quarantine them.
Could deletes skip the check and carry `OrderingValues.getDefault()`, or
should the Summary of the PR be corrected and an impact line for MOR is added?
##########
hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieStreamerUtils.java:
##########
@@ -238,4 +250,145 @@ void testCombinePropertiesWithPropsOverride() {
// propsOverride takes precedence
assertEquals("overrideValue", result.getString("hoodie.overridden.key"));
}
+
+ /**
+ * A null value in the ordering field must be quarantined as a
record-creation failure rather than
+ * flowing into the write.
+ */
+ @Test
+ void testCreateHoodieRecordsWithNullOrderingValue() {
+ HoodieSchema schema = HoodieSchema.parse(NULLABLE_ORDERING_SCHEMA_STRING);
+ JavaRDD<GenericRecord> recordRdd =
jsc.parallelize(Collections.singletonList(1)).map(i -> {
+ GenericRecord genericRecord = new
GenericData.Record(schema.toAvroSchema());
+ genericRecord.put(0, i * 1000L);
+ genericRecord.put(1, "key" + i);
+ genericRecord.put(2, "path" + i);
+ genericRecord.put(3, "rider1");
+ genericRecord.put(4, "driver1");
+ genericRecord.put(5, null);
+ return genericRecord;
+ });
+ HoodieStreamer.Config cfg = new HoodieStreamer.Config();
+ cfg.payloadClassName = DefaultHoodieRecordPayload.class.getName();
+ cfg.sourceOrderingFields = ORDERING_FIELD;
+ TypedProperties props = new TypedProperties();
+ props.put(KeyGeneratorOptions.PARTITIONPATH_FIELD_NAME.key(),
"partition_path");
+ props.put(KeyGeneratorOptions.RECORDKEY_FIELD_NAME.key(), "_row_key");
+ BaseErrorTableWriter errorTableWriter =
Mockito.mock(BaseErrorTableWriter.class);
+ ArgumentCaptor<JavaRDD<?>> errorEventCaptor =
ArgumentCaptor.forClass(JavaRDD.class);
+
doNothing().when(errorTableWriter).addErrorEvents(errorEventCaptor.capture());
+
+ Option<JavaRDD<HoodieRecord>> recordOpt =
HoodieStreamerUtils.createHoodieRecords(
+ cfg, props, Option.of(recordRdd), new SimpleSchemaProvider(jsc,
schema, props),
+ HoodieRecordType.AVRO, false, "000", Option.of(errorTableWriter), new
HoodieTableConfig());
+
+ assertTrue(recordOpt.isPresent());
+ assertEquals(Collections.emptyList(), recordOpt.get().collect());
+ List<ErrorEvent<String>> errorEvents = (List<ErrorEvent<String>>)
errorEventCaptor.getValue().collect();
+ // The record is schema-valid, so it is serialized by the Avro
JsonEncoder, which wraps
+ // nullable union values.
+ ErrorEvent<String> expectedErrorEvent = new ErrorEvent<>(
+
"{\"timestamp\":1000,\"_row_key\":\"key1\",\"partition_path\":{\"string\":\"path1\"},"
+ + "\"rider\":\"rider1\",\"driver\":\"driver1\",\"" +
ORDERING_FIELD + "\":null}",
+ ErrorEvent.ErrorReason.RECORD_CREATION);
+ assertEquals(Collections.singletonList(expectedErrorEvent), errorEvents);
+ }
+
+ /**
+ * With several ordering fields, a null in any one of them is still a
missing ordering value:
+ * OrderingValues.create returns a non-null ArrayComparable holding the
null, so it survives a
+ * plain null check and fails later when the merger compares it.
+ */
+ @Test
+ void testCreateHoodieRecordsWithNullInOneOfSeveralOrderingFields() {
+ HoodieSchema schema = HoodieSchema.parse(NULLABLE_ORDERING_SCHEMA_STRING);
+ JavaRDD<GenericRecord> recordRdd = nullOrderingRecords(schema, false);
+ HoodieStreamer.Config cfg = nullOrderingConfig("timestamp," +
ORDERING_FIELD);
+
+ List<ErrorEvent<String>> errorEvents = quarantinedEvents(cfg, schema,
recordRdd);
+
+ assertEquals(1, errorEvents.size());
+ assertEquals(ErrorEvent.ErrorReason.RECORD_CREATION,
errorEvents.get(0).getReason());
+ }
+
+ /**
+ * A delete carrying a null ordering value is rejected too. It is not a
commit-time-ordering
+ * delete, since that requires the default ordering value rather than a null
one, so the merger
+ * would compare the null and fail.
+ */
+ @Test
+ void testCreateHoodieRecordsWithNullOrderingValueOnDelete() {
+ HoodieSchema schema =
HoodieSchema.parse(NULLABLE_ORDERING_DELETE_SCHEMA_STRING);
+ JavaRDD<GenericRecord> recordRdd = nullOrderingRecords(schema, true);
+ HoodieStreamer.Config cfg = nullOrderingConfig(ORDERING_FIELD);
+
+ List<ErrorEvent<String>> errorEvents = quarantinedEvents(cfg, schema,
recordRdd);
+
+ assertEquals(1, errorEvents.size());
+ assertEquals(ErrorEvent.ErrorReason.RECORD_CREATION,
errorEvents.get(0).getReason());
+ }
+
+ /**
+ * Without an error table there is nowhere to quarantine the record, so the
batch fails instead of
+ * writing an unusable ordering value.
+ */
+ @Test
+ void testCreateHoodieRecordsWithNullOrderingValueFailsWithoutErrorTable() {
+ HoodieSchema schema = HoodieSchema.parse(NULLABLE_ORDERING_SCHEMA_STRING);
+ JavaRDD<GenericRecord> recordRdd = nullOrderingRecords(schema, false);
+ HoodieStreamer.Config cfg = nullOrderingConfig(ORDERING_FIELD);
+
+ Option<JavaRDD<HoodieRecord>> recordOpt =
HoodieStreamerUtils.createHoodieRecords(
+ cfg, nullOrderingProps(), Option.of(recordRdd), new
SimpleSchemaProvider(jsc, schema, nullOrderingProps()),
+ HoodieRecordType.AVRO, false, "000", Option.empty(), new
HoodieTableConfig());
+
+ assertTrue(recordOpt.isPresent());
+ SparkException sparkException = assertThrows(SparkException.class, () ->
recordOpt.get().collect());
+ assertEquals(HoodieRecordCreationException.class,
sparkException.getCause().getClass());
Review Comment:
no-error-table test asserts only the `HoodieRecordCreationException`
wrapper, which any creation failure produces. There are multi field delete
tests above that assert only `RECORD_CREATION` and `ErrorEvent` carries no
throwable.
Can we parameterize the test over three shapes and assert the
`IllegalArgumentException` message instead?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]