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(