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());