Skip to content

Commit 96d21b3

Browse files
authored
refactor: Reuse RowGroupPageIndexReader across columns to improve page-level predicate pushdown performance (alibaba#316)
1 parent 2c4019e commit 96d21b3

2 files changed

Lines changed: 16 additions & 11 deletions

File tree

src/paimon/format/parquet/page_filtered_row_group_reader.cpp

Lines changed: 15 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -138,9 +138,10 @@ std::pair<RowRanges, int64_t> PageFilteredRowGroupReader::ComputeCompressedRowRa
138138
Result<std::shared_ptr<arrow::ChunkedArray>> PageFilteredRowGroupReader::ReadFilteredColumn(
139139
const std::shared_ptr<::parquet::RowGroupReader>& row_group_reader,
140140
::parquet::ParquetFileReader* parquet_reader,
141-
const std::shared_ptr<::parquet::PageIndexReader>& page_index_reader, int32_t row_group_index,
142-
int32_t column_index, const RowRanges& row_ranges, const std::shared_ptr<arrow::Field>& field,
143-
int64_t row_group_row_count, ::arrow::MemoryPool* pool) {
141+
const std::shared_ptr<::parquet::RowGroupPageIndexReader>& rg_page_index_reader,
142+
int32_t row_group_index, int32_t column_index, const RowRanges& row_ranges,
143+
const std::shared_ptr<arrow::Field>& field, int64_t row_group_row_count,
144+
::arrow::MemoryPool* pool) {
144145
auto file_metadata = parquet_reader->metadata();
145146
const auto* col_descriptor = file_metadata->schema()->Column(column_index);
146147

@@ -149,11 +150,8 @@ Result<std::shared_ptr<arrow::ChunkedArray>> PageFilteredRowGroupReader::ReadFil
149150
int64_t effective_row_count = row_group_row_count;
150151

151152
std::shared_ptr<::parquet::OffsetIndex> offset_index;
152-
if (page_index_reader) {
153-
auto rg_page_index_reader = page_index_reader->RowGroup(row_group_index);
154-
if (rg_page_index_reader) {
155-
offset_index = rg_page_index_reader->GetOffsetIndex(column_index);
156-
}
153+
if (rg_page_index_reader) {
154+
offset_index = rg_page_index_reader->GetOffsetIndex(column_index);
157155
}
158156

159157
auto page_reader = row_group_reader->GetColumnPageReader(column_index);
@@ -263,15 +261,22 @@ Result<std::unique_ptr<arrow::RecordBatchReader>> PageFilteredRowGroupReader::Re
263261
int64_t row_group_row_count = rg_metadata->num_rows();
264262
auto page_index_reader = parquet_reader->GetPageIndexReader();
265263

264+
// reuse RowGroupPageIndexReader for multiple columns in the same row group to avoid redundant
265+
// metadata reads
266+
std::shared_ptr<::parquet::RowGroupPageIndexReader> rg_page_index_reader;
267+
if (page_index_reader) {
268+
rg_page_index_reader = page_index_reader->RowGroup(row_group_index);
269+
}
270+
266271
// Read each column with page filtering
267272
std::vector<std::shared_ptr<arrow::ChunkedArray>> columns;
268273
columns.reserve(column_indices.size());
269274

270275
for (size_t i = 0; i < column_indices.size(); ++i) {
271276
PAIMON_ASSIGN_OR_RAISE(
272277
std::shared_ptr<arrow::ChunkedArray> chunked_array,
273-
ReadFilteredColumn(row_group_reader, parquet_reader, page_index_reader, row_group_index,
274-
column_indices[i], row_ranges,
278+
ReadFilteredColumn(row_group_reader, parquet_reader, rg_page_index_reader,
279+
row_group_index, column_indices[i], row_ranges,
275280
arrow_schema->field(static_cast<int>(i)), row_group_row_count,
276281
pool));
277282

src/paimon/format/parquet/page_filtered_row_group_reader.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -90,7 +90,7 @@ class PageFilteredRowGroupReader {
9090
static Result<std::shared_ptr<arrow::ChunkedArray>> ReadFilteredColumn(
9191
const std::shared_ptr<::parquet::RowGroupReader>& row_group_reader,
9292
::parquet::ParquetFileReader* parquet_reader,
93-
const std::shared_ptr<::parquet::PageIndexReader>& page_index_reader,
93+
const std::shared_ptr<::parquet::RowGroupPageIndexReader>& rg_page_index_reader,
9494
int32_t row_group_index, int32_t column_index, const RowRanges& row_ranges,
9595
const std::shared_ptr<arrow::Field>& field, int64_t row_group_row_count,
9696
::arrow::MemoryPool* pool);

0 commit comments

Comments
 (0)