From e65221b00cceaad5f0cd5030e0d4d7e7cf662dfe Mon Sep 17 00:00:00 2001
From: sanshi <1715734693@qq.com>
Date: Mon, 28 Sep 2026 15:16:40 +0800
Subject: [PATCH] [core] Push down bitmap indexes for data evolution merged
reads
---
.../operation/DataEvolutionSplitRead.java | 114 +++++++++++++++---
.../table/DataEvolutionFileIndexTest.java | 42 ++++++-
2 files changed, 135 insertions(+), 21 deletions(-)
diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java b/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
index 696831fbf6bb..60902d352aa2 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
@@ -100,9 +100,10 @@
* A union {@link SplitRead} to read multiple inner files to merge columns.
*
*
Filters can only be pushed down where they can not interfere with the column merging: a file
- * read without merging gets both the file index and the format level push down, while a merged
- * group only uses the file index to skip the whole group, as dropping rows in one of the merged
- * readers would break the positional alignment between them.
+ * read without merging gets both the file index and the format level push down, while an eligible
+ * merged group shares one bitmap selection across all aligned field readers. Other merged groups
+ * only use the file index to skip the whole group, as dropping rows in one reader would break the
+ * positional alignment between them.
*
*
Only a filter whose every field belongs to the read type is pushed down, see {@link
* #readTypeFilters}, and only to the files that wrote those fields, see {@link #fileFilters}.
@@ -267,13 +268,16 @@ && skipByFileIndex(
fileIndexResults, rowRanges, deletionVector)) {
return new EmptyFileRecordReader<>();
}
+ BitmapIndexResult groupSelection =
+ buildGroupSelection(needMergeFiles, fileIndexResults);
return createUnionReader(
needMergeFiles,
partition,
dataFilePathFactory,
rowRanges,
readRowType,
- deletionVector);
+ deletionVector,
+ groupSelection);
});
}
}
@@ -296,7 +300,8 @@ private RecordReader createUnionReader(
DataFilePathFactory dataFilePathFactory,
List rowRanges,
RowType readRowType,
- @Nullable DeletionVectorWithRange deletionVector)
+ @Nullable DeletionVectorWithRange deletionVector,
+ @Nullable BitmapIndexResult groupSelection)
throws IOException {
List vectorRanges =
DataEvolutionVectorReadPlanner.plan(
@@ -341,7 +346,8 @@ private RecordReader createUnionReader(
dataFilePathFactory,
readRanges,
readRowType,
- deletionVector));
+ deletionVector,
+ null));
}
return ConcatRecordReader.create(suppliers);
}
@@ -365,7 +371,8 @@ private RecordReader createUnionReader(
dataFilePathFactory,
rowRanges,
readRowType,
- deletionVector);
+ deletionVector,
+ groupSelection);
}
private RecordReader createUnionReader(
@@ -375,7 +382,8 @@ private RecordReader createUnionReader(
DataFilePathFactory dataFilePathFactory,
List rowRanges,
RowType readRowType,
- @Nullable DeletionVectorWithRange deletionVector)
+ @Nullable DeletionVectorWithRange deletionVector,
+ @Nullable BitmapIndexResult groupSelection)
throws IOException {
long rowCount = fieldsFiles.get(0).rowCount();
@@ -424,7 +432,8 @@ private RecordReader createUnionReader(
formatBuilder,
rowRanges,
readRowType,
- deletionVector);
+ deletionVector,
+ groupSelection);
}
// Build the per-bunch readers from the planned partial read row types.
@@ -464,7 +473,8 @@ private RecordReader createUnionReader(
formatReaderMapping,
rowRanges,
partialReadRowType,
- deletionVector));
+ deletionVector,
+ groupSelection));
}
return nestedFieldEnabled
@@ -482,7 +492,8 @@ private RecordReader createMissingFieldsReader(
Builder formatBuilder,
List rowRanges,
RowType readRowType,
- @Nullable DeletionVectorWithRange deletionVector)
+ @Nullable DeletionVectorWithRange deletionVector,
+ @Nullable BitmapIndexResult groupSelection)
throws IOException {
DataFileMeta firstFile = bunch.files().get(0);
// Use the physical schema: the full table schema may declare columns this file never wrote.
@@ -500,7 +511,8 @@ private RecordReader createMissingFieldsReader(
mapping,
rowRanges,
readRowType,
- deletionVector);
+ deletionVector,
+ groupSelection);
}
private boolean nestedFieldEnabledFor(List files) {
@@ -565,7 +577,8 @@ private RecordReader createFieldBunchReader(
FormatReaderMapping formatReaderMapping,
List rowRanges,
RowType readRowType,
- @Nullable DeletionVectorWithRange deletionVector)
+ @Nullable DeletionVectorWithRange deletionVector,
+ @Nullable BitmapIndexResult groupSelection)
throws IOException {
if (bunch instanceof DataBunch) {
// for data bunch, directly read the single file
@@ -576,7 +589,8 @@ private RecordReader createFieldBunchReader(
formatReaderMapping,
rowRanges,
readRowType,
- deletionVector);
+ deletionVector,
+ groupSelection);
} else if (bunch instanceof VectorFileBunch) {
// for vector bunch, sequential read all data files and concat them
return sequentialReadFiles(
@@ -585,7 +599,8 @@ private RecordReader createFieldBunchReader(
dataFilePathFactory,
formatReaderMapping,
rowRanges,
- deletionVector);
+ deletionVector,
+ groupSelection);
} else if (bunch instanceof BlobFileBunch) {
// for blob bunch, fallback on placeholders
BlobFileBunch blobBunch = (BlobFileBunch) bunch;
@@ -619,7 +634,8 @@ private RecordReader sequentialReadFiles(
DataFilePathFactory dataFilePathFactory,
FormatReaderMapping formatReaderMapping,
List rowRanges,
- @Nullable DeletionVectorWithRange deletionVector)
+ @Nullable DeletionVectorWithRange deletionVector,
+ @Nullable BitmapIndexResult groupSelection)
throws IOException {
List> readerSuppliers = new ArrayList<>();
for (DataFileMeta file : files) {
@@ -636,7 +652,7 @@ private RecordReader sequentialReadFiles(
dataFilePathFactory.toPath(file),
file.fileSize()),
deletionVector,
- null));
+ groupSelection));
}
return ConcatRecordReader.create(readerSuppliers);
}
@@ -727,14 +743,35 @@ private FileRecordReader createFileReader(
return createFileReader(
partition,
file,
+ dataFilePathFactory,
formatReaderMapping,
rowRanges,
readRowType,
- readTarget(file, dataFilePathFactory, rowRanges),
deletionVector,
null);
}
+ private FileRecordReader createFileReader(
+ BinaryRow partition,
+ DataFileMeta file,
+ DataFilePathFactory dataFilePathFactory,
+ FormatReaderMapping formatReaderMapping,
+ List rowRanges,
+ RowType readRowType,
+ @Nullable DeletionVectorWithRange deletionVector,
+ @Nullable BitmapIndexResult groupSelection)
+ throws IOException {
+ return createFileReader(
+ partition,
+ file,
+ formatReaderMapping,
+ rowRanges,
+ readRowType,
+ readTarget(file, dataFilePathFactory, rowRanges),
+ deletionVector,
+ groupSelection);
+ }
+
private FileRecordReader createFileReader(
BinaryRow partition,
DataFileMeta file,
@@ -912,6 +949,47 @@ private boolean skipByFileIndex(
return false;
}
+ /**
+ * Builds one row selection for a normal merged group from the cached file-index results.
+ *
+ * Each result is a necessary condition for the final row, so independent bitmap results are
+ * intersected. An OR predicate is already evaluated as one bitmap by {@link
+ * FileIndexEvaluator}; treating its result as one condition preserves the OR semantics. Results
+ * that cannot provide an exact bitmap remain conservative and do not restrict the group.
+ */
+ @Nullable
+ private BitmapIndexResult buildGroupSelection(
+ List files, List fileIndexResults) {
+ if (!canPushDownGroupSelection(files)) {
+ return null;
+ }
+
+ BitmapIndexResult selection = null;
+ for (FileIndexResultEntry entry : fileIndexResults) {
+ if (entry.result instanceof BitmapIndexResult) {
+ BitmapIndexResult bitmap = (BitmapIndexResult) entry.result;
+ selection = selection == null ? bitmap : (BitmapIndexResult) selection.and(bitmap);
+ }
+ }
+ return selection;
+ }
+
+ private boolean canPushDownGroupSelection(List files) {
+ if (files.isEmpty()) {
+ return false;
+ }
+
+ Range groupRange = files.get(0).nonNullRowIdRange();
+ for (DataFileMeta file : files) {
+ if (isBlobFile(file.fileName())
+ || isVectorStoreFile(file.fileName())
+ || !groupRange.equals(file.nonNullRowIdRange())) {
+ return false;
+ }
+ }
+ return true;
+ }
+
private Builder formatBuilder(
RowType readRowType, @Nullable List filters, boolean nestedFieldEnabled) {
return new Builder(
diff --git a/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionFileIndexTest.java b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionFileIndexTest.java
index 30f4be2dca28..c350c424f6d0 100644
--- a/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionFileIndexTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionFileIndexTest.java
@@ -152,12 +152,47 @@ public void testMergedGroupKeepsColumnsAligned() throws Exception {
writeSplitColumns(table, ROW_COUNT, Collections.emptyMap(), bloomOptions("f2", null));
assertMergedGroup(table);
- // a merged group is never row filtered, but the rows it returns must stay aligned
+ // A non-bitmap index cannot select rows, but the merged fields must stay aligned.
List rows = readWithFilter(table, equalF2(f2(50)));
assertThat(rows).hasSize(ROW_COUNT);
assertAligned(rows);
}
+ @Test
+ public void testMergedGroupBitmapIndexSelectsMatchingRowsOnly() throws Exception {
+ FileStoreTable table = createTable("merged_bitmap", Collections.emptyMap());
+ writeSplitColumns(table, ROW_COUNT, Collections.emptyMap(), bitmapOptions("f2"));
+ assertMergedGroup(table);
+
+ List rows = readWithFilter(table, equalF2(f2(50)));
+ assertThat(rows).hasSize(1);
+ assertRow(rows.get(0), 50);
+ }
+
+ @Test
+ public void testMergedGroupBitmapIndexIntersectsFieldSelections() throws Exception {
+ FileStoreTable table = createTable("merged_bitmap_intersect", Collections.emptyMap());
+ writeSplitColumns(table, ROW_COUNT, bitmapOptions("f1"), bitmapOptions("f2"));
+ assertMergedGroup(table);
+
+ Predicate filter = PredicateBuilder.and(equalF1(f1(50)), equalF2(f2(50)));
+ List rows = readWithFilter(table, filter);
+ assertThat(rows).hasSize(1);
+ assertRow(rows.get(0), 50);
+ }
+
+ @Test
+ public void testMergedGroupBitmapIndexPreservesOrSelection() throws Exception {
+ FileStoreTable table = createTable("merged_bitmap_or", Collections.emptyMap());
+ writeSplitColumns(table, ROW_COUNT, Collections.emptyMap(), bitmapOptions("f2"));
+ assertMergedGroup(table);
+
+ Predicate filter = PredicateBuilder.or(equalF2(f2(50)), equalF2(f2(51)));
+ List rows = readWithFilter(table, filter);
+ assertThat(rows).extracting(row -> row.getInt(0)).containsExactlyInAnyOrder(50, 51);
+ assertAligned(rows);
+ }
+
@Test
public void testQueryMergedGroup() throws Exception {
FileStoreTable table = createTable("merged_group_query", Collections.emptyMap());
@@ -579,11 +614,12 @@ public void testMergedGroupFileIndexComposesWithDeletionVector() throws Exceptio
assertThat(latest.fileIO().delete(anchorPath, false)).isTrue();
// The deleted row is the only bitmap hit in the second merged group. The missing anchor
- // file therefore proves that the group was skipped before any union reader opened it.
+ // file therefore proves that the group was skipped before any union reader opened it;
+ // the first group contributes its matching row through the shared bitmap selection.
RowType readType =
rowTypeWithRowId(rowType()).project(SpecialFields.ROW_ID.name(), "f1", "f2");
List rows = readWithFilter(table, equalF1(f1(50)), readType);
- assertThat(rowIds(rows)).containsExactlyElementsOf(rowIds(0, ROW_COUNT));
+ assertThat(rowIds(rows)).containsExactly(50L);
FileStoreTable neighbour = createTable("merged_bitmap_dv_neighbour", options);
writeSplitColumns(neighbour, ROW_COUNT, bitmapOptions("f1"), Collections.emptyMap());