Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -100,9 +100,10 @@
* A union {@link SplitRead} to read multiple inner files to merge columns.
*
* <p>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.
*
* <p>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}.
Expand Down Expand Up @@ -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);
});
}
}
Expand All @@ -296,7 +300,8 @@ private RecordReader<InternalRow> createUnionReader(
DataFilePathFactory dataFilePathFactory,
List<Range> rowRanges,
RowType readRowType,
@Nullable DeletionVectorWithRange deletionVector)
@Nullable DeletionVectorWithRange deletionVector,
@Nullable BitmapIndexResult groupSelection)
throws IOException {
List<DataEvolutionVectorReadPlanner.ReadRange> vectorRanges =
DataEvolutionVectorReadPlanner.plan(
Expand Down Expand Up @@ -341,7 +346,8 @@ private RecordReader<InternalRow> createUnionReader(
dataFilePathFactory,
readRanges,
readRowType,
deletionVector));
deletionVector,
null));
}
return ConcatRecordReader.create(suppliers);
}
Expand All @@ -365,7 +371,8 @@ private RecordReader<InternalRow> createUnionReader(
dataFilePathFactory,
rowRanges,
readRowType,
deletionVector);
deletionVector,
groupSelection);
}

private RecordReader<InternalRow> createUnionReader(
Expand All @@ -375,7 +382,8 @@ private RecordReader<InternalRow> createUnionReader(
DataFilePathFactory dataFilePathFactory,
List<Range> rowRanges,
RowType readRowType,
@Nullable DeletionVectorWithRange deletionVector)
@Nullable DeletionVectorWithRange deletionVector,
@Nullable BitmapIndexResult groupSelection)
throws IOException {

long rowCount = fieldsFiles.get(0).rowCount();
Expand Down Expand Up @@ -424,7 +432,8 @@ private RecordReader<InternalRow> createUnionReader(
formatBuilder,
rowRanges,
readRowType,
deletionVector);
deletionVector,
groupSelection);
}

// Build the per-bunch readers from the planned partial read row types.
Expand Down Expand Up @@ -464,7 +473,8 @@ private RecordReader<InternalRow> createUnionReader(
formatReaderMapping,
rowRanges,
partialReadRowType,
deletionVector));
deletionVector,
groupSelection));
}

return nestedFieldEnabled
Expand All @@ -482,7 +492,8 @@ private RecordReader<InternalRow> createMissingFieldsReader(
Builder formatBuilder,
List<Range> 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.
Expand All @@ -500,7 +511,8 @@ private RecordReader<InternalRow> createMissingFieldsReader(
mapping,
rowRanges,
readRowType,
deletionVector);
deletionVector,
groupSelection);
}

private boolean nestedFieldEnabledFor(List<DataFileMeta> files) {
Expand Down Expand Up @@ -565,7 +577,8 @@ private RecordReader<InternalRow> createFieldBunchReader(
FormatReaderMapping formatReaderMapping,
List<Range> 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
Expand All @@ -576,7 +589,8 @@ private RecordReader<InternalRow> 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(
Expand All @@ -585,7 +599,8 @@ private RecordReader<InternalRow> createFieldBunchReader(
dataFilePathFactory,
formatReaderMapping,
rowRanges,
deletionVector);
deletionVector,
groupSelection);
} else if (bunch instanceof BlobFileBunch) {
// for blob bunch, fallback on placeholders
BlobFileBunch blobBunch = (BlobFileBunch) bunch;
Expand Down Expand Up @@ -619,7 +634,8 @@ private RecordReader<InternalRow> sequentialReadFiles(
DataFilePathFactory dataFilePathFactory,
FormatReaderMapping formatReaderMapping,
List<Range> rowRanges,
@Nullable DeletionVectorWithRange deletionVector)
@Nullable DeletionVectorWithRange deletionVector,
@Nullable BitmapIndexResult groupSelection)
throws IOException {
List<ReaderSupplier<InternalRow>> readerSuppliers = new ArrayList<>();
for (DataFileMeta file : files) {
Expand All @@ -636,7 +652,7 @@ private RecordReader<InternalRow> sequentialReadFiles(
dataFilePathFactory.toPath(file),
file.fileSize()),
deletionVector,
null));
groupSelection));
}
return ConcatRecordReader.create(readerSuppliers);
}
Expand Down Expand Up @@ -727,14 +743,35 @@ private FileRecordReader<InternalRow> createFileReader(
return createFileReader(
partition,
file,
dataFilePathFactory,
formatReaderMapping,
rowRanges,
readRowType,
readTarget(file, dataFilePathFactory, rowRanges),
deletionVector,
null);
}

private FileRecordReader<InternalRow> createFileReader(
BinaryRow partition,
DataFileMeta file,
DataFilePathFactory dataFilePathFactory,
FormatReaderMapping formatReaderMapping,
List<Range> 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<InternalRow> createFileReader(
BinaryRow partition,
DataFileMeta file,
Expand Down Expand Up @@ -912,6 +949,47 @@ private boolean skipByFileIndex(
return false;
}

/**
* Builds one row selection for a normal merged group from the cached file-index results.
*
* <p>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<DataFileMeta> files, List<FileIndexResultEntry> 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<DataFileMeta> 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<Predicate> filters, boolean nestedFieldEnabled) {
return new Builder(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<InternalRow> 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<InternalRow> 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<InternalRow> 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<InternalRow> 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());
Expand Down Expand Up @@ -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<InternalRow> 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());
Expand Down
Loading