hudi-agent commented on code in PR #19660:
URL: https://github.com/apache/hudi/pull/19660#discussion_r3811749592


##########
hudi-common/src/main/java/org/apache/hudi/common/model/debezium/AbstractDebeziumAvroPayload.java:
##########
@@ -72,15 +75,37 @@ public Option<IndexedRecord> getInsertValue(Schema schema) 
throws IOException {
     return insertValue.isPresent() ? handleDeleteOperation(insertValue.get()) 
: Option.empty();
   }
 
+  @Override
+  public Option<IndexedRecord> getInsertValue(Schema schema, Properties 
properties) throws IOException {
+    // Pin to the Debezium delete-op handling; DefaultHoodieRecordPayload's 
properties-aware variant
+    // (event-time tracking, DELETE_KEY/DELETE_MARKER) must not replace it
+    return getInsertValue(schema);
+  }
+
   @Override
   public Option<IndexedRecord> combineAndGetUpdateValue(IndexedRecord 
currentValue, Schema schema) throws IOException {
+    return combineAndGetUpdateValue(currentValue, schema, new Properties());
+  }
+
+  @Override
+  public Option<IndexedRecord> combineAndGetUpdateValue(IndexedRecord 
currentValue, Schema schema, Properties properties) throws IOException {
     // Step 1: If the time occurrence of the current record in storage is 
higher than the time occurrence of the
     // insert record (including a delete record), pick the current record.
     Option<IndexedRecord> insertValue = getRecord(schema);
     if (!insertValue.isPresent()) {
       return Option.empty();
     }
-    if (shouldPickCurrentRecord(currentValue, insertValue.get(), schema)) {
+    String[] orderingFields = ConfigUtils.getOrderingFields(properties);
+    boolean pickCurrentRecord;
+    if (orderingFields == null || orderingFields.length != 1 || 
orderingFields[0].equals(getConnectorOrderingField())) {
+      // No ordering field configured, a composite ordering (not supported 
yet), or the connector's own column:
+      // use the connector-specific comparison (MySQL's "file.pos" seq needs 
segment-wise numeric compare;
+      // a plain Comparable is lexicographic)
+      pickCurrentRecord = shouldPickCurrentRecord(currentValue, 
insertValue.get(), schema);
+    } else {
+      pickCurrentRecord = !needUpdatingPersistedRecord(currentValue, 
insertValue, properties);

Review Comment:
   🤖 Worth noting the incoming-null case doesn't actually reach a bare NPE 
here: `needUpdatingPersistedRecord` populates `incomingOrderingVals[i]` and 
explicitly checks `if (incomingOrderingVals[i] == null) throw new 
HoodieDebeziumAvroPayloadException(...)` before any `compareTo`, so a null 
incoming ordering value fails loudly rather than NPE-ing (only the persisted 
side reaching the compare, and it's null-guarded via the early `return true`).
   
   The delete-op angle is still a fair thing to check though — the behavior 
change is that a Debezium delete whose flattened after-image doesn't populate 
the configured ordering column would now *throw and fail the merge* on this 
path, whereas `shouldPickCurrentRecord` handled the missing-column case 
gracefully. Since Debezium deletes typically carry the before-image (so data 
columns are populated), this may be fine in practice, but a delete-op test 
through the else branch would be good to pin the intended behavior.



##########
hudi-common/src/test/java/org/apache/hudi/common/model/debezium/TestMySqlDebeziumAvroPayload.java:
##########
@@ -161,6 +163,246 @@ public void testIsCurrentSeqLatest(String currentSeq, 
String newSeq, boolean exp
     assertEquals(expectedResult, 
MySqlDebeziumAvroPayload.isCurrentSeqLatest(currentSeq, newSeq));
   }
 
+  @Test
+  public void testMergeWithConfiguredOrderingFieldStoredRecordWins() throws 
IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // Incoming has a HIGHER seq but a LOWER configured ordering value -> 
stored record must win,
+    // proving the configured field (not the hardcoded seq) decides.
+    MySqlDebeziumAvroPayload payload = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.UPDATE, 
"00005.100", 50L), 50L);
+    GenericRecord existing = createRecordWithOrdering(schema, 1, 
Operation.INSERT, "00001.100", 99L);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps("event_ts"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.INSERT, DebeziumConstants.ADDED_SEQ_COL_NAME, "00001.100", 99L);
+  }
+
+  @Test
+  public void testMergeWithConfiguredOrderingFieldIncomingRecordWins() throws 
IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // Incoming has a LOWER seq but a HIGHER configured ordering value -> 
incoming wins.
+    MySqlDebeziumAvroPayload payload = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.UPDATE, 
"00000.100", 120L), 120L);
+    GenericRecord existing = createRecordWithOrdering(schema, 1, 
Operation.INSERT, "00001.100", 99L);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps("event_ts"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.UPDATE, DebeziumConstants.ADDED_SEQ_COL_NAME, "00000.100", 120L);
+  }
+
+  @Test
+  public void testMergeWithWhitespacePaddedConfiguredOrderingFieldIsTrimmed() 
throws IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // Trimming is what keeps a padded configured field resolving to a real 
column: without it,
+    // " event_ts " resolves no field (null ordering value) and the 
null-incoming guard throws.
+    MySqlDebeziumAvroPayload payload = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.UPDATE, 
"00005.100", 50L), 50L);
+    GenericRecord existing = createRecordWithOrdering(schema, 1, 
Operation.INSERT, "00001.100", 99L);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps(" event_ts "));
+    // Stored row wins on event_ts (99 > 50); the seq order (00005 > 00001) 
would have picked the incoming
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.INSERT, DebeziumConstants.ADDED_SEQ_COL_NAME, "00001.100", 99L);
+  }
+
+  @Test
+  public void testMergeWithCompositeOrderingFieldsSecondFieldDecides() throws 
IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // First ordering field ties (99 == 99); the second decides (7 > 5) -> 
incoming wins,
+    // even though the seq order (00001 < 00005) says otherwise.
+    MySqlDebeziumAvroPayload payload = new MySqlDebeziumAvroPayload(
+        DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.UPDATE, DebeziumConstants.ADDED_SEQ_COL_NAME, "00001.100", 99L, 7L), 
99L);
+    GenericRecord existing = 
DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.INSERT, DebeziumConstants.ADDED_SEQ_COL_NAME, "00005.100", 99L, 5L);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps("event_ts,event_ts2"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.UPDATE, DebeziumConstants.ADDED_SEQ_COL_NAME, "00001.100", 99L, 7L);
+  }
+
+  @Test
+  public void testMergeWithCompositeOrderingFieldsFirstFieldDecides() throws 
IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // First ordering field decides outright when it differs (98 < 99) -> 
stored record kept,
+    // regardless of the second field (100 > 5 would say otherwise).
+    MySqlDebeziumAvroPayload payload = new MySqlDebeziumAvroPayload(
+        DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.UPDATE, DebeziumConstants.ADDED_SEQ_COL_NAME, "00009.100", 98L, 
100L), 98L);
+    GenericRecord existing = 
DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.INSERT, DebeziumConstants.ADDED_SEQ_COL_NAME, "00005.100", 99L, 5L);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps("event_ts,event_ts2"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.INSERT, DebeziumConstants.ADDED_SEQ_COL_NAME, "00005.100", 99L, 5L);
+  }
+
+  @Test
+  public void 
testMergeWithCompositeOrderingNullPersistedSecondFieldFirstDecides() throws 
IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // Persisted second element is null but the FIRST field strictly decides 
(stored 99 > incoming 98):
+    // the element-wise loop returns at the first field without ever touching 
the null.
+    MySqlDebeziumAvroPayload payload = new MySqlDebeziumAvroPayload(
+        DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.UPDATE, DebeziumConstants.ADDED_SEQ_COL_NAME, "00009.100", 98L, 
100L), 98L);
+    GenericRecord existing = 
DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.INSERT, DebeziumConstants.ADDED_SEQ_COL_NAME, "00005.100", 99L, null);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps("event_ts,event_ts2"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.INSERT, DebeziumConstants.ADDED_SEQ_COL_NAME, "00005.100", 99L, null);
+  }
+
+  @Test
+  public void 
testMergeWithCompositeOrderingNullPersistedSecondFieldOnTieTakesIncoming() 
throws IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // First field ties and the persisted second element is null (e.g. 
bootstrapped rows): the incoming
+    // record wins — the null-persisted rule applies per element, no NPE.
+    MySqlDebeziumAvroPayload payload = new MySqlDebeziumAvroPayload(
+        DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.UPDATE, DebeziumConstants.ADDED_SEQ_COL_NAME, "00001.100", 99L, 7L), 
99L);
+    GenericRecord existing = 
DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.INSERT, DebeziumConstants.ADDED_SEQ_COL_NAME, "00005.100", 99L, null);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps("event_ts,event_ts2"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.UPDATE, DebeziumConstants.ADDED_SEQ_COL_NAME, "00001.100", 99L, 7L);
+  }
+
+  @Test
+  public void testMergeWithMySqlNativeCompositeOrderingPosDecides() throws 
IOException {
+    Schema schema = createMySqlNativeOrderingSchema();
+    // The composite every MySQL table actually configures (ORDERING_FIELDS = 
_event_bin_file,_event_pos):
+    // mixed element types (String file, long pos). Same binlog file -> the 
numeric pos decides.
+    GenericRecord incoming = createMySqlNativeOrderingRecord(schema, 1, 
Operation.UPDATE, "mysql-bin.000001", 200L);
+    MySqlDebeziumAvroPayload payload = new MySqlDebeziumAvroPayload(incoming, 
200L);
+    GenericRecord existing = createMySqlNativeOrderingRecord(schema, 1, 
Operation.INSERT, "mysql-bin.000001", 100L);
+    GenericRecord merged = (GenericRecord) 
payload.combineAndGetUpdateValue(existing, schema,
+        orderingProps(MySqlDebeziumAvroPayload.ORDERING_FIELDS)).get();
+    assertEquals(200L, merged.get(DebeziumConstants.FLATTENED_POS_COL_NAME));
+    assertEquals(Operation.UPDATE.op, 
merged.get(DebeziumConstants.FLATTENED_OP_COL_NAME).toString());
+  }
+
+  @Test
+  public void testMergeWithMySqlNativeCompositeOrderingFileDecides() throws 
IOException {
+    Schema schema = createMySqlNativeOrderingSchema();
+    // Binlog file differs -> the String file name decides before pos is 
consulted; stored record kept.
+    GenericRecord incoming = createMySqlNativeOrderingRecord(schema, 1, 
Operation.UPDATE, "mysql-bin.000001", 999L);
+    MySqlDebeziumAvroPayload payload = new MySqlDebeziumAvroPayload(incoming, 
999L);
+    GenericRecord existing = createMySqlNativeOrderingRecord(schema, 1, 
Operation.INSERT, "mysql-bin.000002", 5L);
+    GenericRecord merged = (GenericRecord) 
payload.combineAndGetUpdateValue(existing, schema,
+        orderingProps(MySqlDebeziumAvroPayload.ORDERING_FIELDS)).get();
+    assertEquals(5L, merged.get(DebeziumConstants.FLATTENED_POS_COL_NAME));
+    assertEquals(Operation.INSERT.op, 
merged.get(DebeziumConstants.FLATTENED_OP_COL_NAME).toString());
+  }
+
+  @Test
+  public void testPreCombineWithCompositeOrderingFields() {
+    Schema schema = createSchemaWithOrderingField();
+    Properties props = orderingProps("event_ts,event_ts2");
+    // Composite orderingVal, first field tied: the second field decides the 
dedup winner.
+    MySqlDebeziumAvroPayload lower = new MySqlDebeziumAvroPayload(
+        DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.INSERT, DebeziumConstants.ADDED_SEQ_COL_NAME, "00005.100", 99L, 5L),
+        OrderingValues.create(new Comparable[] {99L, 5L}));
+    MySqlDebeziumAvroPayload higher = new MySqlDebeziumAvroPayload(
+        DebeziumOrderingTestFixtures.recordWithCompositeOrdering(schema, 1, 
Operation.UPDATE, DebeziumConstants.ADDED_SEQ_COL_NAME, "00001.100", 99L, 7L),
+        OrderingValues.create(new Comparable[] {99L, 7L}));
+    assertEquals(higher, higher.preCombine(lower, props));
+    assertEquals(higher, lower.preCombine(higher, props));
+  }
+
+  @Test
+  public void testMergeWithoutOrderingFieldLateRecordLosesOnSeq() throws 
IOException {
+    // Properties present but no ordering field configured -> legacy hardcoded 
seq comparison.
+    GenericRecord lateRecord = createRecord(3, Operation.UPDATE, "00000.222");
+    MySqlDebeziumAvroPayload payload = new 
MySqlDebeziumAvroPayload(lateRecord, "00000.222");
+    GenericRecord existingRecord = createRecord(3, Operation.INSERT, 
"00001.111");
+    Option<IndexedRecord> mergedRecord = 
payload.combineAndGetUpdateValue(existingRecord, avroSchema, new Properties());
+    validateRecord(mergedRecord, 3, Operation.INSERT, "00001.111");
+  }
+
+  @Test
+  public void testMergeWithoutOrderingFieldFreshRecordWinsOnSeq() throws 
IOException {
+    GenericRecord freshRecord = createRecord(3, Operation.UPDATE, "00002.11");
+    MySqlDebeziumAvroPayload payload = new 
MySqlDebeziumAvroPayload(freshRecord, "00002.11");
+    GenericRecord existingRecord = createRecord(3, Operation.INSERT, 
"00001.111");
+    Option<IndexedRecord> mergedRecord = 
payload.combineAndGetUpdateValue(existingRecord, avroSchema, new Properties());
+    validateRecord(mergedRecord, 3, Operation.UPDATE, "00002.11");
+  }
+
+  @Test
+  public void testPreCombineEqualConfiguredOrderingValuesTieGoesToNewer() {
+    Schema schema = createSchemaWithOrderingField();
+    // Two records with EQUAL configured ordering values: the seq parser would 
throw here
+    // ("99".split(".") has no position segment -> 
ArrayIndexOutOfBoundsException); the
+    // configured-field compare must not throw, and the newer payload wins the 
tie.
+    MySqlDebeziumAvroPayload older = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.INSERT, 
"00001.111", 99L), 99L);
+    MySqlDebeziumAvroPayload newer = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.UPDATE, 
"00002.111", 99L), 99L);
+    assertEquals(newer, newer.preCombine(older, orderingProps("event_ts")));
+  }
+
+  @Test
+  public void testPreCombineConfiguredOrderingFieldOverridesSeq() {
+    Schema schema = createSchemaWithOrderingField();
+    // The configured field decides, in both directions, even when the seq 
order says otherwise.
+    MySqlDebeziumAvroPayload lowerTsHigherSeq = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.UPDATE, 
"00005.100", 50L), 50L);
+    MySqlDebeziumAvroPayload higherTsLowerSeq = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.UPDATE, 
"00000.100", 120L), 120L);
+    assertEquals(higherTsLowerSeq, 
higherTsLowerSeq.preCombine(lowerTsHigherSeq, orderingProps("event_ts")));
+    assertEquals(higherTsLowerSeq, 
lowerTsHigherSeq.preCombine(higherTsLowerSeq, orderingProps("event_ts")));
+  }
+
+  @Test
+  public void testPreCombineDeleteRecordKeepsNaturalOrder() {
+    Schema schema = createSchemaWithOrderingField();
+    MySqlDebeziumAvroPayload newer = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.UPDATE, 
"00002.111", 99L), 99L);
+    MySqlDebeziumAvroPayload empty = new 
MySqlDebeziumAvroPayload(Option.empty());
+    assertEquals(newer, newer.preCombine(empty, orderingProps("event_ts")));
+  }
+
+  @Test
+  public void testPreCombineWithoutOrderingFieldUsesSeqCompare() {
+    Schema schema = createSchemaWithOrderingField();
+    // No ordering field configured -> the legacy numeric seq comparison still 
applies. On this
+    // path real ingestion sets orderingVal from the seq column (precombine = 
_event_seq), so
+    // construct accordingly.
+    MySqlDebeziumAvroPayload olderSeq = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.INSERT, 
"00001.111", 99L), "00001.111");
+    MySqlDebeziumAvroPayload newerSeq = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.UPDATE, 
"00002.111", 99L), "00002.111");
+    assertEquals(newerSeq, newerSeq.preCombine(olderSeq, new Properties()));
+  }
+
+  @Test
+  public void testMergeWithSeqOrderingFieldKeepsNumericCompare() throws 
IOException {
+    // Legacy config: ordering field = _event_seq (v8 tables / precombine 
convention). The "file.pos"
+    // segments must be compared numerically. Unpadded values discriminate: 
lexicographically
+    // "2.11" > "10.111", numerically file 2 < file 10 -> the stored record 
must win.
+    Properties props = orderingProps(DebeziumConstants.ADDED_SEQ_COL_NAME);
+    GenericRecord incoming = createRecord(2, Operation.UPDATE, "2.11");
+    MySqlDebeziumAvroPayload payload = new MySqlDebeziumAvroPayload(incoming, 
"2.11");
+    GenericRecord existing = createRecord(2, Operation.INSERT, "10.111");
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
avroSchema, props);
+    validateRecord(merged, 2, Operation.INSERT, "10.111");
+  }
+
+  @Test
+  public void testPreCombineWithSeqOrderingFieldKeepsNumericCompare() {
+    // Same carve-out on the dedup path: orderingVal carries seq strings when 
precombine = _event_seq.
+    // Numerically file 10 > file 9 -> newer wins; lexicographic "9.111" > 
"10.2" would invert it.
+    Schema schema = createSchemaWithOrderingField();
+    MySqlDebeziumAvroPayload olderSeq = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.INSERT, 
"9.111", 99L), "9.111");
+    MySqlDebeziumAvroPayload newerSeq = new 
MySqlDebeziumAvroPayload(createRecordWithOrdering(schema, 1, Operation.UPDATE, 
"10.2", 99L), "10.2");
+    assertEquals(newerSeq, newerSeq.preCombine(olderSeq, 
orderingProps(DebeziumConstants.ADDED_SEQ_COL_NAME)));
+  }
+
+  private Schema createMySqlNativeOrderingSchema() {
+    return Schema.createRecord("test_mysql_native_ordering", null, 
"test_namespace", false, Arrays.asList(
+        new Schema.Field(KEY_FIELD_NAME, Schema.create(Schema.Type.INT), "", 
0),
+        new Schema.Field(DebeziumConstants.FLATTENED_OP_COL_NAME,
+            Schema.createUnion(Schema.create(Schema.Type.NULL), 
Schema.create(Schema.Type.STRING)), "", null),
+        new Schema.Field(DebeziumConstants.FLATTENED_FILE_COL_NAME,
+            Schema.createUnion(Schema.create(Schema.Type.NULL), 
Schema.create(Schema.Type.STRING)), "", null),
+        new Schema.Field(DebeziumConstants.FLATTENED_POS_COL_NAME,
+            Schema.createUnion(Schema.create(Schema.Type.NULL), 
Schema.create(Schema.Type.LONG)), "", null)
+    ));
+  }
+
+  private GenericRecord createMySqlNativeOrderingRecord(Schema schema, int 
key, Operation op, String binlogFile, Long pos) {
+    GenericRecord record = new GenericData.Record(schema);
+    record.put(KEY_FIELD_NAME, key);
+    record.put(DebeziumConstants.FLATTENED_OP_COL_NAME, Objects.toString(op, 
null));
+    // Utf8, matching what the persisted side carries in production (both 
sides come from Avro decoding)
+    record.put(DebeziumConstants.FLATTENED_FILE_COL_NAME, new 
Utf8(binlogFile));
+    record.put(DebeziumConstants.FLATTENED_POS_COL_NAME, pos);
+    return record;
+  }
+
+  private Schema createSchemaWithOrderingField() {
+    return 
DebeziumOrderingTestFixtures.schemaWithOrderingField(DebeziumConstants.ADDED_SEQ_COL_NAME,
 Schema.Type.STRING);
+  }
+
+  private GenericRecord createRecordWithOrdering(Schema schema, int 
primaryKeyValue, @Nullable Operation op,

Review Comment:
   🤖 nit: this private helper is a one-liner that just calls 
`DebeziumOrderingTestFixtures.orderingProps` — could you inline the call sites 
directly? As written, a reader has to chase the indirection to confirm there's 
no MySQL-specific behaviour added here.
   
   <sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag 
quality.</i></sub>



##########
hudi-common/src/test/java/org/apache/hudi/common/model/debezium/TestPostgresDebeziumAvroPayload.java:
##########
@@ -159,6 +159,122 @@ public void testInvalidIncomingRecord() {
         "should have thrown because LSN value of the incoming record is null");
   }
 
