This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new d7c5d42c46 Revert "[core] Fix LookupMergeFunction to use 
sequence.field for picking high level records" (#8354)
d7c5d42c46 is described below

commit d7c5d42c469d4d63a17114ac0883133635545aae
Author: Faiz <[email protected]>
AuthorDate: Sat Jun 27 15:33:27 2026 +0800

    Revert "[core] Fix LookupMergeFunction to use sequence.field for picking 
high level records" (#8354)
    
    Reverts apache/paimon#7221
    
    I believe the scenario fixed by that PR is not reachable under the
    current lookup compaction layout, and the added logic makes the lookup
    merge semantics more complex than necessary.
    
    The original PR is trying to fix this scenario:
    ```text
      L1: sequence = 7
      L2: sequence = 8
      L0: sequence = 6
    ```
    
    However, in lookup mode, high-level files are expected to represent
    materialized merged states. A lower high level, such as L1, should
    shadow older higher levels, such as L2. If an L1 record exists for a key
    after compaction involving L2, it should already have been produced by
    merging the relevant L0 and higher-level state, and therefore should
    carry the latest state according to the configured merge function and
    UserDefinedSeqComparator (the max seqNumber).
    
    I've checked this PR's tests, the PR manually constructs records like
    ```text
      L1: sequence = 7
      L2: sequence = 8
      L0: sequence = 6
    ```
    and feeds them directly into LookupChangelogMergeFunctionWrapper.
    These tests bypass the actual writer, level manager, compaction picker,
    rewrite/upgrade strategy, and lookup compaction flow.
---
 .../LookupChangelogMergeFunctionWrapper.java       |   8 -
 .../mergetree/compact/LookupMergeFunction.java     |  37 +---
 .../LookupChangelogMergeFunctionWrapperTest.java   | 186 ---------------------
 .../apache/paimon/flink/DeletionVectorITCase.java  |  22 +++
 .../paimon/flink/LookupChangelogWithAggITCase.java |  32 ++++
 5 files changed, 60 insertions(+), 225 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java
 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java
index 04d9141602..7283a3030d 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java
@@ -88,14 +88,6 @@ public class LookupChangelogMergeFunctionWrapper<T>
         this.lookupStrategy = lookupStrategy;
         this.deletionVectorsMaintainer = deletionVectorsMaintainer;
         this.comparator = createSequenceComparator(userDefinedSeqComparator);
-        // Only set sequence comparator when user-defined sequence field is 
configured
-        // to preserve original behavior (pick by level) when sequence.field 
is not set.
-        // Note: We use the same comparator for both insertInto and 
pickHighLevel.
-        // The comparator's semantics (ascending/descending) are already 
handled correctly
-        // by UserDefinedSeqComparator based on sequence.field.sort-order 
configuration.
-        if (userDefinedSeqComparator != null) {
-            this.mergeFunction.setSequenceComparator(this.comparator);
-        }
     }
 
     @Override
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeFunction.java
 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeFunction.java
index eb2c8859b7..cb494b86d0 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeFunction.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeFunction.java
@@ -27,9 +27,7 @@ import org.apache.paimon.utils.CloseableIterator;
 
 import javax.annotation.Nullable;
 
-import java.util.ArrayList;
 import java.util.Comparator;
