peter-toth commented on code in PR #58895:
URL: https://github.com/apache/spark/pull/58895#discussion_r4106892022
##########
sql/core/src/main/java/org/apache/spark/sql/execution/datasources/parquet/VectorizedParquetRecordReader.java:
##########
@@ -492,6 +930,630 @@ private void checkEndOfRowGroup() throws IOException {
totalCountLoadedSoFar += pages.getRowCount();
}
+ /**
+ * Loads the next row group using the three-phase late-materialization
pattern, all driven by the
+ * single {@link #lateMatReader} with its requested schema mutated per phase:
+ * - Phase 0 (full schema): compute {@code pushedFilterRanges} from the
pushed data filter via
+ * column index (metadata-only) using {@link
ParquetFileReader#getRowRanges}.
+ * - Phase 1 (key-only schema): read key-column pages restricted to {@code
pushedFilterRanges},
+ * evaluate the storage filter per row, build {@code finalRanges}.
+ * - Phase 2: read the non-key columns restricted to {@code finalRanges}.
A row group that gave
+ * splicing up reads the whole projection instead, still under {@code
finalRanges}, and one
+ * that gave the filter up reads it under {@code pushedFilterRanges},
which is what a plain
+ * scan reads. Skipped entirely only for an all-keys projection that is
still splicing, since
+ * emit then builds every batch from the key queues alone.
+ *
+ * Row groups for which {@code finalRanges} is empty are skipped entirely
(no phase-2 IO).
+ * Sets {@link #hitEndOfData} when all row groups have been processed.
+ */
+ private void loadNextRowGroupWithLateMaterialization() throws IOException {
+ while (nextBlockIndex < totalBlockCount) {
+ int blockIdx = nextBlockIndex++;
+ long blockRowCount =
lateMatReader.getRowGroups().get(blockIdx).getRowCount();
+ if (blockRowCount == 0) {
+ // parquet-mr never writes these, but RowRanges.createSingle(0) would
build Range(0, -1) and
+ // trip parquet's own `from <= to` assertion. The plain read path
skips them too.
+ continue;
+ }
+ // Splicing buffers one key value per surviving row of the whole row
group before it can emit
+ // the first batch, and that buffer is outside any MemoryConsumer, so
phase 1 counts what it
+ // holds against `maxSplicedRowGroupBytes` together with the row ranges
phase 2 will hold.
+ // Past that it gives splicing up, and past it again the filter itself,
which is what
+ // `filterGivenUp` says. A file already known to have no offset index
starts there.
+ filterGivenUp = fileHasNoOffsetIndex;
+ spliceCurrentRowGroup = !filterGivenUp;
+ splicedBytes = 0L;
+
+ // Phase 0: rows allowed by the pushed data filter, at column-index
granularity. The full
+ // requestedSchema goes back on first, because phases 1 and 2 narrow it
and
+ // ParquetFileReader.getRowRanges computes ranges against the reader's
current paths.
+ lateMatReader.setRequestedSchema(requestedColumns);
+ // getRowRanges checks only whether a filter is pushed, not
options.useColumnIndexFilter(),
+ // so calling it unconditionally would keep applying column-index
filtering after a user
+ // turned it off, which is the escape hatch for a file whose column
index is wrong. Every
+ // phase below reads within these ranges, so a wrong column index would
cost rows the plain
+ // path would have returned. Phase 2 is unaffected: it selects pages
through the offset index,
+ // a separate structure this conf says nothing about.
+ RowRanges pushedFilterRanges = useColumnIndexFilter
+ ? lateMatReader.getRowRanges(blockIdx)
+ : RowRanges.createSingle(blockRowCount);
+ // RowRanges.rowCount() walks every range, so resolve each range set's
count once.
+ long baselineRows = pushedFilterRanges.rowCount();
+ if (baselineRows == 0) {
+ // Pushed data filter rejects this block entirely via column index.
Not a storage-filter
+ // skip, so we don't increment storage-filter metrics.
+ continue;
+ }
+
+ // What this feature can avoid reading is the non-key columns of the
rows the storage filter
+ // rejects, so that is the baseline both byte metrics are measured
against: the non-key bytes
+ // a plain read of this projection would transfer for every row the
pushed filter kept. The
+ // null checks only skip work for a caller that drives this reader
without a scan's metrics;
+ // FileSourceScanLike creates all five whenever storageFilters is
non-empty.
+ // compressedBytesForRowRanges never does IO of its own. A row group
whose filter is already
+ // given up reports nothing either way, so it does not pay for the
baseline at all.
+ StorageFilterMetrics m = storageFilter.metrics();
+ SQLMetric bytesAvoidedRg = m.bytesAvoidedByRowGroup();
+ SQLMetric bytesAvoidedPf = m.bytesAvoidedByPageFiltering();
+ boolean needBytes = (bytesAvoidedRg != null || bytesAvoidedPf != null)
&& !filterGivenUp;
+ Map<ColumnPath, ColumnChunkMetaData> blockChunks =
+ needBytes ? chunksByPath(lateMatReader, blockIdx) : null;
+ long nonKeyBaselineBytes = needBytes
+ ? compressedBytesForRowRanges(lateMatReader, blockIdx, blockChunks,
nonKeyColumns,
+ pushedFilterRanges, baselineRows)
+ : 0L;
+
+ // Phase 1: switch to key-only schema, read key columns under
pushedFilterRanges, evaluate the
+ // storage filter per row. Skipped for a row group the filter is already
given up for, which
+ // leaves every row of `pushedFilterRanges` to emit, exactly what a
plain read would.
+ RowRanges finalRanges = pushedFilterRanges;
+ long finalRowCount = baselineRows;
+ if (!filterGivenUp) {
+ lateMatReader.setRequestedSchema(keyOnlyColumns);
+ PageReadStore keyPages = lateMatReader.readFilteredRowGroup(blockIdx,
pushedFilterRanges);
Review Comment:
Done: try-with-resources around the phase-1 acquisition and evaluation, so
the store is closed on every exit, including skip, give-up and exception.
The ownership point is right and worth writing down, which the code now
does: `readFilteredRowGroup` hands out a store the file reader does not track,
unlike `readNextRowGroup`, so nothing else would close it. Phase 2's store is
closed by this reader too, at the top of the next row group's load and in
`close()`. c29625f
##########
sql/core/src/main/java/org/apache/spark/sql/execution/datasources/parquet/VectorizedParquetRecordReader.java:
##########
@@ -492,6 +930,630 @@ private void checkEndOfRowGroup() throws IOException {
totalCountLoadedSoFar += pages.getRowCount();
}
+ /**
+ * Loads the next row group using the three-phase late-materialization
pattern, all driven by the
+ * single {@link #lateMatReader} with its requested schema mutated per phase:
+ * - Phase 0 (full schema): compute {@code pushedFilterRanges} from the
pushed data filter via
+ * column index (metadata-only) using {@link
ParquetFileReader#getRowRanges}.
+ * - Phase 1 (key-only schema): read key-column pages restricted to {@code
pushedFilterRanges},
+ * evaluate the storage filter per row, build {@code finalRanges}.
+ * - Phase 2: read the non-key columns restricted to {@code finalRanges}.
A row group that gave
+ * splicing up reads the whole projection instead, still under {@code
finalRanges}, and one
+ * that gave the filter up reads it under {@code pushedFilterRanges},
which is what a plain
+ * scan reads. Skipped entirely only for an all-keys projection that is
still splicing, since
+ * emit then builds every batch from the key queues alone.
+ *
+ * Row groups for which {@code finalRanges} is empty are skipped entirely
(no phase-2 IO).
+ * Sets {@link #hitEndOfData} when all row groups have been processed.
+ */
+ private void loadNextRowGroupWithLateMaterialization() throws IOException {
+ while (nextBlockIndex < totalBlockCount) {
+ int blockIdx = nextBlockIndex++;
+ long blockRowCount =
lateMatReader.getRowGroups().get(blockIdx).getRowCount();
+ if (blockRowCount == 0) {
+ // parquet-mr never writes these, but RowRanges.createSingle(0) would
build Range(0, -1) and
+ // trip parquet's own `from <= to` assertion. The plain read path
skips them too.
+ continue;
+ }
+ // Splicing buffers one key value per surviving row of the whole row
group before it can emit
+ // the first batch, and that buffer is outside any MemoryConsumer, so
phase 1 counts what it
+ // holds against `maxSplicedRowGroupBytes` together with the row ranges
phase 2 will hold.
+ // Past that it gives splicing up, and past it again the filter itself,
which is what
+ // `filterGivenUp` says. A file already known to have no offset index
starts there.
+ filterGivenUp = fileHasNoOffsetIndex;
+ spliceCurrentRowGroup = !filterGivenUp;
+ splicedBytes = 0L;
+
+ // Phase 0: rows allowed by the pushed data filter, at column-index
granularity. The full
+ // requestedSchema goes back on first, because phases 1 and 2 narrow it
and
+ // ParquetFileReader.getRowRanges computes ranges against the reader's
current paths.
+ lateMatReader.setRequestedSchema(requestedColumns);
+ // getRowRanges checks only whether a filter is pushed, not
options.useColumnIndexFilter(),
+ // so calling it unconditionally would keep applying column-index
filtering after a user
+ // turned it off, which is the escape hatch for a file whose column
index is wrong. Every
+ // phase below reads within these ranges, so a wrong column index would
cost rows the plain
+ // path would have returned. Phase 2 is unaffected: it selects pages
through the offset index,
+ // a separate structure this conf says nothing about.
+ RowRanges pushedFilterRanges = useColumnIndexFilter
+ ? lateMatReader.getRowRanges(blockIdx)
+ : RowRanges.createSingle(blockRowCount);
+ // RowRanges.rowCount() walks every range, so resolve each range set's
count once.
+ long baselineRows = pushedFilterRanges.rowCount();
+ if (baselineRows == 0) {
+ // Pushed data filter rejects this block entirely via column index.
Not a storage-filter
+ // skip, so we don't increment storage-filter metrics.
+ continue;
+ }
+
+ // What this feature can avoid reading is the non-key columns of the
rows the storage filter
+ // rejects, so that is the baseline both byte metrics are measured
against: the non-key bytes
+ // a plain read of this projection would transfer for every row the
pushed filter kept. The
+ // null checks only skip work for a caller that drives this reader
without a scan's metrics;
+ // FileSourceScanLike creates all five whenever storageFilters is
non-empty.
+ // compressedBytesForRowRanges never does IO of its own. A row group
whose filter is already
+ // given up reports nothing either way, so it does not pay for the
baseline at all.
+ StorageFilterMetrics m = storageFilter.metrics();
+ SQLMetric bytesAvoidedRg = m.bytesAvoidedByRowGroup();
+ SQLMetric bytesAvoidedPf = m.bytesAvoidedByPageFiltering();
+ boolean needBytes = (bytesAvoidedRg != null || bytesAvoidedPf != null)
&& !filterGivenUp;
+ Map<ColumnPath, ColumnChunkMetaData> blockChunks =
+ needBytes ? chunksByPath(lateMatReader, blockIdx) : null;
+ long nonKeyBaselineBytes = needBytes
+ ? compressedBytesForRowRanges(lateMatReader, blockIdx, blockChunks,
nonKeyColumns,
+ pushedFilterRanges, baselineRows)
+ : 0L;
+
+ // Phase 1: switch to key-only schema, read key columns under
pushedFilterRanges, evaluate the
+ // storage filter per row. Skipped for a row group the filter is already
given up for, which
+ // leaves every row of `pushedFilterRanges` to emit, exactly what a
plain read would.
+ RowRanges finalRanges = pushedFilterRanges;
+ long finalRowCount = baselineRows;
+ if (!filterGivenUp) {
+ lateMatReader.setRequestedSchema(keyOnlyColumns);
+ PageReadStore keyPages = lateMatReader.readFilteredRowGroup(blockIdx,
pushedFilterRanges);
Review Comment:
This one cannot arise, and there is now a test for the file that would
produce it.
Phase 0 sets the full projection before it asks parquet for the ranges, so
the block's `ColumnIndexStore` is built over every projected column, key
columns included. `ColumnIndexStoreImpl.create` returns `EMPTY` when any of
them has no offset index, and with an empty store `ColumnIndexFilter` narrows
nothing, so `pushedFilterRanges` covers the whole block.
`readFilteredRowGroup(int, RowRanges)` then short-circuits to
`internalReadRowGroup` and consults no index at all. The mixed-index file
therefore reaches phase 1 with whole-block ranges: one column without an offset
index degrades the whole block, which is what makes the phase-1 read safe.
The new test writes a file with no offset index for any column and asserts
phase 1 still runs on it. The phase-2 boundary keeps the catch, because that is
where a strict subset is asked for, and the retry there no longer trusts the
argument above: if `pushedFilterRanges` were ever a strict subset on such a
file, it fails with a message that says so rather than throwing from inside
parquet. c29625f
##########
sql/core/src/main/java/org/apache/spark/sql/execution/datasources/parquet/VectorizedParquetRecordReader.java:
##########
@@ -492,6 +930,630 @@ private void checkEndOfRowGroup() throws IOException {
totalCountLoadedSoFar += pages.getRowCount();
}
+ /**
+ * Loads the next row group using the three-phase late-materialization
pattern, all driven by the
+ * single {@link #lateMatReader} with its requested schema mutated per phase:
+ * - Phase 0 (full schema): compute {@code pushedFilterRanges} from the
pushed data filter via
+ * column index (metadata-only) using {@link
ParquetFileReader#getRowRanges}.
+ * - Phase 1 (key-only schema): read key-column pages restricted to {@code
pushedFilterRanges},
+ * evaluate the storage filter per row, build {@code finalRanges}.
+ * - Phase 2: read the non-key columns restricted to {@code finalRanges}.
A row group that gave
+ * splicing up reads the whole projection instead, still under {@code
finalRanges}, and one
+ * that gave the filter up reads it under {@code pushedFilterRanges},
which is what a plain
+ * scan reads. Skipped entirely only for an all-keys projection that is
still splicing, since
+ * emit then builds every batch from the key queues alone.
+ *
+ * Row groups for which {@code finalRanges} is empty are skipped entirely
(no phase-2 IO).
+ * Sets {@link #hitEndOfData} when all row groups have been processed.
+ */
+ private void loadNextRowGroupWithLateMaterialization() throws IOException {
+ while (nextBlockIndex < totalBlockCount) {
+ int blockIdx = nextBlockIndex++;
+ long blockRowCount =
lateMatReader.getRowGroups().get(blockIdx).getRowCount();
+ if (blockRowCount == 0) {
+ // parquet-mr never writes these, but RowRanges.createSingle(0) would
build Range(0, -1) and
+ // trip parquet's own `from <= to` assertion. The plain read path
skips them too.
+ continue;
+ }
+ // Splicing buffers one key value per surviving row of the whole row
group before it can emit
+ // the first batch, and that buffer is outside any MemoryConsumer, so
phase 1 counts what it
+ // holds against `maxSplicedRowGroupBytes` together with the row ranges
phase 2 will hold.
+ // Past that it gives splicing up, and past it again the filter itself,
which is what
+ // `filterGivenUp` says. A file already known to have no offset index
starts there.
+ filterGivenUp = fileHasNoOffsetIndex;
+ spliceCurrentRowGroup = !filterGivenUp;
+ splicedBytes = 0L;
+
+ // Phase 0: rows allowed by the pushed data filter, at column-index
granularity. The full
+ // requestedSchema goes back on first, because phases 1 and 2 narrow it
and
+ // ParquetFileReader.getRowRanges computes ranges against the reader's
current paths.
+ lateMatReader.setRequestedSchema(requestedColumns);
+ // getRowRanges checks only whether a filter is pushed, not
options.useColumnIndexFilter(),
+ // so calling it unconditionally would keep applying column-index
filtering after a user
+ // turned it off, which is the escape hatch for a file whose column
index is wrong. Every
+ // phase below reads within these ranges, so a wrong column index would
cost rows the plain
+ // path would have returned. Phase 2 is unaffected: it selects pages
through the offset index,
+ // a separate structure this conf says nothing about.
+ RowRanges pushedFilterRanges = useColumnIndexFilter
+ ? lateMatReader.getRowRanges(blockIdx)
+ : RowRanges.createSingle(blockRowCount);
+ // RowRanges.rowCount() walks every range, so resolve each range set's
count once.
+ long baselineRows = pushedFilterRanges.rowCount();
+ if (baselineRows == 0) {
+ // Pushed data filter rejects this block entirely via column index.
Not a storage-filter
+ // skip, so we don't increment storage-filter metrics.
+ continue;
+ }
+
+ // What this feature can avoid reading is the non-key columns of the
rows the storage filter
+ // rejects, so that is the baseline both byte metrics are measured
against: the non-key bytes
+ // a plain read of this projection would transfer for every row the
pushed filter kept. The
+ // null checks only skip work for a caller that drives this reader
without a scan's metrics;
+ // FileSourceScanLike creates all five whenever storageFilters is
non-empty.
+ // compressedBytesForRowRanges never does IO of its own. A row group
whose filter is already
+ // given up reports nothing either way, so it does not pay for the
baseline at all.
+ StorageFilterMetrics m = storageFilter.metrics();
+ SQLMetric bytesAvoidedRg = m.bytesAvoidedByRowGroup();
+ SQLMetric bytesAvoidedPf = m.bytesAvoidedByPageFiltering();
+ boolean needBytes = (bytesAvoidedRg != null || bytesAvoidedPf != null)
&& !filterGivenUp;
+ Map<ColumnPath, ColumnChunkMetaData> blockChunks =
+ needBytes ? chunksByPath(lateMatReader, blockIdx) : null;
+ long nonKeyBaselineBytes = needBytes
+ ? compressedBytesForRowRanges(lateMatReader, blockIdx, blockChunks,
nonKeyColumns,
+ pushedFilterRanges, baselineRows)
+ : 0L;
+
+ // Phase 1: switch to key-only schema, read key columns under
pushedFilterRanges, evaluate the
+ // storage filter per row. Skipped for a row group the filter is already
given up for, which
+ // leaves every row of `pushedFilterRanges` to emit, exactly what a
plain read would.
+ RowRanges finalRanges = pushedFilterRanges;
+ long finalRowCount = baselineRows;
+ if (!filterGivenUp) {
+ lateMatReader.setRequestedSchema(keyOnlyColumns);
+ PageReadStore keyPages = lateMatReader.readFilteredRowGroup(blockIdx,
pushedFilterRanges);
+ if (keyPages == null) {
+ // Unreachable: readFilteredRowGroup returns null only for an empty
block, and we already
+ // know pushedFilterRanges selects at least one row. Skipping the
block here would drop
+ // its surviving rows from the output, so assert rather than
`continue`.
+ throw new IllegalStateException(
+ "No key pages for row group " + blockIdx + " despite " +
baselineRows
+ + " rows selected by the pushed filter");
+ }
+ RowRanges survivors = evaluateStorageFilter(keyPages,
pushedFilterRanges);
+ if (!filterGivenUp
+ && rowRangeStateBytes(survivorRangeCount) >
storageFilter.maxSplicedRowGroupBytes()) {
+ // Phase 1 weighs the budget once per accumulator, so a row group
whose survivors fit in
+ // a single one is only caught here, with its survivors buffered.
Those are released,
+ // since the ranges they were spliced against are about to be thrown
away.
+ giveUpFilter();
+ }
+ if (!filterGivenUp) {
+ finalRanges = survivors;
+ finalRowCount = survivors.rowCount();
+ if (finalRowCount == 0) {
+ // Every surviving row was rejected by the storage filter; skip
the block entirely,
+ // which avoids the whole non-key baseline. Phase 1 still paid to
read the key columns,
+ // and that cost is not part of the baseline, so nothing is
subtracted from it here.
+ recordRowGroupSkipped(m, baselineRows, nonKeyBaselineBytes);
+ continue;
+ }
+ }
+ }
+
+ // Phase 2 reads the non-key columns under the surviving rows, or the
whole projection under
+ // `pushedFilterRanges` for a row group whose filter was given up. It is
skipped only when the
+ // projection is all keys and their values were buffered, since emit
then builds every batch
+ // from the key queues alone.
+ long keptRows;
+ long phase2Bytes;
+ PageReadStore dataPages = null;
+ if (nonKeyColumns == null && spliceCurrentRowGroup) {
+ keptRows = finalRowCount;
+ phase2Bytes = 0L;
+ } else {
+ lateMatReader.setRequestedSchema(
+ spliceCurrentRowGroup ? nonKeyColumns : requestedColumns);
+ // Reading a strict subset of a block's rows needs a Parquet offset
index, and parquet
+ // enforces that itself: it resolves every requested column's offset
index before reading
+ // anything, and a column without one makes its column index store
throw
+ // MissingOffsetIndexException. Files written before parquet-mr 1.11,
or by a writer that
+ // omits the page index (pyarrow's `write_table` defaults to
`write_page_index=False`), have
+ // none. The filter is then given up for this row group and the read
retried over
+ // `pushedFilterRanges`, which is what a plain scan reads. That retry
cannot hit the same
+ // wall: a store missing one column's offset index reports no column
index either, so
+ // `getRowRanges` could not have narrowed anything and the ranges
cover the whole block.
+ //
+ // Nothing is checked up front, so a file with no page index still
reads with the filter
+ // applied wherever the filter keeps a row group whole
(`readFilteredRowGroup` degrades to a
+ // plain read when the ranges cover the block) or rejects one whole.
+ try {
+ dataPages = lateMatReader.readFilteredRowGroup(blockIdx,
finalRanges);
+ } catch (MissingOffsetIndexException e) {
+ LOG.warn("Not applying the storage filter to {}: reading part of a
row group needs a "
+ + "Parquet offset index, and this file was written without a
page index for at least "
+ + "one projected column", e, MDC.of(LogKeys.PATH,
lateMatReader.getFile()));
+ fileHasNoOffsetIndex = true;
+ giveUpFilter();
+ finalRanges = pushedFilterRanges;
+ finalRowCount = baselineRows;
+ lateMatReader.setRequestedSchema(requestedColumns);
+ dataPages = lateMatReader.readFilteredRowGroup(blockIdx,
finalRanges);
+ }
+ if (dataPages == null) {
+ // Unreachable: readFilteredRowGroup returns null only for an empty
block or empty ranges,
+ // both excluded above. Match phase 1 and fail with a message rather
than an NPE.
+ throw new IllegalStateException(
+ "No data pages for row group " + blockIdx + " despite " +
finalRowCount
+ + " rows to read");
+ }
+ keptRows = dataPages.getRowCount();
+ // Nothing is computed for a row group whose filter was given up: it
read what a plain scan
+ // reads, so the answer is a certain zero. `needBytes`, not just
`bytesAvoidedPf != null`,
+ // because that is what built `blockChunks`.
+ if (needBytes && bytesAvoidedPf != null && !filterGivenUp) {
+ phase2Bytes = compressedBytesForRowRanges(lateMatReader, blockIdx,
blockChunks,
+ nonKeyColumns, finalRanges, finalRowCount);
+ if (!spliceCurrentRowGroup) {
+ // This row group gave splicing up, so phase 2 read the key
columns a second time. The
+ // baseline counts them once, in phase 1, so the extra read is a
cost against it.
+ phase2Bytes += compressedBytesForRowRanges(lateMatReader,
blockIdx, blockChunks,
+ keyOnlyColumns, finalRanges, finalRowCount);
+ }
+ } else {
+ phase2Bytes = 0L;
+ }
+ }
+ long filteredRows = baselineRows - keptRows;
+ SQLMetric rowsExcludedWithinRg = m.rowsExcludedWithinRowGroup();
+ if (rowsExcludedWithinRg != null && filteredRows > 0)
rowsExcludedWithinRg.add(filteredRows);
+ if (bytesAvoidedPf != null && !filterGivenUp) {
+ // `SQLMetric.add` ignores a negative value, so a row group that read
more than the baseline
+ // after giving splicing up contributes nothing rather than
subtracting.
+ bytesAvoidedPf.add(nonKeyBaselineBytes - phase2Bytes);
+ }
+
+ if (dataPages != null) {
+ if (rowIndexGenerator != null) {
+ rowIndexGenerator.initFromPageReadStore(dataPages);
+ }
+ for (int i = 0; i < columnVectors.length; i++) {
+ if (spliceCurrentRowGroup && isKeyTopLevel[i]) {
+ // Key columns are sourced from the queues during emit; skip
phase-2 reader init.
+ continue;
+ }
+ initColumnReader(dataPages, columnVectors[i]);
+ }
+ }
+ totalCountLoadedSoFar += keptRows;
+ return;
+ }
+ hitEndOfData = true;
+ }
+
+
+ /** Counts a row group whose data columns the filter kept the reader from
touching at all. */
+ private static void recordRowGroupSkipped(
+ StorageFilterMetrics m, long excludedRows, long avoidedBytes) {
+ SQLMetric rgSkipped = m.rowGroupsSkipped();
+ if (rgSkipped != null) rgSkipped.add(1L);
+ SQLMetric rowsExcluded = m.rowsExcludedByRowGroup();
+ if (rowsExcluded != null) rowsExcluded.add(excludedRows);
+ SQLMetric bytesAvoided = m.bytesAvoidedByRowGroup();
+ if (bytesAvoided != null) bytesAvoided.add(avoidedBytes);
+ }
+
+ /**
+ * Counts a file the filter rejects whole, which happens when every key
column is missing from it
+ * and the predicate is constant-false for the value the reader would have
materialized. Every row
+ * group counts as skipped and every projected byte as avoided, which is
what the counters mean
+ * for a row group the filter empties.
+ */
+ private void recordFileSkipped() {
+ StorageFilterMetrics m = storageFilter.metrics();
+ boolean needBytes = m.bytesAvoidedByRowGroup() != null;
+ if (m.rowGroupsSkipped() == null && m.rowsExcludedByRowGroup() == null &&
!needBytes) return;
+ List<ColumnDescriptor> projected = requestedSchema.getColumns();
+ List<BlockMetaData> blocks = lateMatReader.getRowGroups();
+ for (int blockIdx = 0; blockIdx < blocks.size(); blockIdx++) {
+ // Measured against the rows the pushed data filter kept, which is the
baseline every other
+ // skip path uses: the rows its column index already excluded were never
this filter's to
+ // save. `getRowRanges` is a cache hit whenever the two can differ,
because
+ // `getFilteredRecordCount()` at initialize resolved every block's
ranges then. It has to be
+ // guarded the same way phase 0 guards it, since it consults the pushed
filter but not the
+ // conf that turns column-index filtering off.
+ long blockRowCount = blocks.get(blockIdx).getRowCount();
+ if (blockRowCount == 0) continue;
+ RowRanges blockRanges = useColumnIndexFilter
+ ? lateMatReader.getRowRanges(blockIdx)
+ : RowRanges.createSingle(blockRowCount);
+ long survivingRows = blockRanges.rowCount();
+ if (survivingRows == 0) continue;
+ // The key columns are missing from this file, so they contribute
nothing to the walk, and the
+ // whole projection is what a plain read would have transferred.
+ long avoidedBytes = needBytes
+ ? compressedBytesForRowRanges(lateMatReader, blockIdx,
+ chunksByPath(lateMatReader, blockIdx), projected, blockRanges,
survivingRows)
+ : 0L;
+ recordRowGroupSkipped(m, survivingRows, avoidedBytes);
+ }
+ }
+
+ /**
+ * The block's column chunks by path, built once per row group and shared by
the byte-metric calls
+ * that consume it, since {@link BlockMetaData} offers no lookup of its own.
+ */
+ private static Map<ColumnPath, ColumnChunkMetaData> chunksByPath(
+ ParquetFileReader reader, int blockIndex) {
+ Map<ColumnPath, ColumnChunkMetaData> chunks = new HashMap<>();
+ for (ColumnChunkMetaData chunk :
reader.getRowGroups().get(blockIndex).getColumns()) {
+ chunks.put(chunk.getPath(), chunk);
+ }
+ return chunks;
+ }
+
+ /**
+ * Compressed bytes the reader transfers for the given leaf {@code columns}
when it reads exactly
+ * {@code rowRanges} of the given block. Page headers and the dictionary
page are included, since
+ * both are read whenever any page of a chunk is read. {@code rowRangeCount}
is
+ * {@code rowRanges.rowCount()}, passed in because that walks every range
and the caller has it.
+ *
+ * <p>Two sources, chosen so this never causes IO of its own:
+ * <ul>
+ * <li>{@code rowRanges} covers the whole block: the answer is the sum of
the chunks'
+ * {@code getTotalSize()}, which is already in the footer. This is the
case that matters:
+ * whenever nothing else has built the block's {@link
ColumnIndexStore}, {@code rowRanges}
+ * is necessarily the whole block, because a narrower range can only
come from column-index
+ * filtering, which builds the store as a side effect.
+ * <li>{@code rowRanges} is a strict subset: walk the offset index, as
parquet's own read path
+ * does, and add the dictionary page the way {@code
calculateOffsetRanges} does. The store
+ * is guaranteed to exist here, so the walk is pure metadata
arithmetic. For the ranges the
+ * storage filter narrowed, which column-index filtering had no hand
in, that guarantee is
+ * an ordering one: phase 2's own read of those ranges built the store
first.
+ * </ul>
+ *
+ * <p>Columns absent from this physical file (schema evolution) contribute
nothing, which is
+ * correct: the reader transfers nothing for them.
+ */
+ private static long compressedBytesForRowRanges(
+ ParquetFileReader reader,
+ int blockIndex,
+ Map<ColumnPath, ColumnChunkMetaData> chunks,
+ List<ColumnDescriptor> columns,
+ RowRanges rowRanges,
+ long rowRangeCount) {
+ if (columns == null || columns.isEmpty() || rowRangeCount == 0) {
+ return 0L;
+ }
+ long blockRowCount = reader.getRowGroups().get(blockIndex).getRowCount();
+ boolean wholeBlock = rowRangeCount == blockRowCount;
+ ColumnIndexStore ciStore = wholeBlock ? null :
reader.getColumnIndexStore(blockIndex);
+ long total = 0L;
+ for (ColumnDescriptor column : columns) {
+ ColumnPath path = ColumnPath.get(column.getPath());
+ ColumnChunkMetaData chunk = chunks.get(path);
+ if (chunk == null) {
+ // Column is in the (clipped) requested schema but not in this file.
+ continue;
+ }
+ if (wholeBlock) {
+ total += chunk.getTotalSize();
+ continue;
+ }
+ OffsetIndex offsetIndex;
+ try {
+ offsetIndex = ciStore.getOffsetIndex(path);
+ } catch (MissingOffsetIndexException e) {
+ continue;
+ }
+ if (offsetIndex == null) {
+ continue;
+ }
+ // The dictionary page is read whenever any data page of the chunk is,
so count it here the
+ // same way parquet's ColumnIndexFilterUtils.calculateOffsetRanges does.
+ total += dictionaryPageSize(chunk);
+ int pageCount = offsetIndex.getPageCount();
+ for (int i = 0; i < pageCount; i++) {
+ long from = offsetIndex.getFirstRowIndex(i);
+ long to = offsetIndex.getLastRowIndex(i, blockRowCount);
+ if (rowRanges.isOverlapping(from, to)) {
+ total += offsetIndex.getCompressedPageSize(i);
+ }
+ }
+ }
+ return total;
+ }
+
+ /**
+ * Compressed size of a chunk's dictionary page, or 0 if it has none.
+ * {@link ColumnChunkMetaData#getStartingPos()} already resolves to the
dictionary page offset
+ * when there is a valid one, so the gap up to the first data page is
exactly the dictionary page.
+ */
+ private static long dictionaryPageSize(ColumnChunkMetaData chunk) {
+ long startingPos = chunk.getStartingPos();
+ long firstDataPageOffset = chunk.getFirstDataPageOffset();
+ return startingPos < firstDataPageOffset ? firstDataPageOffset -
startingPos : 0L;
+ }
+
+ /**
+ * Evaluates the storage filter over every row of a key-only {@link
PageReadStore}, in
+ * capacity-sized chunks, and returns the surviving rows as {@link
RowRanges} in block-row
+ * coordinates. The result is a subset of {@code pushedFilterRanges}: rows
outside it were never
+ * read.
+ *
+ * <p>Each survivor's key values are appended to {@link
#currentKeyAccumulators} for the emit path
+ * to splice, until the buffer passes its cap. From there the row group is
evaluated without
+ * buffering and {@link #spliceCurrentRowGroup} is false, so its phase 2
reads the key columns
+ * again along with everything else.
+ *
+ * <p>Returns null once the budget makes the reader give the filter up for
this row group: the
+ * ranges built so far are then incomplete, and the caller reads the row
group the plain way.
+ */
+ private RowRanges evaluateStorageFilter(
+ PageReadStore keyPages,
+ RowRanges pushedFilterRanges) throws IOException {
+ ensureKeyScratchAllocated();
+ VectorizedColumnReader[] readers = new
VectorizedColumnReader[keyDescriptors.length];
+ for (int i = 0; i < readers.length; i++) {
+ readers[i] = new VectorizedColumnReader(
+ keyDescriptors[i], keyRequired[i], keyPages, convertTz,
datetimeRebaseMode,
+ datetimeRebaseTz, int96RebaseMode, int96RebaseTz, writerVersion);
+ }
+ ensureCurrentKeyAccumulatorsAllocated();
+
+ PrimitiveIterator.OfLong rowIndexIter = pushedFilterRanges.iterator();
+ RowRanges.Builder finalRangesBuilder = RowRanges.builder();
+ survivorRangeCount = 0L;
+ long previousSurvivor = -2L;
+ // Recomputed rather than taken from the caller: a count that disagreed
with this iterator would
+ // silently drop surviving rows, and no post-scan Filter is left to catch
that.
+ long remaining = pushedFilterRanges.rowCount();
+ boolean accumulate = true;
+ while (remaining > 0) {
+ int num = (int) Math.min((long) capacity, remaining);
+ for (int i = 0; i < keyScratchVectors.length; i++) {
+ keyScratchVectors[i].reset();
+ readers[i].readBatch(num, keyScratchVectors[i], null, null);
+ }
+ keyScratchBatch.setNumRows(num);
+ for (int r = 0; r < num; r++) {
+ long blockRow = rowIndexIter.nextLong();
+ if (storageFilter.test(keyScratchBatch.getRow(r))) {
+ finalRangesBuilder.addSelectedRow(blockRow);
+ if (blockRow != previousSurvivor + 1) survivorRangeCount++;
+ previousSurvivor = blockRow;
+ if (accumulate) {
+ accumulate = appendSurvivorRowToAccumulators(r);
Review Comment:
Done as described. One check, in the per-survivor loop, over `splicedBytes +
rowRangeStateBytes(survivorRangeCount)`, and it runs whether or not splicing is
still on. Crossing it releases the key accumulators first, since that is the
cheaper concession, and gives the filter up when the ranges alone exceed the
cap, which returns immediately so the rest of the range list is never built.
Nothing partial is left observable: `giveUpFilter` closes every buffered
vector, and the caller then reads the row group over `pushedFilterRanges`,
which is what a plain read reads.
The conf doc is aligned with it: two pools, both growing with the survivor
count, examined after every surviving row, and what is counted is the values
and their per-row overhead rather than the backing arrays.
Tests cover both boundaries, a row group that crosses the cap inside the
loop and one whose ranges alone do, plus a file whose first row group splices
while a later one does not, since the batch is one object and its key slots
have to go back to the vectors phase 2 reads into.
The range half also got much cheaper at the root, from @dongjoon-hyun's
comment on the same loop: `ParquetReadState` no longer copies the range list
per column reader. c29625f
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]