+  @Test
+  public void testMergeWithConfiguredOrderingFieldStoredRecordWins() throws 
IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // Incoming has a HIGHER LSN but a LOWER configured ordering value -> 
stored record must win,
+    // proving the configured field (not the hardcoded LSN) decides.
+    PostgresDebeziumAvroPayload payload = new 
PostgresDebeziumAvroPayload(createRecordWithOrdering(schema, 1, 
Operation.UPDATE, 200L, 50L), 50L);
+    GenericRecord existing = createRecordWithOrdering(schema, 1, 
Operation.INSERT, 100L, 99L);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps("event_ts"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.INSERT, DebeziumConstants.FLATTENED_LSN_COL_NAME, 100L, 99L);
+  }
+
+  @Test
+  public void testMergeWithConfiguredOrderingFieldIncomingRecordWins() throws 
IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // Incoming has a LOWER LSN but a HIGHER configured ordering value -> 
incoming wins.
+    PostgresDebeziumAvroPayload payload = new 
PostgresDebeziumAvroPayload(createRecordWithOrdering(schema, 1, 
Operation.UPDATE, 90L, 120L), 120L);
+    GenericRecord existing = createRecordWithOrdering(schema, 1, 
Operation.INSERT, 100L, 99L);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps("event_ts"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.UPDATE, DebeziumConstants.FLATTENED_LSN_COL_NAME, 90L, 120L);
+  }
+
+  @Test
+  public void testMergeWithConfiguredOrderingFieldTieGoesToIncoming() throws 
IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // Tie on the configured ordering value -> incoming wins (mirrors legacy 
LSN tie semantics).
+    PostgresDebeziumAvroPayload payload = new 
PostgresDebeziumAvroPayload(createRecordWithOrdering(schema, 1, 
Operation.UPDATE, 90L, 99L), 99L);
+    GenericRecord existing = createRecordWithOrdering(schema, 1, 
Operation.INSERT, 100L, 99L);
+    Option<IndexedRecord> merged = payload.combineAndGetUpdateValue(existing, 
schema, orderingProps("event_ts"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.UPDATE, DebeziumConstants.FLATTENED_LSN_COL_NAME, 90L, 99L);
+  }
+
+  @Test
+  public void testMergeWithConfiguredOrderingFieldNullPersistedTakesIncoming() 
throws IOException {
+    Schema schema = createSchemaWithOrderingField();
+    // Stored record has a null ordering value (e.g. bootstrapped rows) -> 
incoming wins.
+    PostgresDebeziumAvroPayload payload = new 
PostgresDebeziumAvroPayload(createRecordWithOrdering(schema, 1, 
Operation.UPDATE, 90L, 10L), 10L);
+    GenericRecord bootstrapped = createRecordWithOrdering(schema, 1, null, 
null, null);
+    Option<IndexedRecord> merged = 
payload.combineAndGetUpdateValue(bootstrapped, schema, 
orderingProps("event_ts"));
+    DebeziumOrderingTestFixtures.validateOrderingRecord(merged, 1, 
Operation.UPDATE, DebeziumConstants.FLATTENED_LSN_COL_NAME, 90L, 10L);
+  }
+
+  @Test
+  public void testMergeWithConfiguredOrderingFieldNullIncomingThrows() {
+    Schema schema = createSchemaWithOrderingField();
+    GenericRecord incoming = createRecordWithOrdering(schema, 2, 
Operation.UPDATE, 200L, null);
+    PostgresDebeziumAvroPayload payload = new 
PostgresDebeziumAvroPayload(incoming, 0L);
+    GenericRecord existing = createRecordWithOrdering(schema, 2, 
Operation.INSERT, 100L, 99L);
+    assertThrows(HoodieDebeziumAvroPayloadException.class,
+        () -> payload.combineAndGetUpdateValue(existing, schema, 
orderingProps("event_ts")),
+        "null configured ordering value in the incoming record must fail 
loudly, not NPE");
+  }
+
+  @Test
+  public void testMergeWithoutOrderingFieldLateRecordLosesOnLsn() throws 
IOException {
+    // Properties present but no ordering field configured -> legacy hardcoded 
LSN comparison.
+    GenericRecord lateRecord = createRecord(3, Operation.UPDATE, 98L);
+    PostgresDebeziumAvroPayload payload = new 
PostgresDebeziumAvroPayload(lateRecord, 98L);
+    GenericRecord existingRecord = createRecord(3, Operation.INSERT, 99L);
+    Option<IndexedRecord> mergedRecord = 
payload.combineAndGetUpdateValue(existingRecord, avroSchema, new Properties());
+    validateRecord(mergedRecord, 3, Operation.INSERT, 99L);
+  }
+
+  @Test
+  public void testMergeWithoutOrderingFieldFreshRecordWinsOnLsn() throws 
IOException {
+    GenericRecord freshRecord = createRecord(3, Operation.UPDATE, 100L);
+    PostgresDebeziumAvroPayload payload = new 
PostgresDebeziumAvroPayload(freshRecord, 100L);
+    GenericRecord existingRecord = createRecord(3, Operation.INSERT, 99L);
+    Option<IndexedRecord> mergedRecord = 
payload.combineAndGetUpdateValue(existingRecord, avroSchema, new Properties());
+    validateRecord(mergedRecord, 3, Operation.UPDATE, 100L);
+  }
+
+  @Test
+  public void testMergeWithToastedValuesUnderConfiguredOrdering() throws 
IOException {
+    // Toast-column merging must keep working when the ordering decision goes 
through the
+    // configured ordering field instead of the hardcoded LSN.
+    Schema schema = SchemaBuilder.builder()
+        .record("test_toast_ordering")
+        .namespace("test_namespace")
+        .fields()
+        
.name(DebeziumConstants.FLATTENED_LSN_COL_NAME).type().longType().noDefault()
+        .name("event_ts").type().longType().noDefault()
+        .name("string_col").type().stringType().noDefault()
+        .endRecord();
+
+    GenericRecord oldVal = new GenericData.Record(schema);
+    oldVal.put(DebeziumConstants.FLATTENED_LSN_COL_NAME, 100L);
+    oldVal.put("event_ts", 1L);
+    oldVal.put("string_col", "valid string value");
+
+    // Incoming loses on LSN but wins on the configured ordering field, and 
carries a toasted column
+    GenericRecord newVal = new GenericData.Record(schema);
+    newVal.put(DebeziumConstants.FLATTENED_LSN_COL_NAME, 90L);
+    newVal.put("event_ts", 2L);
+    newVal.put("string_col", 
PostgresDebeziumAvroPayload.DEBEZIUM_TOASTED_VALUE);
+
+    PostgresDebeziumAvroPayload payload = new 
PostgresDebeziumAvroPayload(Option.of(newVal));
+    GenericRecord merged = (GenericRecord) payload
+        .combineAndGetUpdateValue(oldVal, schema, 
orderingProps("event_ts")).get();
+
+    assertEquals(2L, (long) merged.get("event_ts"));
+    assertEquals("valid string value", merged.get("string_col"));
+  }
+
+  private Schema createSchemaWithOrderingField() {
+    return 
DebeziumOrderingTestFixtures.schemaWithOrderingField(DebeziumConstants.FLATTENED_LSN_COL_NAME,
 Schema.Type.LONG);
+  }
+
+  private GenericRecord createRecordWithOrdering(Schema schema, int 
primaryKeyValue, @Nullable Operation op,
+                                                 @Nullable Long lsnValue, 
@Nullable Long orderingValue) {
+    return DebeziumOrderingTestFixtures.recordWithOrdering(schema, 
primaryKeyValue, op, DebeziumConstants.FLATTENED_LSN_COL_NAME, lsnValue, 
orderingValue);

Review Comment:
   🤖 nit: same delegation-only wrapper as in the MySQL test — could you call 
`DebeziumOrderingTestFixtures.orderingProps` directly at each use site to make 
it clear there's no Postgres-specific logic being applied?
   
   <sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag 
quality.</i></sub>



-- 
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