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);
+        }
+    }
 }

Reply via email to