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 a5e87ebf6b [core] Optimize BTree negative predicate complement (#8945)
a5e87ebf6b is described below
commit a5e87ebf6b9b63cfa46980d05fff2edbe5c873f8
Author: zhoulii <[email protected]>
AuthorDate: Thu Aug 6 13:58:02 2026 +0800
[core] Optimize BTree negative predicate complement (#8945)
---
.../paimon/globalindex/GlobalIndexReader.java | 5 +
.../paimon/globalindex/GlobalIndexResult.java | 7 ++
.../globalindex/OffsetGlobalIndexReader.java | 39 ++++++++
.../globalindex/btree/LazyFilteredBTreeReader.java | 5 +
.../globalindex/GlobalIndexEvaluatorTest.java | 103 +++++++++++++++++++++
.../btree/LazyFilteredBTreeIndexReaderTest.java | 58 ++++++++++++
6 files changed, 217 insertions(+)
diff --git
a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexReader.java
b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexReader.java
index b857fdc1a1..694d367aaa 100644
---
a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexReader.java
+++
b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexReader.java
@@ -36,6 +36,11 @@ import java.util.concurrent.CompletableFuture;
public interface GlobalIndexReader
extends
FunctionVisitor<CompletableFuture<Optional<GlobalIndexResult>>>, Closeable {
+ /** Whether this reader can answer negative predicates by complementing a
known row range. */
+ default boolean supportsRangeComplement() {
+ return false;
+ }
+
@Override
default CompletableFuture<Optional<GlobalIndexResult>> visitIsNaN(FieldRef
fieldRef) {
return CompletableFuture.completedFuture(Optional.empty());
diff --git
a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexResult.java
b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexResult.java
index b92990a839..e9fe61a6e7 100644
---
a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexResult.java
+++
b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexResult.java
@@ -50,6 +50,13 @@ public interface GlobalIndexResult {
return create(RoaringNavigableMap64.or(this.results(),
other.results()));
}
+ default GlobalIndexResult andNot(GlobalIndexResult other) {
+ RoaringNavigableMap64 result = new RoaringNavigableMap64();
+ result.or(this.results());
+ result.andNot(other.results());
+ return create(result);
+ }
+
/** Returns an empty {@link GlobalIndexResult}. */
static GlobalIndexResult createEmpty() {
return create(new RoaringNavigableMap64());
diff --git
a/paimon-common/src/main/java/org/apache/paimon/globalindex/OffsetGlobalIndexReader.java
b/paimon-common/src/main/java/org/apache/paimon/globalindex/OffsetGlobalIndexReader.java
index 23d38d3a94..064893bf2e 100644
---
a/paimon-common/src/main/java/org/apache/paimon/globalindex/OffsetGlobalIndexReader.java
+++
b/paimon-common/src/main/java/org/apache/paimon/globalindex/OffsetGlobalIndexReader.java
@@ -23,6 +23,7 @@ import org.apache.paimon.predicate.FieldRef;
import org.apache.paimon.predicate.FullTextSearch;
import org.apache.paimon.predicate.TopN;
import org.apache.paimon.predicate.VectorSearch;
+import org.apache.paimon.utils.Range;
import java.io.IOException;
import java.util.ArrayList;
@@ -48,6 +49,9 @@ public class OffsetGlobalIndexReader implements
GlobalIndexReader {
@Override
public CompletableFuture<Optional<GlobalIndexResult>>
visitIsNotNull(FieldRef fieldRef) {
+ if (wrapped.supportsRangeComplement()) {
+ return complement(wrapped.visitIsNull(fieldRef));
+ }
return wrapped.visitIsNotNull(fieldRef).thenApply(this::applyOffset);
}
@@ -95,6 +99,12 @@ public class OffsetGlobalIndexReader implements
GlobalIndexReader {
@Override
public CompletableFuture<Optional<GlobalIndexResult>> visitNotEqual(
FieldRef fieldRef, Object literal) {
+ if (literal == null) {
+ return
CompletableFuture.completedFuture(Optional.of(GlobalIndexResult.createEmpty()));
+ }
+ if (wrapped.supportsRangeComplement()) {
+ return complement(wrapped.visitIsNull(fieldRef),
wrapped.visitEqual(fieldRef, literal));
+ }
return wrapped.visitNotEqual(fieldRef,
literal).thenApply(this::applyOffset);
}
@@ -125,6 +135,15 @@ public class OffsetGlobalIndexReader implements
GlobalIndexReader {
@Override
public CompletableFuture<Optional<GlobalIndexResult>> visitNotIn(
FieldRef fieldRef, List<Object> literals) {
+ for (Object literal : literals) {
+ if (literal == null) {
+ return CompletableFuture.completedFuture(
+ Optional.of(GlobalIndexResult.createEmpty()));
+ }
+ }
+ if (wrapped.supportsRangeComplement()) {
+ return complement(wrapped.visitIsNull(fieldRef),
wrapped.visitIn(fieldRef, literals));
+ }
return wrapped.visitNotIn(fieldRef,
literals).thenApply(this::applyOffset);
}
@@ -178,6 +197,26 @@ public class OffsetGlobalIndexReader implements
GlobalIndexReader {
return result.map(r -> r.offset(offset));
}
+ @SafeVarargs
+ private final CompletableFuture<Optional<GlobalIndexResult>> complement(
+ CompletableFuture<Optional<GlobalIndexResult>>... excludeFutures) {
+ return CompletableFuture.allOf(excludeFutures)
+ .thenApply(
+ ignored -> {
+ GlobalIndexResult result =
+ GlobalIndexResult.fromRange(new
Range(offset, to));
+ for
(CompletableFuture<Optional<GlobalIndexResult>> future :
+ excludeFutures) {
+ Optional<GlobalIndexResult> excluded =
future.join();
+ if (!excluded.isPresent()) {
+ return Optional.empty();
+ }
+ result =
result.andNot(excluded.get().offset(offset));
+ }
+ return Optional.of(result);
+ });
+ }
+
@Override
public void close() throws IOException {
wrapped.close();
diff --git
a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeReader.java
b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeReader.java
index d346d868e7..f97d7a0735 100644
---
a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeReader.java
+++
b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeReader.java
@@ -58,6 +58,11 @@ public class LazyFilteredBTreeReader extends
SortedFileGlobalIndexReader<BTreeIn
this.keySerializer = keySerializer;
}
+ @Override
+ public boolean supportsRangeComplement() {
+ return true;
+ }
+
@Override
protected Optional<GlobalIndexResult> visitIsNotNull(BTreeIndexReader
reader) {
return reader.visitIsNotNull();
diff --git
a/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java
b/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java
index 126ff3682c..a480ae35af 100644
---
a/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java
+++
b/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java
@@ -806,6 +806,109 @@ class GlobalIndexEvaluatorTest {
assertThat(reader.visitNotBetween(fieldRef, 1,
2).join()).contains(expected);
}
+ @Test
+ void testOffsetRangeComplementForNegativePredicates() {
+ FieldRef fieldRef = new FieldRef(0, "a", DataTypes.INT());
+ AtomicBoolean isNotNullVisited = new AtomicBoolean();
+ AtomicBoolean notEqualVisited = new AtomicBoolean();
+ AtomicBoolean notInVisited = new AtomicBoolean();
+ GlobalIndexReader delegate =
+ new StubGlobalIndexReader(null) {
+ @Override
+ public boolean supportsRangeComplement() {
+ return true;
+ }
+
+ @Override
+ public CompletableFuture<Optional<GlobalIndexResult>>
visitIsNull(
+ FieldRef fieldRef) {
+ return
CompletableFuture.completedFuture(Optional.of(resultOf(2, 4)));
+ }
+
+ @Override
+ public CompletableFuture<Optional<GlobalIndexResult>>
visitEqual(
+ FieldRef fieldRef, Object literal) {
+ return
CompletableFuture.completedFuture(Optional.of(resultOf(1, 3)));
+ }
+
+ @Override
+ public CompletableFuture<Optional<GlobalIndexResult>>
visitIn(
+ FieldRef fieldRef, List<Object> literals) {
+ return
CompletableFuture.completedFuture(Optional.of(resultOf(0, 5)));
+ }
+
+ @Override
+ public CompletableFuture<Optional<GlobalIndexResult>>
visitIsNotNull(
+ FieldRef fieldRef) {
+ isNotNullVisited.set(true);
+ return
CompletableFuture.completedFuture(Optional.of(resultOf(999)));
+ }
+
+ @Override
+ public CompletableFuture<Optional<GlobalIndexResult>>
visitNotEqual(
+ FieldRef fieldRef, Object literal) {
+ notEqualVisited.set(true);
+ return
CompletableFuture.completedFuture(Optional.of(resultOf(999)));
+ }
+
+ @Override
+ public CompletableFuture<Optional<GlobalIndexResult>>
visitNotIn(
+ FieldRef fieldRef, List<Object> literals) {
+ notInVisited.set(true);
+ return
CompletableFuture.completedFuture(Optional.of(resultOf(999)));
+ }
+ };
+
+ GlobalIndexReader reader = new OffsetGlobalIndexReader(delegate, 10L,
15L);
+
+ assertBitmapContainsExactly(
+ reader.visitIsNotNull(fieldRef).join().get().results(), 10L,
11L, 13L, 15L);
+ assertBitmapContainsExactly(
+ reader.visitNotEqual(fieldRef, 5).join().get().results(), 10L,
15L);
+ assertBitmapContainsExactly(
+ reader.visitNotIn(fieldRef, Arrays.asList(5,
6)).join().get().results(), 11L, 13L);
+ assertThat(isNotNullVisited).isFalse();
+ assertThat(notEqualVisited).isFalse();
+ assertThat(notInVisited).isFalse();
+ }
+
+ @Test
+ void testOffsetRangeComplementNullAndUnsupportedPredicates() {
+ FieldRef fieldRef = new FieldRef(0, "a", DataTypes.INT());
+ GlobalIndexReader unsupportedEqual =
+ new StubGlobalIndexReader(null) {
+ @Override
+ public boolean supportsRangeComplement() {
+ return true;
+ }
+
+ @Override
+ public CompletableFuture<Optional<GlobalIndexResult>>
visitIsNull(
+ FieldRef fieldRef) {
+ return CompletableFuture.completedFuture(
+ Optional.of(GlobalIndexResult.createEmpty()));
+ }
+
+ @Override
+ public CompletableFuture<Optional<GlobalIndexResult>>
visitEqual(
+ FieldRef fieldRef, Object literal) {
+ return
CompletableFuture.completedFuture(Optional.empty());
+ }
+ };
+
+ GlobalIndexReader reader = new
OffsetGlobalIndexReader(unsupportedEqual, 10L, 15L);
+
+ assertThat(reader.visitNotEqual(fieldRef, 5).join()).isEmpty();
+ assertThat(reader.visitNotEqual(fieldRef,
null).join().get().results().isEmpty()).isTrue();
+ assertThat(
+ reader.visitNotIn(fieldRef, Arrays.asList(5, null))
+ .join()
+ .get()
+ .results()
+ .isEmpty())
+ .isTrue();
+ }
+
@Test
void testNullPredicate() {
RowType rowType = rowType();
diff --git
a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeIndexReaderTest.java
b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeIndexReaderTest.java
index d8321d7684..1fd4c509e9 100644
---
a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeIndexReaderTest.java
+++
b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeIndexReaderTest.java
@@ -19,18 +19,22 @@
package org.apache.paimon.globalindex.btree;
import org.apache.paimon.fs.Path;
+import org.apache.paimon.fs.SeekableInputStream;
import org.apache.paimon.globalindex.GlobalIndexIOMeta;
import org.apache.paimon.globalindex.GlobalIndexReader;
import org.apache.paimon.globalindex.GlobalIndexResult;
import org.apache.paimon.globalindex.GlobalIndexSingleColumnWriter;
+import org.apache.paimon.globalindex.OffsetGlobalIndexReader;
import org.apache.paimon.globalindex.ResultEntry;
import org.apache.paimon.globalindex.btree.BTreeIndexReader.KeyRowIds;
+import org.apache.paimon.globalindex.io.GlobalIndexFileReader;
import org.apache.paimon.options.MemorySize;
import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.FieldRef;
import org.apache.paimon.predicate.TopN;
import
org.apache.paimon.testutils.junit.parameterized.ParameterizedTestExtension;
import org.apache.paimon.types.DataField;
+import org.apache.paimon.types.IntType;
import org.apache.paimon.utils.Pair;
import org.apache.paimon.utils.SemaphoredDelegatingExecutor;
@@ -41,6 +45,7 @@ import org.junit.jupiter.api.extension.ExtendWith;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Comparator;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -179,6 +184,34 @@ public class LazyFilteredBTreeIndexReaderTest extends
AbstractIndexReaderTest {
}
}
+ @TestTemplate
+ public void testOffsetRangeComplementAvoidsOpeningAllBTreeFiles() throws
Exception {
+ if (!(dataType instanceof IntType)) {
+ // The test uses Integer literals; file pruning itself is
independent of the key type.
+ return;
+ }
+
+ List<GlobalIndexIOMeta> written = new ArrayList<>();
+ written.add(writeData(singletonData(1, 0L)));
+ written.add(writeData(singletonData(100, 1L)));
+ written.add(writeData(singletonData(200, 2L)));
+ written.add(writeData(singletonData(null, 3L)));
+
+ CountingGlobalIndexFileReader countingReader = new
CountingGlobalIndexFileReader();
+ FieldRef ref = new FieldRef(1, "testField", dataType);
+ try (GlobalIndexReader reader =
+ new OffsetGlobalIndexReader(
+ globalIndexer.createReader(
+ countingReader, written,
newDirectExecutorService()),
+ 1000L,
+ 1003L)) {
+ GlobalIndexResult result = reader.visitNotEqual(ref,
100).join().get();
+
+ assertRows(result, 1000L, 1002L);
+ assertThat(countingReader.openedFiles).hasSize(2);
+ }
+ }
+
@TestTemplate
public void testUnorderedIterator() throws Exception {
// Set some null values
@@ -342,6 +375,20 @@ public class LazyFilteredBTreeIndexReaderTest extends
AbstractIndexReaderTest {
resultEntry.meta());
}
+ private List<Pair<Object, Long>> singletonData(Object key, long rowId) {
+ List<Pair<Object, Long>> result = new ArrayList<>();
+ result.add(Pair.of(key, rowId));
+ return result;
+ }
+
+ private void assertRows(GlobalIndexResult indexResult, Long... expected) {
+ List<Long> actual = new ArrayList<>();
+ for (Long rowId : indexResult.results()) {
+ actual.add(rowId);
+ }
+ assertThat(actual).containsExactlyInAnyOrder(expected);
+ }
+
/**
* Regression test for deadlock when using {@link
SemaphoredDelegatingExecutor}. Before the fix,
* {@code LazyFilteredBTreeReader.visitParallel} submitted tasks to the
executor (acquiring
@@ -514,4 +561,15 @@ public class LazyFilteredBTreeIndexReaderTest extends
AbstractIndexReaderTest {
break;
}
}
+
+ private class CountingGlobalIndexFileReader implements
GlobalIndexFileReader {
+
+ private final Set<Path> openedFiles = new HashSet<>();
+
+ @Override
+ public SeekableInputStream getInputStream(GlobalIndexIOMeta meta)
throws IOException {
+ openedFiles.add(meta.filePath());
+ return fileReader.getInputStream(meta);
+ }
+ }
}