-import java.util.List;
 
 /**
  * A {@link MergeFunction} for lookup, this wrapper only considers the latest 
high level record,
@@ -43,7 +41,6 @@ public class LookupMergeFunction implements 
MergeFunction<KeyValue> {
     private final KeyValueBuffer candidates;
     private boolean containLevel0;
     private InternalRow currentKey;
-    @Nullable private Comparator<KeyValue> sequenceComparator;
 
     public LookupMergeFunction(
             MergeFunction<KeyValue> mergeFunction,
@@ -55,11 +52,6 @@ public class LookupMergeFunction implements 
MergeFunction<KeyValue> {
         this.candidates = KeyValueBuffer.createHybridBuffer(options, keyType, 
valueType, ioManager);
     }
 
-    /** Set the sequence comparator for picking high level records. */
-    public void setSequenceComparator(@Nullable Comparator<KeyValue> 
sequenceComparator) {
-        this.sequenceComparator = sequenceComparator;
-    }
-
     @Override
     public void reset() {
         candidates.reset();
@@ -91,16 +83,9 @@ public class LookupMergeFunction implements 
MergeFunction<KeyValue> {
                 if (kv.level() <= 0) {
                     continue;
                 }
-                if (highLevel == null) {
-                    highLevel = kv;
-                } else if (sequenceComparator != null) {
-                    // When sequence comparator is set, use it to pick the 
record with highest
-                    // sequence value, which represents the latest record
-                    if (sequenceComparator.compare(kv, highLevel) > 0) {
-                        highLevel = kv;
-                    }
-                } else if (kv.level() < highLevel.level()) {
-                    // Without sequence comparator, fall back to picking the 
minimum level
+                // For high-level comparison logic (not involving Level 0), 
only the value of the
+                // minimum Level should be selected
+                if (highLevel == null || kv.level() < highLevel.level()) {
                     highLevel = kv;
                 }
             }
@@ -122,28 +107,18 @@ public class LookupMergeFunction implements 
MergeFunction<KeyValue> {
     public KeyValue getResult() {
         mergeFunction.reset();
         KeyValue highLevel = pickHighLevel();
-
-        // Collect records to merge: level-0 records and the picked high level 
record
-        List<KeyValue> toMerge = new ArrayList<>();
         try (CloseableIterator<KeyValue> iterator = candidates.iterator()) {
             while (iterator.hasNext()) {
                 KeyValue kv = iterator.next();
+                // records that has not been stored on the disk yet, such as 
the data in the write
+                // buffer being at level -1
                 if (kv.level() <= 0 || kv == highLevel) {
-                    toMerge.add(kv);
+                    mergeFunction.add(kv);
                 }
             }
         } catch (Exception e) {
             throw new RuntimeException(e);
         }
-
-        // When sequence comparator is set, sort by sequence so highest 
sequence is added last
-        if (sequenceComparator != null) {
-            toMerge.sort(sequenceComparator);
-        }
-
-        for (KeyValue kv : toMerge) {
-            mergeFunction.add(kv);
-        }
         return mergeFunction.getResult();
     }
 
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java
index 470e84525a..57d99557ca 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java
@@ -556,190 +556,4 @@ public class LookupChangelogMergeFunctionWrapperTest {
         kv = result.result();
         assertThat(kv.value().getInt(0)).isEqualTo(3);
     }
-
-    /**
-     * Test that sequence.field is correctly used to pick the high level 
record with the highest
-     * sequence value, even when it's at a higher level number.
-     *
-     * <p>Scenario: L1 has older sequence (7), L2 has newer sequence (8), L0 
has oldest (6). The
-     * correct behavior should pick L2 (sequence=8) as the high level record.
-     */
-    @Test
-    public void testSequenceFieldWithMultipleLevels() {
-        // Define value type with sequence field as the second column
-        RowType valueType =
-                RowType.builder()
-                        .fields(
-                                new DataType[] {DataTypes.INT(), 
DataTypes.INT()},
-                                new String[] {"value", "sequence"})
-                        .build();
-
-        // Create user-defined sequence comparator on the second field
-        UserDefinedSeqComparator userDefinedSeqComparator =
-                UserDefinedSeqComparator.create(
-                        valueType,
-                        CoreOptions.fromMap(ImmutableMap.of("sequence.field", 
"sequence")));
-        assertThat(userDefinedSeqComparator).isNotNull();
-
-        Map<InternalRow, KeyValue> highLevel = new HashMap<>();
-
-        LookupChangelogMergeFunctionWrapper function =
-                new LookupChangelogMergeFunctionWrapper(
-                        LookupMergeFunction.wrap(
-                                DeduplicateMergeFunction.factory(), null, 
null, null),
-                        highLevel::get,
-                        null,
-                        LookupStrategy.from(false, true, false, false),
-                        null,
-                        userDefinedSeqComparator);
-
-        // Test scenario:
-        // L1: (key=1, value=100, sequence=7)  <- Level 1, but older sequence
-        // L2: (key=1, value=200, sequence=8)  <- Level 2, but newer sequence 
(should be picked!)
-        // L0: (key=1, value=50,  sequence=6)  <- Level 0, oldest sequence
-
-        function.reset();
-        function.add(
-                new KeyValue()
-                        .replace(row(1), 1, INSERT, row(100, 7))
-                        .setLevel(1)); // Level 1, seq=7
-        function.add(
-                new KeyValue()
-                        .replace(row(1), 1, INSERT, row(200, 8))
-                        .setLevel(2)); // Level 2, seq=8
-        function.add(
-                new KeyValue()
-                        .replace(row(1), 2, INSERT, row(50, 6))
-                        .setLevel(0)); // Level 0, seq=6
-
-        ChangelogResult result = function.getResult();
-        assertThat(result).isNotNull();
-
-        KeyValue kv = result.result();
-        assertThat(kv).isNotNull();
-
-        // Should return the record with highest sequence (seq=8 from L2)
-        int actualSequence = kv.value().getInt(1);
-        int actualValue = kv.value().getInt(0);
-
-        assertThat(actualSequence)
-                .as("Should return record with highest sequence field (8)")
-                .isEqualTo(8);
-        assertThat(actualValue).isEqualTo(200);
-
-        // Verify changelog: before should be L2 (seq=8), after should be 
merged result
-        List<KeyValue> changelogs = result.changelogs();
-        assertThat(changelogs).hasSize(2);
-        assertThat(changelogs.get(0).valueKind()).isEqualTo(UPDATE_BEFORE);
-        assertThat(changelogs.get(0).value().getInt(1)).isEqualTo(8); // 
before is L2
-        assertThat(changelogs.get(1).valueKind()).isEqualTo(UPDATE_AFTER);
-    }
-
-    /**
-     * Test that without sequence.field, the original behavior is preserved: 
pick the record with
-     * the lowest level number.
-     */
-    @Test
-    public void testWithoutSequenceFieldPreservesOriginalBehavior() {
-        Map<InternalRow, KeyValue> highLevel = new HashMap<>();
-
-        // No userDefinedSeqComparator (null)
-        LookupChangelogMergeFunctionWrapper function =
-                new LookupChangelogMergeFunctionWrapper(
-                        LookupMergeFunction.wrap(
-                                DeduplicateMergeFunction.factory(), null, 
null, null),
-                        highLevel::get,
-                        null,
-                        LookupStrategy.from(false, true, false, false),
-                        null,
-                        null); // No sequence comparator
-
-        // L1: value=100, L2: value=200
-        // Without sequence.field, should pick L1 (level 1 < level 2)
-        function.reset();
-        function.add(new KeyValue().replace(row(1), 1, INSERT, 
row(100)).setLevel(1));
-        function.add(new KeyValue().replace(row(1), 1, INSERT, 
row(200)).setLevel(2));
-        function.add(new KeyValue().replace(row(1), 2, INSERT, 
row(50)).setLevel(0));
-
-        ChangelogResult result = function.getResult();
-        assertThat(result).isNotNull();
-
-        // Without sequence.field, L1 is picked as highLevel, and L0 is the 
latest
-        // So the result should be L0's value (50) since 
DeduplicateMergeFunction keeps the last
-        KeyValue kv = result.result();
-        assertThat(kv).isNotNull();
-        assertThat(kv.value().getInt(0)).isEqualTo(50);
-
-        // Changelog: before=L1(100), after=L0(50)
-        List<KeyValue> changelogs = result.changelogs();
-        assertThat(changelogs).hasSize(2);
-        assertThat(changelogs.get(0).value().getInt(0)).isEqualTo(100); // 
before is L1
-        assertThat(changelogs.get(1).value().getInt(0)).isEqualTo(50); // 
after is L0
-    }
-
-    /**
-     * Test sequence.field with descending sort order. When 
sort-order=descending, smaller sequence
-     * values are considered "newer".
-     *
-     * <p>Note: We use a 3-field schema with sequence at index 2 to avoid 
cache collision with other
-     * tests, because CodeGenUtils caches comparators by field types and 
indices without considering
-     * sort order.
-     */
-    @Test
-    public void testSequenceFieldWithDescendingSortOrder() {
-        RowType valueType =
-                RowType.builder()
-                        .fields(
-                                new DataType[] {DataTypes.INT(), 
DataTypes.INT(), DataTypes.INT()},
-                                new String[] {"value", "extra", "sequence"})
-                        .build();
-
-        // Create comparator with descending order
-        UserDefinedSeqComparator userDefinedSeqComparator =
-                UserDefinedSeqComparator.create(
-                        valueType,
-                        CoreOptions.fromMap(
-                                ImmutableMap.of(
-                                        "sequence.field", "sequence",
-                                        "sequence.field.sort-order", 
"descending")));
-        assertThat(userDefinedSeqComparator).isNotNull();
-
-        Map<InternalRow, KeyValue> highLevel = new HashMap<>();
-
-        LookupChangelogMergeFunctionWrapper function =
-                new LookupChangelogMergeFunctionWrapper(
-                        LookupMergeFunction.wrap(
-                                DeduplicateMergeFunction.factory(), null, 
null, null),
-                        highLevel::get,
-                        null,
-                        LookupStrategy.from(false, true, false, false),
-                        null,
-                        userDefinedSeqComparator);
-
-        // With descending order, smaller sequence = newer
-        // L1: (key=1, value=100, extra=0, sequence=7)  <- Level 1, newer (7 < 
8)
-        // L2: (key=1, value=200, extra=0, sequence=8)  <- Level 2, older (8 > 
7)
-        // L0: (key=1, value=50,  extra=0, sequence=9)  <- Level 0, oldest (9 
> 8 > 7)
-
-        function.reset();
-        function.add(new KeyValue().replace(row(1), 1, INSERT, row(100, 0, 
7)).setLevel(1));
-        function.add(new KeyValue().replace(row(1), 1, INSERT, row(200, 0, 
8)).setLevel(2));
-        function.add(new KeyValue().replace(row(1), 2, INSERT, row(50, 0, 
9)).setLevel(0));
-
-        ChangelogResult result = function.getResult();
-        assertThat(result).isNotNull();
-
-        KeyValue kv = result.result();
-        assertThat(kv).isNotNull();
-
-        int actualSequence = kv.value().getInt(2);
-        int actualValue = kv.value().getInt(0);
-
-        // With descending order, L1 (seq=7) is the newest high-level record
-        // The result should be L1's value (100) since it's the newest
-        assertThat(actualSequence)
-                .as("With descending order, should return record with smallest 
sequence (7)")
-                .isEqualTo(7);
-        assertThat(actualValue).isEqualTo(100);
-    }
 }
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/DeletionVectorITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/DeletionVectorITCase.java
index f42f83a28e..4221d05571 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/DeletionVectorITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/DeletionVectorITCase.java
@@ -342,6 +342,28 @@ public class DeletionVectorITCase extends 
CatalogITCaseBase {
                 .containsExactlyInAnyOrder(Row.of(1, 3, "1_2"), Row.of(2, 2, 
"2_1"));
     }
 
