lokeshj1703 commented on code in PR #20049:
URL: https://github.com/apache/hudi/pull/20049#discussion_r4182894602
##########
hudi-common/src/main/java/org/apache/hudi/common/util/OrderingValues.java:
##########
@@ -91,6 +92,21 @@ public static Comparable getDefault() {
/**
* Returns whether the given {@code orderingValue} is default.
*/
+ /**
Review Comment:
Fixed in ea87bd7 — `isMissing` now sits below `isDefault`, so each javadoc
is attached to its own method.
##########
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:
Good catch, this was a real hole. Fixed in ea87bd7: both the merge mode and
the payload class now prefer the table config, with `cfg` only as the fallback.
Worth noting the payload half mattered as much as the merge-mode half. My
first attempt fixed only `recordMergeMode` and left `payloadClassName` deriving
from the stale `cfg.recordMergeMode`, which resolves to
`OverwriteWithLatestAvroPayload` and skipped the check anyway —
`testNullOrderingValueRejectedWhenOnlyTheTableSaysEventTime` failed against
that half-fix with "Expected SparkException to be thrown, but nothing was
thrown", which is what caught it.
##########
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:
Confirmed and fixed in ea87bd7, including the payload point — fixing only
the merge-mode clause was not enough, since `payloadClassName` derived from the
same stale value and resolved to `OverwriteWithLatestAvroPayload`. Both now
come from the table config. The new test fails against the merge-mode-only fix.
##########
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:
Done in ea87bd7 — the no-error-table case now asserts the
`IllegalArgumentException` and its message, not just the
`HoodieRecordCreationException` wrapper.
##########
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:
Right, the false side had no coverage. ea87bd7 adds
`testNullOrderingValueAllowedWhenNoOrderingValueRequired`, parameterized over
`COMMIT_TIME_ORDERING` and `OverwriteWithLatestAvroPayload`, asserting the
null-ordering record is returned and nothing is quarantined. There is also a
case that pins the merge mode coming from the table config rather than from
`cfg`.
--
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]