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]

Reply via email to