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 b9915d4922 [core] Read source-backed indexes in data evolution tables
(#8794)
b9915d4922 is described below
commit b9915d49222f9633025e60ff2f7da9866c47d2f9
Author: Jingsong Lee <[email protected]>
AuthorDate: Wed Jul 22 14:39:35 2026 +0800
[core] Read source-backed indexes in data evolution tables (#8794)
- Rename the Java and Python GlobalIndexScanner components to
DataEvolutionGlobalIndexScanner.
- Rename the Java and Python GlobalIndexCoverage components to
DataEvolutionGlobalIndexCoverage.
- Remove primary-key and sourceMeta filtering from the Data Evolution
scanner and coverage calculation.
- Let Data Evolution full-text and vector paths consume the renamed
components.
- Add regression coverage for source-backed index selection, mixed index
coverage, and full-text reads.
---
.../paimon/globalindex/DataEvolutionBatchScan.java | 6 +-
....java => DataEvolutionGlobalIndexCoverage.java} | 9 +--
...r.java => DataEvolutionGlobalIndexScanner.java} | 41 ++++++-------
.../paimon/table/source/AbstractVectorRead.java | 14 ++---
.../table/source/DataEvolutionFullTextScan.java | 7 ++-
.../table/source/DataEvolutionVectorScan.java | 7 ++-
.../paimon/table/BitmapGlobalIndexTableTest.java | 8 ++-
.../paimon/table/BtreeGlobalIndexTableTest.java | 71 ++++++++++++----------
.../table/source/FullTextSearchBuilderTest.java | 13 ++--
.../table/source/VectorSearchBuilderTest.java | 3 +-
paimon-python/pypaimon/globalindex/__init__.py | 6 +-
....py => data_evolution_global_index_coverage.py} | 8 +--
...r.py => data_evolution_global_index_scanner.py} | 18 +++---
.../pypaimon/read/scanner/file_scanner.py | 4 +-
.../pypaimon/table/source/full_text_scan.py | 4 +-
.../table/source/primary_key_sorted_index_scan.py | 2 +-
.../pypaimon/table/source/vector_search_read.py | 8 +--
.../pypaimon/table/source/vector_search_scan.py | 6 +-
.../pypaimon/tests/e2e/java_py_read_write_test.py | 8 +--
.../pypaimon/tests/global_index_build_test.py | 6 +-
paimon-python/pypaimon/tests/global_index_test.py | 36 +++++------
.../pypaimon/tests/vector_search_filter_test.py | 66 ++++++++++----------
22 files changed, 181 insertions(+), 170 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java
index ebdfb0894f..aacae5ce68 100644
---
a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java
+++
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java
@@ -295,14 +295,14 @@ public class DataEvolutionBatchScan implements
DataTableScan {
PartitionPredicate partitionFilter =
batchScan.snapshotReader().manifestsReader().partitionFilter();
long totalStart = System.nanoTime();
- Optional<GlobalIndexScanner> optionalScanner =
- GlobalIndexScanner.create(table, partitionFilter,
globalIndexFilter);
+ Optional<DataEvolutionGlobalIndexScanner> optionalScanner =
+ DataEvolutionGlobalIndexScanner.create(table, partitionFilter,
globalIndexFilter);
long metadataDuration = System.nanoTime() - totalStart;
if (!optionalScanner.isPresent()) {
return Optional.empty();
}
- try (GlobalIndexScanner scanner = optionalScanner.get()) {
+ try (DataEvolutionGlobalIndexScanner scanner = optionalScanner.get()) {
long lookupStart = System.nanoTime();
Optional<GlobalIndexResult> result =
scanner.scan(globalIndexFilter);
long lookupDuration = System.nanoTime() - lookupStart;
diff --git
a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexCoverage.java
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexCoverage.java
similarity index 96%
rename from
paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexCoverage.java
rename to
paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexCoverage.java
index 6353b6616e..55678824e6 100644
---
a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexCoverage.java
+++
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexCoverage.java
@@ -45,15 +45,15 @@ import java.util.Map;
import static org.apache.paimon.predicate.PredicateVisitor.collectFieldIds;
import static org.apache.paimon.utils.Preconditions.checkNotNull;
-/** Row ranges covered and not covered by global index files. */
-public class GlobalIndexCoverage {
+/** Row ranges covered and not covered by global index files on data-evolution
tables. */
+public class DataEvolutionGlobalIndexCoverage {
private final FileStoreTable table;
@Nullable private final Snapshot snapshot;
@Nullable private final PartitionPredicate partitionFilter;
private final Map<Integer, List<Range>> coverageByField;
- public GlobalIndexCoverage(
+ public DataEvolutionGlobalIndexCoverage(
FileStoreTable table,
@Nullable Snapshot snapshot,
@Nullable PartitionPredicate partitionFilter,
@@ -64,9 +64,6 @@ public class GlobalIndexCoverage {
this.coverageByField = new HashMap<>();
for (IndexFileMeta indexFile : indexFiles) {
GlobalIndexMeta meta = checkNotNull(indexFile.globalIndexMeta());
- if (meta.sourceMeta() != null) {
- continue;
- }
Range range = new Range(meta.rowRangeStart(), meta.rowRangeEnd());
addCoverage(meta.indexFieldId(), range);
if (meta.extraFieldIds() != null) {
diff --git
a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexScanner.java
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java
similarity index 92%
rename from
paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexScanner.java
rename to
paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java
index df9a306120..5f07e8a17e 100644
---
a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexScanner.java
+++
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java
@@ -63,20 +63,21 @@ import static
org.apache.paimon.table.source.snapshot.TimeTravelUtil.tryTravelOr
import static org.apache.paimon.utils.Preconditions.checkArgument;
import static org.apache.paimon.utils.Preconditions.checkNotNull;
-/** Scanner for shard-based global indexes. */
-public class GlobalIndexScanner implements Closeable {
+/** Scanner for shard-based global indexes on data-evolution tables. */
+public class DataEvolutionGlobalIndexScanner implements Closeable {
- private static final Logger LOG =
LoggerFactory.getLogger(GlobalIndexScanner.class);
+ private static final Logger LOG =
+ LoggerFactory.getLogger(DataEvolutionGlobalIndexScanner.class);
private final Options options;
private final RowType rowType;
private final ExecutorService executor;
private final GlobalIndexEvaluator globalIndexEvaluator;
private final IndexPathFactory indexPathFactory;
- private final GlobalIndexCoverage coverage;
+ private final DataEvolutionGlobalIndexCoverage coverage;
private final FileStoreTable table;
- private GlobalIndexScanner(
+ private DataEvolutionGlobalIndexScanner(
FileStoreTable table,
@Nullable Snapshot snapshot,
@Nullable PartitionPredicate partitionFilter,
@@ -91,7 +92,8 @@ public class GlobalIndexScanner implements Closeable {
this.executor =
GlobalIndexReadThreadPool.getExecutorService(options.get(GLOBAL_INDEX_THREAD_NUM));
this.indexPathFactory = indexPathFactory;
- this.coverage = new GlobalIndexCoverage(table, snapshot,
partitionFilter, indexFiles);
+ this.coverage =
+ new DataEvolutionGlobalIndexCoverage(table, snapshot,
partitionFilter, indexFiles);
GlobalIndexFileReader indexFileReader = meta ->
fileIO.newInputStream(meta.filePath());
Map<Integer, IndexMetaFileGroup> indexMetas = new HashMap<>();
Map<Integer, List<IndexMetaFileGroup>> extraIndexMetas = new
HashMap<>();
@@ -186,21 +188,21 @@ public class GlobalIndexScanner implements Closeable {
}
}
- public static Optional<GlobalIndexScanner> create(
+ public static Optional<DataEvolutionGlobalIndexScanner> create(
FileStoreTable table, Collection<IndexFileMeta> indexFiles) {
return create(table, null, indexFiles);
}
- public static Optional<GlobalIndexScanner> create(
+ public static Optional<DataEvolutionGlobalIndexScanner> create(
FileStoreTable table,
@Nullable PartitionPredicate partitionFilter,
Collection<IndexFileMeta> indexFiles) {
- List<IndexFileMeta> ordinaryIndexFiles =
ordinaryGlobalIndexFiles(indexFiles);
- if (ordinaryIndexFiles.isEmpty()) {
+ List<IndexFileMeta> globalIndexFiles = globalIndexFiles(indexFiles);
+ if (globalIndexFiles.isEmpty()) {
return Optional.empty();
}
return Optional.of(
- new GlobalIndexScanner(
+ new DataEvolutionGlobalIndexScanner(
table,
tryTravelOrLatest(table),
partitionFilter,
@@ -208,10 +210,10 @@ public class GlobalIndexScanner implements Closeable {
table.rowType(),
table.fileIO(),
table.store().pathFactory().globalIndexFileFactory(),
- ordinaryIndexFiles));
+ globalIndexFiles));
}
- public static Optional<GlobalIndexScanner> create(
+ public static Optional<DataEvolutionGlobalIndexScanner> create(
FileStoreTable table,
@Nullable PartitionPredicate partitionFilter,
@Nullable Predicate filter) {
@@ -225,7 +227,7 @@ public class GlobalIndexScanner implements Closeable {
return Optional.empty();
}
return Optional.of(
- new GlobalIndexScanner(
+ new DataEvolutionGlobalIndexScanner(
table,
snapshot,
partitionFilter,
@@ -250,7 +252,7 @@ public class GlobalIndexScanner implements Closeable {
return false;
}
GlobalIndexMeta globalIndex =
entry.indexFile().globalIndexMeta();
- if (globalIndex == null || globalIndex.sourceMeta() !=
null) {
+ if (globalIndex == null) {
return false;
}
// Collect indexes whose primary column is filtered, and
also multi-column
@@ -270,14 +272,9 @@ public class GlobalIndexScanner implements Closeable {
return indexFileFilter;
}
- private static List<IndexFileMeta> ordinaryGlobalIndexFiles(
- Collection<IndexFileMeta> indexFiles) {
+ private static List<IndexFileMeta>
globalIndexFiles(Collection<IndexFileMeta> indexFiles) {
return indexFiles.stream()
- .filter(
- indexFile -> {
- GlobalIndexMeta meta = indexFile.globalIndexMeta();
- return meta != null && meta.sourceMeta() == null;
- })
+ .filter(indexFile -> indexFile.globalIndexMeta() != null)
.collect(Collectors.toList());
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractVectorRead.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractVectorRead.java
index 931d78389a..326f814428 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractVectorRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractVectorRead.java
@@ -22,10 +22,10 @@ import org.apache.paimon.data.InternalArray;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.InternalVector;
import org.apache.paimon.fs.FileIO;
+import org.apache.paimon.globalindex.DataEvolutionGlobalIndexScanner;
import org.apache.paimon.globalindex.GlobalIndexIOMeta;
import org.apache.paimon.globalindex.GlobalIndexReader;
import org.apache.paimon.globalindex.GlobalIndexResult;
-import org.apache.paimon.globalindex.GlobalIndexScanner;
import org.apache.paimon.globalindex.GlobalIndexer;
import org.apache.paimon.globalindex.GlobalIndexerFactoryUtils;
import org.apache.paimon.globalindex.OffsetGlobalIndexReader;
@@ -172,13 +172,13 @@ public abstract class AbstractVectorRead implements
Serializable {
scalarIndexFiles.addAll(split.scalarIndexFiles());
}
- Optional<GlobalIndexScanner> optionalScanner =
- GlobalIndexScanner.create(table, partitionFilter,
scalarIndexFiles);
+ Optional<DataEvolutionGlobalIndexScanner> optionalScanner =
+ DataEvolutionGlobalIndexScanner.create(table, partitionFilter,
scalarIndexFiles);
if (!optionalScanner.isPresent()) {
return new RoaringNavigableMap64();
}
- try (GlobalIndexScanner scanner = optionalScanner.get()) {
+ try (DataEvolutionGlobalIndexScanner scanner = optionalScanner.get()) {
Optional<GlobalIndexResult> result = scanner.scan(filter);
if (!result.isPresent()) {
return new RoaringNavigableMap64();
@@ -205,14 +205,14 @@ public abstract class AbstractVectorRead implements
Serializable {
for (RawVectorSearchSplit split : splits) {
scalarIndexFiles.addAll(split.scalarIndexFiles());
}
- Optional<GlobalIndexScanner> optionalScanner =
- GlobalIndexScanner.create(table, partitionFilter,
scalarIndexFiles);
+ Optional<DataEvolutionGlobalIndexScanner> optionalScanner =
+ DataEvolutionGlobalIndexScanner.create(table, partitionFilter,
scalarIndexFiles);
if (!optionalScanner.isPresent()) {
return null;
}
RoaringNavigableMap64 include = new RoaringNavigableMap64();
- try (GlobalIndexScanner scanner = optionalScanner.get()) {
+ try (DataEvolutionGlobalIndexScanner scanner = optionalScanner.get()) {
Optional<GlobalIndexResult> result = scanner.scan(filter);
if (!result.isPresent()) {
return null;
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionFullTextScan.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionFullTextScan.java
index bb250dbe60..87ef271d4b 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionFullTextScan.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionFullTextScan.java
@@ -20,7 +20,7 @@ package org.apache.paimon.table.source;
import org.apache.paimon.Snapshot;
import org.apache.paimon.annotation.VisibleForTesting;
-import org.apache.paimon.globalindex.GlobalIndexCoverage;
+import org.apache.paimon.globalindex.DataEvolutionGlobalIndexCoverage;
import org.apache.paimon.globalindex.GlobalIndexerFactory;
import org.apache.paimon.globalindex.GlobalIndexerFactoryUtils;
import org.apache.paimon.index.GlobalIndexMeta;
@@ -94,7 +94,7 @@ public class DataEvolutionFullTextScan implements
FullTextScan {
return false;
}
GlobalIndexMeta globalIndex =
entry.indexFile().globalIndexMeta();
- if (globalIndex == null || globalIndex.sourceMeta() !=
null) {
+ if (globalIndex == null) {
return false;
}
return !matchedTextColumnIds(globalIndex,
textColumnIds).isEmpty()
@@ -120,7 +120,8 @@ public class DataEvolutionFullTextScan implements
FullTextScan {
if (!allIndexFiles.isEmpty()) {
List<Range> rawRowRanges =
- new GlobalIndexCoverage(table, snapshot, partitionFilter,
allIndexFiles)
+ new DataEvolutionGlobalIndexCoverage(
+ table, snapshot, partitionFilter,
allIndexFiles)
.unindexedRanges(textColumnIds);
if (!rawRowRanges.isEmpty()) {
splits.add(new RawFullTextSearchSplit(rawRowRanges));
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorScan.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorScan.java
index 6f123632c1..317ffc6d52 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorScan.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorScan.java
@@ -19,7 +19,7 @@
package org.apache.paimon.table.source;
import org.apache.paimon.Snapshot;
-import org.apache.paimon.globalindex.GlobalIndexCoverage;
+import org.apache.paimon.globalindex.DataEvolutionGlobalIndexCoverage;
import org.apache.paimon.index.GlobalIndexMeta;
import org.apache.paimon.index.IndexFileHandler;
import org.apache.paimon.index.IndexFileMeta;
@@ -145,14 +145,15 @@ public class DataEvolutionVectorScan implements
VectorScan {
}
List<Range> rawRowRanges =
- new GlobalIndexCoverage(table, snapshot, partitionFilter,
vectorIndexFiles)
+ new DataEvolutionGlobalIndexCoverage(
+ table, snapshot, partitionFilter,
vectorIndexFiles)
.unindexedRanges(vectorColumn.id());
if (filter != null) {
rawRowRanges =
Range.sortAndMergeOverlap(
addAll(
rawRowRanges,
- new GlobalIndexCoverage(
+ new DataEvolutionGlobalIndexCoverage(
table,
snapshot,
partitionFilter,
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/BitmapGlobalIndexTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/BitmapGlobalIndexTableTest.java
index dc276b8d98..fb8f980dcc 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/BitmapGlobalIndexTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/BitmapGlobalIndexTableTest.java
@@ -22,8 +22,8 @@ import org.apache.paimon.data.BinaryRow;
import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.globalindex.DataEvolutionGlobalIndexScanner;
import org.apache.paimon.globalindex.GlobalIndexBuilderUtils;
-import org.apache.paimon.globalindex.GlobalIndexScanner;
import org.apache.paimon.globalindex.GlobalIndexSingleColumnWriter;
import org.apache.paimon.globalindex.GlobalIndexWriter;
import org.apache.paimon.globalindex.IndexedSplit;
@@ -293,8 +293,10 @@ public class BitmapGlobalIndexTableTest extends
DataEvolutionTestBase {
private RoaringNavigableMap64 globalIndexScan(FileStoreTable table,
Predicate predicate)
throws Exception {
- try (GlobalIndexScanner scanner =
- GlobalIndexScanner.create(table,
PartitionPredicate.ALWAYS_TRUE, predicate).get()) {
+ try (DataEvolutionGlobalIndexScanner scanner =
+ DataEvolutionGlobalIndexScanner.create(
+ table, PartitionPredicate.ALWAYS_TRUE,
predicate)
+ .get()) {
return scanner.scan(predicate).get().results();
}
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java
index 4e306768c1..c0d8ea01aa 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java
@@ -23,9 +23,9 @@ import org.apache.paimon.Snapshot;
import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.GenericRow;
import org.apache.paimon.globalindex.DataEvolutionBatchScan;
-import org.apache.paimon.globalindex.GlobalIndexCoverage;
+import org.apache.paimon.globalindex.DataEvolutionGlobalIndexCoverage;
+import org.apache.paimon.globalindex.DataEvolutionGlobalIndexScanner;
import org.apache.paimon.globalindex.GlobalIndexResult;
-import org.apache.paimon.globalindex.GlobalIndexScanner;
import org.apache.paimon.globalindex.IndexedSplit;
import org.apache.paimon.globalindex.btree.BTreeIndexOptions;
import org.apache.paimon.globalindex.sorted.SortedGlobalIndexBuilder;
@@ -181,7 +181,7 @@ public class BtreeGlobalIndexTableTest extends
DataEvolutionTestBase {
.setLayout(PatternLayout.newBuilder().withPattern("%level %msg%n").build())
.build();
Logger scanLogger = (Logger)
LogManager.getLogger(DataEvolutionBatchScan.class);
- Logger scannerLogger = (Logger)
LogManager.getLogger(GlobalIndexScanner.class);
+ Logger scannerLogger = (Logger)
LogManager.getLogger(DataEvolutionGlobalIndexScanner.class);
Level previousScanLevel = scanLogger.getLevel();
Level previousScannerLevel = scannerLogger.getLevel();
@@ -259,7 +259,7 @@ public class BtreeGlobalIndexTableTest extends
DataEvolutionTestBase {
}
@Test
- public void testGlobalIndexScannerKeepsUnindexedRowsSeparate() throws
Exception {
+ public void
testDataEvolutionGlobalIndexScannerKeepsUnindexedRowsSeparate() throws
Exception {
write(500L);
createIndex("f1");
appendRows(500, 1000);
@@ -274,8 +274,10 @@ public class BtreeGlobalIndexTableTest extends
DataEvolutionTestBase {
BinaryString.fromString("a100"),
BinaryString.fromString("a700")));
- try (GlobalIndexScanner scanner =
- GlobalIndexScanner.create(table,
PartitionPredicate.ALWAYS_TRUE, predicate).get()) {
+ try (DataEvolutionGlobalIndexScanner scanner =
+ DataEvolutionGlobalIndexScanner.create(
+ table, PartitionPredicate.ALWAYS_TRUE,
predicate)
+ .get()) {
assertThat(scanner.scan(predicate).get().results().toRangeList())
.containsExactly(new Range(100L, 100L));
assertThat(scanner.unindexedRows(predicate).results().toRangeList())
@@ -284,10 +286,11 @@ public class BtreeGlobalIndexTableTest extends
DataEvolutionTestBase {
}
@Test
- public void testSourceBackedIndexIsExcludedFromGlobalRowIdScan() throws
Exception {
+ public void
testDataEvolutionSourceBackedIndexParticipatesInGlobalRowIdScan() throws
Exception {
write(10L);
FileStoreTable table =
tableWithSearchMode((FileStoreTable)
catalog.getTable(identifier()), "full");
+ Snapshot snapshot = table.snapshotManager().latestSnapshot();
IndexFileMeta sourceBacked =
new IndexFileMeta(
"btree",
@@ -297,44 +300,48 @@ public class BtreeGlobalIndexTableTest extends
DataEvolutionTestBase {
new GlobalIndexMeta(0, 9, 1, null, null, new byte[]
{1}),
null);
- assertThat(GlobalIndexScanner.create(table,
Collections.singletonList(sourceBacked)))
- .isEmpty();
+ assertThat(
+ DataEvolutionGlobalIndexScanner.create(
+ table,
Collections.singletonList(sourceBacked)))
+ .isPresent();
- GlobalIndexCoverage coverage =
- new GlobalIndexCoverage(
+ DataEvolutionGlobalIndexCoverage coverage =
+ new DataEvolutionGlobalIndexCoverage(
table,
- table.snapshotManager().latestSnapshot(),
+ snapshot,
PartitionPredicate.ALWAYS_TRUE,
Collections.singletonList(sourceBacked));
- assertThat(coverage.unindexedRanges(1)).containsExactly(new Range(0,
9));
+ assertThat(coverage.unindexedRanges(1)).isEmpty();
}
@Test
- public void testOrdinaryAndSourceBackedBTreeIndexesCanCoexist() throws
Exception {
+ public void testOrdinaryAndSourceBackedBTreeIndexCoverageCanCoexist()
throws Exception {
write(10L);
- createIndex("f1");
- FileStoreTable table = (FileStoreTable) catalog.getTable(identifier());
+ FileStoreTable table =
+ tableWithSearchMode((FileStoreTable)
catalog.getTable(identifier()), "full");
Snapshot snapshot = table.snapshotManager().latestSnapshot();
- List<IndexFileMeta> mixedIndexes =
- table.store().newIndexFileHandler().scan(snapshot,
"btree").stream()
- .map(IndexManifestEntry::indexFile)
- .collect(Collectors.toCollection(ArrayList::new));
+ List<IndexFileMeta> mixedIndexes = new ArrayList<>();
+ mixedIndexes.add(
+ new IndexFileMeta(
+ "btree",
+ "ordinary-index",
+ 0,
+ 5,
+ new GlobalIndexMeta(0, 4, 1, null, null),
+ null));
mixedIndexes.add(
new IndexFileMeta(
"btree",
"source-backed-index",
0,
- 10,
- new GlobalIndexMeta(0, 9, 1, null, null, new byte[]
{1}),
+ 5,
+ new GlobalIndexMeta(5, 9, 1, null, null, new byte[]
{1}),
null));
- Predicate predicate =
- new PredicateBuilder(table.rowType()).equal(1,
BinaryString.fromString("a7"));
- try (GlobalIndexScanner scanner =
- GlobalIndexScanner.create(table,
mixedIndexes).orElseThrow(AssertionError::new)) {
- assertThat(scanner.scan(predicate).get().results().toRangeList())
- .containsExactly(new Range(7, 7));
- }
+ DataEvolutionGlobalIndexCoverage coverage =
+ new DataEvolutionGlobalIndexCoverage(
+ table, snapshot, PartitionPredicate.ALWAYS_TRUE,
mixedIndexes);
+ assertThat(coverage.unindexedRanges(1)).isEmpty();
}
@Test
@@ -577,8 +584,10 @@ public class BtreeGlobalIndexTableTest extends
DataEvolutionTestBase {
private RoaringNavigableMap64 globalIndexScan(FileStoreTable table,
Predicate predicate)
throws Exception {
- try (GlobalIndexScanner scanner =
- GlobalIndexScanner.create(table,
PartitionPredicate.ALWAYS_TRUE, predicate).get()) {
+ try (DataEvolutionGlobalIndexScanner scanner =
+ DataEvolutionGlobalIndexScanner.create(
+ table, PartitionPredicate.ALWAYS_TRUE,
predicate)
+ .get()) {
return scanner.scan(predicate).get().results();
}
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/source/FullTextSearchBuilderTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/source/FullTextSearchBuilderTest.java
index 5c43f43470..12af712f1d 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/source/FullTextSearchBuilderTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/source/FullTextSearchBuilderTest.java
@@ -747,7 +747,7 @@ public class FullTextSearchBuilderTest extends
TableTestBase {
}
@Test
- public void testOrdinaryFullTextScanSkipsSourceBackedPrimaryKeyArchive()
throws Exception {
+ public void testDataEvolutionFullTextScanReadsSourceBackedIndex() throws
Exception {
createTableDefault();
FileStoreTable table = getTableDefault();
@@ -755,14 +755,15 @@ public class FullTextSearchBuilderTest extends
TableTestBase {
writeDocuments(table, documents);
buildAndCommitSourceBackedIndex(table, documents);
- FullTextScan.Plan plan =
+ FullTextSearchBuilder searchBuilder =
table.newFullTextSearchBuilder()
.withQuery(TEXT_FIELD_NAME, matchQuery("Paimon"))
- .withLimit(2)
- .newFullTextScan()
- .scan();
+ .withLimit(2);
+ FullTextScan.Plan plan = searchBuilder.newFullTextScan().scan();
- assertThat(plan.splits()).isEmpty();
+ assertThat(plan.splits()).hasSize(1);
+ assertThat(searchBuilder.executeLocal().results().toRangeList())
+ .containsExactly(new Range(0, 0));
}
@Test
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java
index f704e1f370..63c2d1834d 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java
@@ -1267,7 +1267,8 @@ public class VectorSearchBuilderTest extends
TableTestBase {
buildAndCommitVectorIndexWithFields(
table, vectors, new Range(0, 3), Arrays.asList(vectorField,
idField));
// The dedicated scalar index covers only the head. The vector index's
scalar extra field
- // must still serve the tail; otherwise GlobalIndexScanner's
primary-field preference drops
+ // must still serve the tail; otherwise
DataEvolutionGlobalIndexScanner's primary-field
+ // preference drops
// rows [2, 3] even though coverage planning treats them as indexed.
buildAndCommitBTreeIndex(table, new int[] {0, 1}, new Range(0, 1));
diff --git a/paimon-python/pypaimon/globalindex/__init__.py
b/paimon-python/pypaimon/globalindex/__init__.py
index 842b06b9e8..b9ef7c9a28 100644
--- a/paimon-python/pypaimon/globalindex/__init__.py
+++ b/paimon-python/pypaimon/globalindex/__init__.py
@@ -26,7 +26,9 @@ from pypaimon.globalindex.vector_search_result import (
)
from pypaimon.globalindex.global_index_meta import GlobalIndexMeta,
GlobalIndexIOMeta
from pypaimon.globalindex.global_index_evaluator import GlobalIndexEvaluator
-from pypaimon.globalindex.global_index_scanner import GlobalIndexScanner
+from pypaimon.globalindex.data_evolution_global_index_scanner import (
+ DataEvolutionGlobalIndexScanner,
+)
from pypaimon.globalindex.key_serializer import KeySerializer
from pypaimon.globalindex.memory_slice_input import MemorySliceInput
from pypaimon.globalindex.offset_global_index_reader import
OffsetGlobalIndexReader
@@ -55,7 +57,7 @@ __all__ = [
'GlobalIndexMeta',
'GlobalIndexIOMeta',
'GlobalIndexEvaluator',
- 'GlobalIndexScanner',
+ 'DataEvolutionGlobalIndexScanner',
'KeySerializer',
'MemorySliceInput',
'OffsetGlobalIndexReader',
diff --git a/paimon-python/pypaimon/globalindex/global_index_coverage.py
b/paimon-python/pypaimon/globalindex/data_evolution_global_index_coverage.py
similarity index 95%
rename from paimon-python/pypaimon/globalindex/global_index_coverage.py
rename to
paimon-python/pypaimon/globalindex/data_evolution_global_index_coverage.py
index d0a05cb50b..07a9213bb8 100644
--- a/paimon-python/pypaimon/globalindex/global_index_coverage.py
+++ b/paimon-python/pypaimon/globalindex/data_evolution_global_index_coverage.py
@@ -15,7 +15,7 @@
# specific language governing permissions and limitations
# under the License.
-"""Row ranges covered and not covered by global index files."""
+"""Global-index coverage for data-evolution tables."""
from typing import Collection, Dict, List, Optional, Union
@@ -27,7 +27,7 @@ from pypaimon.schema.data_types import DataField
from pypaimon.utils.range import Range
-class GlobalIndexCoverage:
+class DataEvolutionGlobalIndexCoverage:
"""Computes global-index coverage by field id."""
def __init__(
@@ -107,8 +107,8 @@ class GlobalIndexCoverage:
return Range.sort_and_merge_overlap(unindexed, True)
def _data_ranges_by_data_files(self) -> List[Range]:
- if hasattr(self._table, "data_ranges_for_global_index_coverage"):
- return self._table.data_ranges_for_global_index_coverage(
+ if hasattr(self._table,
"data_ranges_for_data_evolution_global_index_coverage"):
+ return
self._table.data_ranges_for_data_evolution_global_index_coverage(
self._snapshot,
self._partition_filter,
)
diff --git a/paimon-python/pypaimon/globalindex/global_index_scanner.py
b/paimon-python/pypaimon/globalindex/data_evolution_global_index_scanner.py
similarity index 96%
rename from paimon-python/pypaimon/globalindex/global_index_scanner.py
rename to
paimon-python/pypaimon/globalindex/data_evolution_global_index_scanner.py
index 5d82ebb570..4ee6b94ced 100644
--- a/paimon-python/pypaimon/globalindex/global_index_scanner.py
+++ b/paimon-python/pypaimon/globalindex/data_evolution_global_index_scanner.py
@@ -15,7 +15,7 @@
# specific language governing permissions and limitations
# under the License.
-"""Scanner for shard-based global indexes."""
+"""Scanner for shard-based global indexes on data-evolution tables."""
from concurrent.futures import ThreadPoolExecutor
from typing import Collection, Optional
@@ -27,13 +27,13 @@ from pypaimon.globalindex.global_index_result import
GlobalIndexResult
from pypaimon.common.options.core_options import CoreOptions
from pypaimon.common.options.options import Options
from pypaimon.common.predicate import Predicate
-from pypaimon.globalindex.global_index_coverage import GlobalIndexCoverage
+from pypaimon.globalindex.data_evolution_global_index_coverage import
DataEvolutionGlobalIndexCoverage
from pypaimon.read.push_down_utils import _get_all_fields
from pypaimon.schema.data_types import DataField
from pypaimon.utils.range import Range
-class GlobalIndexScanner:
+class DataEvolutionGlobalIndexScanner:
"""Scanner for shard-based global indexes."""
def __init__(
@@ -54,7 +54,7 @@ class GlobalIndexScanner:
)
self._fields = fields
self._coverage = (
- GlobalIndexCoverage(table, snapshot, partition_filter, index_files)
+ DataEvolutionGlobalIndexCoverage(table, snapshot,
partition_filter, index_files)
if table is not None
else None
)
@@ -132,8 +132,8 @@ class GlobalIndexScanner:
@staticmethod
def create(table, index_files=None, partition_filter=None, predicate=None,
- snapshot=None) -> Optional['GlobalIndexScanner']:
- """Create a GlobalIndexScanner.
+ snapshot=None) -> Optional['DataEvolutionGlobalIndexScanner']:
+ """Create a DataEvolutionGlobalIndexScanner.
Can be called in two ways:
1. create(table, index_files) - with explicit index files
@@ -148,7 +148,7 @@ class GlobalIndexScanner:
if len(index_files) == 0:
return None
core_options = _core_options(table)
- return GlobalIndexScanner(
+ return DataEvolutionGlobalIndexScanner(
fields=table.fields,
file_io=table.file_io,
index_path=table.path_factory().global_index_path_factory().index_path(),
@@ -193,7 +193,7 @@ class GlobalIndexScanner:
if len(scanned_index_files) == 0:
return None
core_options = _core_options(table)
- return GlobalIndexScanner(
+ return DataEvolutionGlobalIndexScanner(
fields=table.fields,
file_io=table.file_io,
index_path=table.path_factory().global_index_path_factory().index_path(),
@@ -221,7 +221,7 @@ class GlobalIndexScanner:
self._evaluator.close()
self._executor.shutdown(wait=False)
- def __enter__(self) -> 'GlobalIndexScanner':
+ def __enter__(self) -> 'DataEvolutionGlobalIndexScanner':
return self
def __exit__(self, exc_type, exc_val, exc_tb) -> None:
diff --git a/paimon-python/pypaimon/read/scanner/file_scanner.py
b/paimon-python/pypaimon/read/scanner/file_scanner.py
index 0a99cb6ee5..e952a9f3a4 100755
--- a/paimon-python/pypaimon/read/scanner/file_scanner.py
+++ b/paimon-python/pypaimon/read/scanner/file_scanner.py
@@ -483,10 +483,10 @@ class FileScanner:
if not self.table.options.global_index_enabled():
return None
- from pypaimon.globalindex.global_index_scanner import
GlobalIndexScanner
+ from pypaimon.globalindex.data_evolution_global_index_scanner import
DataEvolutionGlobalIndexScanner
try:
- scanner = GlobalIndexScanner.create(
+ scanner = DataEvolutionGlobalIndexScanner.create(
self.table,
partition_filter=self.partition_key_predicate,
predicate=self.predicate,
diff --git a/paimon-python/pypaimon/table/source/full_text_scan.py
b/paimon-python/pypaimon/table/source/full_text_scan.py
index 6caecea363..5e0038a184 100644
--- a/paimon-python/pypaimon/table/source/full_text_scan.py
+++ b/paimon-python/pypaimon/table/source/full_text_scan.py
@@ -21,7 +21,7 @@ from abc import ABC, abstractmethod
from collections import defaultdict
from typing import List
-from pypaimon.globalindex.global_index_coverage import GlobalIndexCoverage
+from pypaimon.globalindex.data_evolution_global_index_coverage import
DataEvolutionGlobalIndexCoverage
from pypaimon.globalindex.full_text.native_full_text_global_index_reader
import (
FULL_TEXT_IDENTIFIER,
)
@@ -117,7 +117,7 @@ class DataEvolutionFullTextScan(FullTextScan):
column_name, range_key.from_, range_key.to, files))
if all_index_files:
- raw_row_ranges = GlobalIndexCoverage(
+ raw_row_ranges = DataEvolutionGlobalIndexCoverage(
self._table,
snapshot,
partition_filter,
diff --git
a/paimon-python/pypaimon/table/source/primary_key_sorted_index_scan.py
b/paimon-python/pypaimon/table/source/primary_key_sorted_index_scan.py
index 3fe565c022..5a2c91ad6a 100644
--- a/paimon-python/pypaimon/table/source/primary_key_sorted_index_scan.py
+++ b/paimon-python/pypaimon/table/source/primary_key_sorted_index_scan.py
@@ -21,7 +21,7 @@ from pypaimon.globalindex.global_index_evaluator import
GlobalIndexEvaluator
from pypaimon.globalindex.global_index_meta import GlobalIndexIOMeta
from pypaimon.globalindex.global_index_reader import GlobalIndexReader,
_map_future
from pypaimon.globalindex.global_index_result import GlobalIndexResult
-from pypaimon.globalindex.global_index_scanner import _create_inner_readers
+from pypaimon.globalindex.data_evolution_global_index_scanner import
_create_inner_readers
from pypaimon.common.options.core_options import CoreOptions
from pypaimon.index.pk.primary_key_index_source_file import
PrimaryKeyIndexSourceFile
from pypaimon.index.pksorted.pk_sorted_bucket_index_state import
PkSortedBucketIndexState
diff --git a/paimon-python/pypaimon/table/source/vector_search_read.py
b/paimon-python/pypaimon/table/source/vector_search_read.py
index aa3fa1065f..eba11094f3 100644
--- a/paimon-python/pypaimon/table/source/vector_search_read.py
+++ b/paimon-python/pypaimon/table/source/vector_search_read.py
@@ -128,8 +128,8 @@ class AbstractVectorSearchReadImpl:
if not scalar_files:
return RoaringBitmap64()
- from pypaimon.globalindex.global_index_scanner import
GlobalIndexScanner
- scanner = GlobalIndexScanner.create(
+ from pypaimon.globalindex.data_evolution_global_index_scanner import
DataEvolutionGlobalIndexScanner
+ scanner = DataEvolutionGlobalIndexScanner.create(
self._table,
index_files=scalar_files,
partition_filter=self._partition_filter,
@@ -175,8 +175,8 @@ class AbstractVectorSearchReadImpl:
if not scalar_files:
return None
- from pypaimon.globalindex.global_index_scanner import
GlobalIndexScanner
- scanner = GlobalIndexScanner.create(
+ from pypaimon.globalindex.data_evolution_global_index_scanner import
DataEvolutionGlobalIndexScanner
+ scanner = DataEvolutionGlobalIndexScanner.create(
self._table,
index_files=scalar_files,
partition_filter=self._partition_filter,
diff --git a/paimon-python/pypaimon/table/source/vector_search_scan.py
b/paimon-python/pypaimon/table/source/vector_search_scan.py
index 68d57754cd..e9f65a5b61 100644
--- a/paimon-python/pypaimon/table/source/vector_search_scan.py
+++ b/paimon-python/pypaimon/table/source/vector_search_scan.py
@@ -20,7 +20,7 @@
from abc import ABC, abstractmethod
from collections import defaultdict
-from pypaimon.globalindex.global_index_coverage import GlobalIndexCoverage
+from pypaimon.globalindex.data_evolution_global_index_coverage import
DataEvolutionGlobalIndexCoverage
from pypaimon.table.source.vector_search_split import (
IndexVectorSearchSplit,
RawVectorSearchSplit,
@@ -163,7 +163,7 @@ class DataEvolutionVectorScan(VectorSearchScan):
)
)
- raw_row_ranges = GlobalIndexCoverage(
+ raw_row_ranges = DataEvolutionGlobalIndexCoverage(
self._table,
snapshot,
partition_filter,
@@ -177,7 +177,7 @@ class DataEvolutionVectorScan(VectorSearchScan):
if self._filter is not None:
raw_row_ranges = Range.sort_and_merge_overlap(
raw_row_ranges
- + GlobalIndexCoverage(
+ + DataEvolutionGlobalIndexCoverage(
self._table,
snapshot,
partition_filter,
diff --git a/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py
b/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py
index d12b23c130..bc48ef3ccc 100644
--- a/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py
+++ b/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py
@@ -26,7 +26,7 @@ import pyarrow as pa
from parameterized import parameterized
from pypaimon.catalog.catalog_factory import CatalogFactory
from pypaimon.data.generic_variant import GenericVariant
-from pypaimon.globalindex.global_index_scanner import GlobalIndexScanner
+from pypaimon.globalindex.data_evolution_global_index_scanner import
DataEvolutionGlobalIndexScanner
from pypaimon.schema.data_types import VectorType
from pypaimon.schema.schema import Schema
from pypaimon.read.read_builder import ReadBuilder
@@ -758,7 +758,7 @@ class JavaPyReadWriteTest(unittest.TestCase):
read_builder = table.new_read_builder()
predicate = predicate_factory(read_builder.new_predicate_builder())
- scanner = GlobalIndexScanner.create(table, predicate=predicate)
+ scanner = DataEvolutionGlobalIndexScanner.create(table,
predicate=predicate)
self.assertIsNotNone(scanner)
with scanner:
result = scanner.scan(predicate)
@@ -786,7 +786,7 @@ class JavaPyReadWriteTest(unittest.TestCase):
.new_predicate_builder()
.greater_or_equal('k', 'key-295'))
- scanner = GlobalIndexScanner.create(table, predicate=predicate)
+ scanner = DataEvolutionGlobalIndexScanner.create(table,
predicate=predicate)
self.assertIsNotNone(scanner)
with scanner:
result = scanner.scan(predicate)
@@ -804,7 +804,7 @@ class JavaPyReadWriteTest(unittest.TestCase):
self.assertEqual(expected, actual)
disabled_table = table.copy({budget_key: '0 b'})
- disabled_scanner = GlobalIndexScanner.create(disabled_table,
predicate=predicate)
+ disabled_scanner =
DataEvolutionGlobalIndexScanner.create(disabled_table, predicate=predicate)
self.assertIsNotNone(disabled_scanner)
with disabled_scanner:
disabled_result = disabled_scanner.scan(predicate)
diff --git a/paimon-python/pypaimon/tests/global_index_build_test.py
b/paimon-python/pypaimon/tests/global_index_build_test.py
index 59b3fcd352..d03eb4b2f0 100644
--- a/paimon-python/pypaimon/tests/global_index_build_test.py
+++ b/paimon-python/pypaimon/tests/global_index_build_test.py
@@ -43,7 +43,7 @@ from pypaimon.globalindex.vindex.vindex_vector_index_writer
import (
VindexVectorIndexWriter,
native_options,
)
-from pypaimon.globalindex.global_index_scanner import GlobalIndexScanner
+from pypaimon.globalindex.data_evolution_global_index_scanner import
DataEvolutionGlobalIndexScanner
from pypaimon.index.index_file_handler import IndexFileHandler
from pypaimon.schema.data_types import ArrayType, AtomicType, RowType
from pypaimon.tests.data_evolution_test_helpers import (
@@ -211,7 +211,7 @@ class GlobalIndexBuildTest(
read_builder = table.new_read_builder()
predicate = read_builder.new_predicate_builder().equal('id', 2)
- with GlobalIndexScanner.create(
+ with DataEvolutionGlobalIndexScanner.create(
table,
predicate=predicate,
snapshot=snapshot) as scanner:
@@ -372,7 +372,7 @@ class GlobalIndexBuildTest(
[Range(0, 1), Range(3, 3)]),
]
for predicate, expected in cases:
- with GlobalIndexScanner.create(
+ with DataEvolutionGlobalIndexScanner.create(
table,
predicate=predicate,
snapshot=snapshot) as scanner:
diff --git a/paimon-python/pypaimon/tests/global_index_test.py
b/paimon-python/pypaimon/tests/global_index_test.py
index 135cdcc6c0..4d7a1b142b 100644
--- a/paimon-python/pypaimon/tests/global_index_test.py
+++ b/paimon-python/pypaimon/tests/global_index_test.py
@@ -93,7 +93,7 @@ class _CoverageTable:
]
self._data_ranges = data_ranges or []
- def data_ranges_for_global_index_coverage(self, snapshot,
partition_filter):
+ def data_ranges_for_data_evolution_global_index_coverage(self, snapshot,
partition_filter):
return self._data_ranges
@@ -118,13 +118,13 @@ def _coverage_index_file(field_id, start, end,
extra_field_ids=None):
)
-class GlobalIndexCoverageTest(unittest.TestCase):
+class DataEvolutionGlobalIndexCoverageTest(unittest.TestCase):
def test_fast_mode_does_not_return_unindexed_ranges(self):
- from pypaimon.globalindex.global_index_coverage import
GlobalIndexCoverage
+ from pypaimon.globalindex.data_evolution_global_index_coverage import
DataEvolutionGlobalIndexCoverage
table = _CoverageTable(GlobalIndexSearchMode.FAST)
- coverage = GlobalIndexCoverage(
+ coverage = DataEvolutionGlobalIndexCoverage(
table,
_CoverageSnapshot(10),
None,
@@ -134,10 +134,10 @@ class GlobalIndexCoverageTest(unittest.TestCase):
self.assertEqual([], coverage.unindexed_ranges(0))
def test_full_mode_uses_snapshot_next_row_id(self):
- from pypaimon.globalindex.global_index_coverage import
GlobalIndexCoverage
+ from pypaimon.globalindex.data_evolution_global_index_coverage import
DataEvolutionGlobalIndexCoverage
table = _CoverageTable("full")
- coverage = GlobalIndexCoverage(
+ coverage = DataEvolutionGlobalIndexCoverage(
table,
_CoverageSnapshot(10),
None,
@@ -147,10 +147,10 @@ class GlobalIndexCoverageTest(unittest.TestCase):
self.assertEqual([Range(5, 9)], coverage.unindexed_ranges(0))
def test_full_mode_accepts_multiple_field_ids(self):
- from pypaimon.globalindex.global_index_coverage import
GlobalIndexCoverage
+ from pypaimon.globalindex.data_evolution_global_index_coverage import
DataEvolutionGlobalIndexCoverage
table = _CoverageTable("full")
- coverage = GlobalIndexCoverage(
+ coverage = DataEvolutionGlobalIndexCoverage(
table,
_CoverageSnapshot(10),
None,
@@ -163,10 +163,10 @@ class GlobalIndexCoverageTest(unittest.TestCase):
self.assertEqual([Range(5, 9)], coverage.unindexed_ranges([0, 1]))
def test_full_mode_intersects_coverage_for_all_predicate_fields(self):
- from pypaimon.globalindex.global_index_coverage import
GlobalIndexCoverage
+ from pypaimon.globalindex.data_evolution_global_index_coverage import
DataEvolutionGlobalIndexCoverage
table = _CoverageTable("full")
- coverage = GlobalIndexCoverage(
+ coverage = DataEvolutionGlobalIndexCoverage(
table,
_CoverageSnapshot(10),
None,
@@ -188,10 +188,10 @@ class GlobalIndexCoverageTest(unittest.TestCase):
)
def test_extra_fields_count_as_index_coverage(self):
- from pypaimon.globalindex.global_index_coverage import
GlobalIndexCoverage
+ from pypaimon.globalindex.data_evolution_global_index_coverage import
DataEvolutionGlobalIndexCoverage
table = _CoverageTable("full")
- coverage = GlobalIndexCoverage(
+ coverage = DataEvolutionGlobalIndexCoverage(
table,
_CoverageSnapshot(10),
None,
@@ -201,10 +201,10 @@ class GlobalIndexCoverageTest(unittest.TestCase):
self.assertEqual([], coverage.unindexed_ranges(1))
def test_detail_mode_uses_table_data_ranges(self):
- from pypaimon.globalindex.global_index_coverage import
GlobalIndexCoverage
+ from pypaimon.globalindex.data_evolution_global_index_coverage import
DataEvolutionGlobalIndexCoverage
table = _CoverageTable("detail", data_ranges=[Range(0, 2), Range(7,
9)])
- coverage = GlobalIndexCoverage(
+ coverage = DataEvolutionGlobalIndexCoverage(
table,
_CoverageSnapshot(10),
None,
@@ -241,7 +241,7 @@ class GlobalIndexScalarFallbackTest(unittest.TestCase):
fake_scanner.__exit__.return_value = None
with unittest.mock.patch(
-
"pypaimon.globalindex.global_index_scanner.GlobalIndexScanner.create",
+
"pypaimon.globalindex.data_evolution_global_index_scanner.DataEvolutionGlobalIndexScanner.create",
return_value=fake_scanner):
result = scanner._eval_global_index(snapshot=object())
@@ -272,7 +272,7 @@ class GlobalIndexScalarFallbackTest(unittest.TestCase):
fake_scanner.__exit__.return_value = None
with unittest.mock.patch(
-
"pypaimon.globalindex.global_index_scanner.GlobalIndexScanner.create",
+
"pypaimon.globalindex.data_evolution_global_index_scanner.DataEvolutionGlobalIndexScanner.create",
return_value=fake_scanner):
result = scanner._eval_global_index(snapshot=object())
@@ -314,7 +314,7 @@ class PlanSnapshotFetchRegressionTest(
1, call_count[0],
msg=f"Plan fetched latest snapshot {call_count[0]} times — "
"duplicate from #7513: manifest_scanner + "
- "GlobalIndexScanner.create both fetch independently.")
+ "DataEvolutionGlobalIndexScanner.create both fetch
independently.")
def test_time_travel_plan(self):
table = self._create_table()
@@ -349,6 +349,6 @@ class PlanSnapshotFetchRegressionTest(
msg=f"Global index evaluated against snapshot "
f"{seen_snapshot_ids[0]}, expected time-travel snapshot "
f"{snapshot_1_id}. Before #7513 was fixed, "
- "GlobalIndexScanner.create self-fetched latest snapshot, "
+ "DataEvolutionGlobalIndexScanner.create self-fetched latest
snapshot, "
"so global index used latest while manifest used the "
"time-travel snapshot — silent correctness bug.")
diff --git a/paimon-python/pypaimon/tests/vector_search_filter_test.py
b/paimon-python/pypaimon/tests/vector_search_filter_test.py
index a116d410bf..55fe9bb07b 100644
--- a/paimon-python/pypaimon/tests/vector_search_filter_test.py
+++ b/paimon-python/pypaimon/tests/vector_search_filter_test.py
@@ -931,7 +931,7 @@ class VectorSearchFilterTest(unittest.TestCase):
return _FakeReader()
with mock.patch(
-
"pypaimon.globalindex.global_index_scanner.GlobalIndexScanner.create",
+
"pypaimon.globalindex.data_evolution_global_index_scanner.DataEvolutionGlobalIndexScanner.create",
return_value=scanner), \
mock.patch(
"pypaimon.table.source.vector_search_read._create_vector_reader",
@@ -1282,9 +1282,9 @@ class VectorSearchFilterTest(unittest.TestCase):
)
def test_scanner_threads_external_path_to_btree_reader(self):
- """GlobalIndexScanner (backing _pre_filter) must thread external_path
+ """DataEvolutionGlobalIndexScanner (backing _pre_filter) must thread
external_path
onto the GlobalIndexIOMeta handed to the btree reader factory."""
- from pypaimon.globalindex.global_index_scanner import
GlobalIndexScanner
+ from pypaimon.globalindex.data_evolution_global_index_scanner import
DataEvolutionGlobalIndexScanner
scalar_file = self.entries[2].index_file
@@ -1301,7 +1301,7 @@ class VectorSearchFilterTest(unittest.TestCase):
with mock.patch(
"pypaimon.globalindex.btree.lazy_filtered_btree_reader.LazyFilteredBTreeReader",
_FakeLazyReader):
- scanner = GlobalIndexScanner(
+ scanner = DataEvolutionGlobalIndexScanner(
fields=self.table.fields,
file_io=self.table.file_io,
index_path="/unused/index-path",
@@ -1488,7 +1488,7 @@ class VectorSearchFilterTest(unittest.TestCase):
class VectorSearchMultiShardScalarTest(unittest.TestCase):
"""Scalar pre-filter across multiple btree shards of the same field.
- Exercises the real GlobalIndexScanner reader-construction path (with
+ Exercises the real DataEvolutionGlobalIndexScanner reader-construction
path (with
OffsetGlobalIndexReader + UnionGlobalIndexReader wrapping) so that:
- Local row ids from each shard are rebased to the global row-id space
before being unioned.
@@ -1500,8 +1500,8 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
def test_hit_only_in_later_shard_returns_global_row_id(self):
from pypaimon.globalindex.global_index_result import GlobalIndexResult
- from pypaimon.globalindex.global_index_scanner import (
- GlobalIndexScanner,
+ from pypaimon.globalindex.data_evolution_global_index_scanner import (
+ DataEvolutionGlobalIndexScanner,
)
id_field = _field(0, "id")
@@ -1545,7 +1545,7 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
with mock.patch(
"pypaimon.globalindex.sorted_file_global_index_reader.SortedIndexFileMeta.deserialize",
return_value=wide_meta):
- scanner = GlobalIndexScanner(
+ scanner = DataEvolutionGlobalIndexScanner(
fields=table.fields,
file_io=table.file_io,
index_path="/unused",
@@ -1566,8 +1566,8 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
def test_extra_field_groups_are_padded_before_and(self):
from pypaimon.globalindex.global_index_reader import GlobalIndexReader
- from pypaimon.globalindex.global_index_scanner import (
- GlobalIndexScanner,
+ from pypaimon.globalindex.data_evolution_global_index_scanner import (
+ DataEvolutionGlobalIndexScanner,
)
a_field = _field(0, "a")
@@ -1606,9 +1606,9 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
return [_StubReader(io_meta.file_name) for io_meta in io_metas]
with mock.patch(
-
"pypaimon.globalindex.global_index_scanner._create_inner_readers",
+
"pypaimon.globalindex.data_evolution_global_index_scanner._create_inner_readers",
side_effect=_stub_create_inner_readers):
- scanner = GlobalIndexScanner(
+ scanner = DataEvolutionGlobalIndexScanner(
fields=[a_field, b_field, c_field],
file_io=object(),
index_path="/unused",
@@ -1629,8 +1629,8 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
def test_extra_field_groups_pad_missing_coverage_before_and(self):
from pypaimon.globalindex.global_index_reader import GlobalIndexReader
- from pypaimon.globalindex.global_index_scanner import (
- GlobalIndexScanner,
+ from pypaimon.globalindex.data_evolution_global_index_scanner import (
+ DataEvolutionGlobalIndexScanner,
)
a_field = _field(0, "a")
@@ -1675,9 +1675,9 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
return [_StubReader(io_meta.file_name) for io_meta in io_metas]
with mock.patch(
-
"pypaimon.globalindex.global_index_scanner._create_inner_readers",
+
"pypaimon.globalindex.data_evolution_global_index_scanner._create_inner_readers",
side_effect=_stub_create_inner_readers):
- scanner = GlobalIndexScanner(
+ scanner = DataEvolutionGlobalIndexScanner(
fields=[a_field, b_field, c_field],
file_io=object(),
index_path="/unused",
@@ -1697,8 +1697,8 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
def test_extra_field_padding_does_not_convert_none_to_hits(self):
from pypaimon.globalindex.global_index_reader import GlobalIndexReader
- from pypaimon.globalindex.global_index_scanner import (
- GlobalIndexScanner,
+ from pypaimon.globalindex.data_evolution_global_index_scanner import (
+ DataEvolutionGlobalIndexScanner,
)
a_field = _field(0, "a", "STRING")
@@ -1727,9 +1727,9 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
return [_StubReader() for _ in io_metas]
with mock.patch(
-
"pypaimon.globalindex.global_index_scanner._create_inner_readers",
+
"pypaimon.globalindex.data_evolution_global_index_scanner._create_inner_readers",
side_effect=_stub_create_inner_readers):
- scanner = GlobalIndexScanner(
+ scanner = DataEvolutionGlobalIndexScanner(
fields=[a_field, b_field, c_field],
file_io=object(),
index_path="/unused",
@@ -1746,12 +1746,12 @@ class
VectorSearchMultiShardScalarTest(unittest.TestCase):
def test_native_fulltext_index_is_dispatched_by_scanner(self):
"""Non-btree scalar global indexes (full-text, etc.) must be
- instantiated by GlobalIndexScanner — previously only 'btree' was
+ instantiated by DataEvolutionGlobalIndexScanner — previously only
'btree' was
handled and everything else was silently dropped, making text-column
pre-filter a no-op."""
from pypaimon.globalindex.global_index_result import GlobalIndexResult
- from pypaimon.globalindex.global_index_scanner import (
- GlobalIndexScanner,
+ from pypaimon.globalindex.data_evolution_global_index_scanner import (
+ DataEvolutionGlobalIndexScanner,
)
name_field = _field(0, "name", "STRING")
@@ -1785,7 +1785,7 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
with mock.patch(
"pypaimon.globalindex.full_text.NativeFullTextGlobalIndexReader",
_StubFullTextReader):
- scanner = GlobalIndexScanner(
+ scanner = DataEvolutionGlobalIndexScanner(
fields=table.fields,
file_io=table.file_io,
index_path="/unused",
@@ -1814,8 +1814,8 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
the pre-filter is silently skipped and vector search returns rows
that violate the predicate."""
from pypaimon.globalindex.global_index_result import GlobalIndexResult
- from pypaimon.globalindex.global_index_scanner import (
- GlobalIndexScanner,
+ from pypaimon.globalindex.data_evolution_global_index_scanner import (
+ DataEvolutionGlobalIndexScanner,
)
name_field = _field(0, "name", "STRING")
@@ -1847,7 +1847,7 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
with mock.patch(
"pypaimon.globalindex.sorted_file_global_index_reader.SortedIndexFileMeta.deserialize",
return_value=BTreeIndexMeta(first_key=b'',
last_key=b'zzzz', has_nulls=False)):
- scanner = GlobalIndexScanner(
+ scanner = DataEvolutionGlobalIndexScanner(
fields=table.fields,
file_io=table.file_io,
index_path="/unused",
@@ -1867,8 +1867,8 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
def test_scanner_reports_unindexed_rows_for_full_mode(self):
from pypaimon.common.options.core_options import CoreOptions
from pypaimon.common.options.options import Options
- from pypaimon.globalindex.global_index_scanner import (
- GlobalIndexScanner,
+ from pypaimon.globalindex.data_evolution_global_index_scanner import (
+ DataEvolutionGlobalIndexScanner,
)
class _Options:
@@ -1893,7 +1893,7 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
table.options = _Options()
table.snapshot_manager = lambda: _Snapshots()
- scanner = GlobalIndexScanner.create(table, index_files=[indexed])
+ scanner = DataEvolutionGlobalIndexScanner.create(table,
index_files=[indexed])
try:
result = scanner.unindexed_rows(
Predicate(method="equal", index=0, field="id", literals=[7]))
@@ -1903,8 +1903,8 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
self.assertEqual([Range(5, 9)], result.results().to_range_list())
def test_scanner_create_selects_extra_field_indexes(self):
- from pypaimon.globalindex.global_index_scanner import (
- GlobalIndexScanner,
+ from pypaimon.globalindex.data_evolution_global_index_scanner import (
+ DataEvolutionGlobalIndexScanner,
)
name_field = _field(0, "name", "STRING")
@@ -1940,7 +1940,7 @@ class VectorSearchMultiShardScalarTest(unittest.TestCase):
with mock.patch(
"pypaimon.globalindex.sorted_file_global_index_reader.SortedIndexFileMeta.deserialize",
return_value=BTreeIndexMeta(first_key=b'',
last_key=b'zzzz', has_nulls=False)):
- scanner = GlobalIndexScanner.create(
+ scanner = DataEvolutionGlobalIndexScanner.create(
table,
predicate=Predicate(method="equal", index=1, field="id",
literals=[3]),