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]

Reply via email to