+    @ParameterizedTest
+    @ValueSource(booleans = {true, false})
+    public void 
testBatchReadDVTableWithOutOfOrderSequenceFieldAndAggregation(boolean 
dvBitmap64) {
+        sql(
+                String.format(
+                        "CREATE TABLE T (id INT PRIMARY KEY NOT ENFORCED, 
sequence INT, v INT) "
+                                + "WITH ("
+                                + "'deletion-vectors.enabled' = 'true', "
+                                + "'deletion-vectors.bitmap64' = '%s', "
+                                + "'changelog-producer' = 'none', "
+                                + "'sequence.field' = 'sequence', "
+                                + "'merge-engine' = 'aggregation', "
+                                + "'fields.v.aggregate-function' = 'sum')",
+                        dvBitmap64));
+
+        sql("INSERT INTO T /*+ OPTIONS('write-only' = 'true') */ VALUES (1, 7, 
7)");
+        sql("INSERT INTO T /*+ OPTIONS('write-only' = 'true') */ VALUES (1, 8, 
8)");
+        sql("INSERT INTO T /*+ OPTIONS('write-only' = 'false') */ VALUES (1, 
6, 6)");
+
+        assertThat(batchSql("SELECT * FROM T")).containsExactly(Row.of(1, 8, 
21));
+    }
+
     @ParameterizedTest
     @ValueSource(booleans = {true, false})
     public void testReadTagWithDv(boolean dvBitmap64) {
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupChangelogWithAggITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupChangelogWithAggITCase.java
index e288f8bcbf..289c857ae7 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupChangelogWithAggITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupChangelogWithAggITCase.java
@@ -104,6 +104,38 @@ public class LookupChangelogWithAggITCase extends 
CatalogITCaseBase {
                 .containsExactlyInAnyOrder(Row.of(1, 19), Row.of(2, 8), 
Row.of(3, 6));
     }
 
+    @Test
+    public void testLookupChangelogProducerWithOutOfOrderSequenceField() 
throws Exception {
+        sql(
+                "CREATE TABLE T (k INT PRIMARY KEY NOT ENFORCED, seq INT, v 
INT) WITH ("
+                        + "'bucket'='1', "
+                        + "'changelog-producer'='lookup', "
+                        + "'merge-engine'='aggregation', "
+                        + "'sequence.field'='seq', "
+                        + "'fields.v.aggregate-function'='sum', "
+                        + "'num-sorted-run.compaction-trigger'='2')");
+        BlockingIterator<Row, Row> iterator = streamSqlBlockIter("SELECT * 
FROM T");
+
+        sql("INSERT INTO T VALUES (1, 7, 7)");
+        assertThat(iterator.collect(1)).containsExactly(Row.of(1, 7, 7));
+
+        sql("INSERT INTO T VALUES (1, 8, 8)");
+        assertThat(iterator.collect(2))
+                .containsExactlyInAnyOrder(
+                        Row.ofKind(RowKind.UPDATE_BEFORE, 1, 7, 7),
+                        Row.ofKind(RowKind.UPDATE_AFTER, 1, 8, 15));
+
+        sql("INSERT INTO T VALUES (1, 6, 6)");
+        assertThat(iterator.collect(2))
+                .containsExactlyInAnyOrder(
+                        Row.ofKind(RowKind.UPDATE_BEFORE, 1, 8, 15),
+                        Row.ofKind(RowKind.UPDATE_AFTER, 1, 8, 21));
+
+        iterator.close();
+
+        assertThat(sql("SELECT * FROM T")).containsExactly(Row.of(1, 8, 21));
+    }
+
     @Test
     public void testLookupChangelogProducerWithProjection() {
         sql(

Reply via email to