Skip to content

Commit 8a97335

Browse files
zhf999SteNicholas
authored andcommitted
fix(parquet): reuse RowGroupPageIndexReader in FileReaderWrapper layer (#166)
1 parent 9d19960 commit 8a97335

10 files changed

Lines changed: 129 additions & 92 deletions

src/paimon/format/parquet/column_index_filter.cpp

Lines changed: 3 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -38,15 +38,9 @@ namespace paimon::parquet {
3838

3939
Result<RowRanges> ColumnIndexFilter::CalculateRowRanges(
4040
const std::shared_ptr<Predicate>& predicate,
41-
const std::shared_ptr<::parquet::PageIndexReader>& page_index_reader,
42-
const std::map<std::string, int32_t>& column_name_to_index, int32_t row_group_index,
43-
int64_t row_group_row_count) {
44-
if (!predicate || !page_index_reader) {
45-
return RowRanges::CreateSingle(row_group_row_count);
46-
}
47-
48-
auto rg_page_index_reader = page_index_reader->RowGroup(row_group_index);
49-
if (!rg_page_index_reader) {
41+
const std::shared_ptr<::parquet::RowGroupPageIndexReader>& rg_page_index_reader,
42+
const std::map<std::string, int32_t>& column_name_to_index, int64_t row_group_row_count) {
43+
if (!predicate || !rg_page_index_reader) {
5044
return RowRanges::CreateSingle(row_group_row_count);
5145
}
5246

src/paimon/format/parquet/column_index_filter.h

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -57,16 +57,14 @@ class ColumnIndexFilter {
5757

5858
/// Calculate row ranges based on predicate and column indices.
5959
/// @param predicate The predicate to evaluate.
60-
/// @param page_index_reader The page index reader for the file.
60+
/// @param rg_page_index_reader The page index reader of target row group for the file.
6161
/// @param column_name_to_index Map from column name to column index.
62-
/// @param row_group_index The row group index to filter.
6362
/// @param row_group_row_count The number of rows in the row group.
6463
/// @return RowRanges that may contain matching rows.
6564
static Result<RowRanges> CalculateRowRanges(
6665
const std::shared_ptr<Predicate>& predicate,
67-
const std::shared_ptr<::parquet::PageIndexReader>& page_index_reader,
68-
const std::map<std::string, int32_t>& column_name_to_index, int32_t row_group_index,
69-
int64_t row_group_row_count);
66+
const std::shared_ptr<::parquet::RowGroupPageIndexReader>& rg_page_index_reader,
67+
const std::map<std::string, int32_t>& column_name_to_index, int64_t row_group_row_count);
7068

7169
private:
7270
/// Visit a predicate and calculate row ranges.

src/paimon/format/parquet/column_index_filter_test.cpp

Lines changed: 24 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -305,9 +305,8 @@ class ColumnIndexFilterTest : public ::testing::Test {
305305
}
306306

307307
Result<RowRanges> Filter(const std::shared_ptr<Predicate>& predicate) {
308-
return ColumnIndexFilter::CalculateRowRanges(predicate, page_index_reader_,
309-
column_name_to_index_, /*row_group_index=*/0,
310-
row_group_row_count_);
308+
return ColumnIndexFilter::CalculateRowRanges(predicate, page_index_reader_->RowGroup(0),
309+
column_name_to_index_, row_group_row_count_);
311310
}
312311

313312
std::shared_ptr<arrow::MemoryPool> arrow_pool_;
@@ -553,19 +552,19 @@ TEST_F(ColumnIndexFilterTest, SignedZeroUsesJavaOrderForFloatingPointPages) {
553552
auto less_negative_zero = PredicateBuilder::LessThan(
554553
/*field_index=*/0, /*field_name=*/"value", field_type,
555554
field_type == FieldType::FLOAT ? Literal(-0.0f) : Literal(-0.0));
556-
ASSERT_OK_AND_ASSIGN(
557-
auto ranges, ColumnIndexFilter::CalculateRowRanges(
558-
less_negative_zero, page_index_reader, {{"value", 0}},
559-
/*row_group_index=*/0, reader->metadata()->RowGroup(0)->num_rows()));
555+
ASSERT_OK_AND_ASSIGN(auto ranges,
556+
ColumnIndexFilter::CalculateRowRanges(
557+
less_negative_zero, page_index_reader->RowGroup(0), {{"value", 0}},
558+
reader->metadata()->RowGroup(0)->num_rows()));
560559
ASSERT_TRUE(ranges.IsEmpty()) << "field type: " << static_cast<int32_t>(field_type);
561560

562561
auto less_positive_zero = PredicateBuilder::LessThan(
563562
/*field_index=*/0, /*field_name=*/"value", field_type,
564563
field_type == FieldType::FLOAT ? Literal(0.0f) : Literal(0.0));
565-
ASSERT_OK_AND_ASSIGN(
566-
ranges, ColumnIndexFilter::CalculateRowRanges(
567-
less_positive_zero, page_index_reader, {{"value", 0}},
568-
/*row_group_index=*/0, reader->metadata()->RowGroup(0)->num_rows()));
564+
ASSERT_OK_AND_ASSIGN(ranges,
565+
ColumnIndexFilter::CalculateRowRanges(
566+
less_positive_zero, page_index_reader->RowGroup(0), {{"value", 0}},
567+
reader->metadata()->RowGroup(0)->num_rows()));
569568
ASSERT_EQ(20, ranges.RowCount());
570569
ASSERT_EQ(1, ranges.GetRanges().size());
571570
ASSERT_EQ(0, ranges.GetRanges()[0].from);
@@ -576,35 +575,35 @@ TEST_F(ColumnIndexFilterTest, SignedZeroUsesJavaOrderForFloatingPointPages) {
576575
field_type == FieldType::FLOAT ? Literal(-0.0f) : Literal(-0.0));
577576
ASSERT_OK_AND_ASSIGN(
578577
ranges, ColumnIndexFilter::CalculateRowRanges(
579-
greater_negative_zero, page_index_reader, {{"value", 0}},
580-
/*row_group_index=*/0, reader->metadata()->RowGroup(0)->num_rows()));
578+
greater_negative_zero, page_index_reader->RowGroup(0), {{"value", 0}},
579+
reader->metadata()->RowGroup(0)->num_rows()));
581580
ASSERT_EQ(30, ranges.RowCount());
582581

583582
auto not_equal_negative_zero = PredicateBuilder::NotEqual(
584583
/*field_index=*/0, /*field_name=*/"value", field_type,
585584
field_type == FieldType::FLOAT ? Literal(-0.0f) : Literal(-0.0));
586585
ASSERT_OK_AND_ASSIGN(
587586
ranges, ColumnIndexFilter::CalculateRowRanges(
588-
not_equal_negative_zero, page_index_reader, {{"value", 0}},
589-
/*row_group_index=*/0, reader->metadata()->RowGroup(0)->num_rows()));
587+
not_equal_negative_zero, page_index_reader->RowGroup(0), {{"value", 0}},
588+
reader->metadata()->RowGroup(0)->num_rows()));
590589
ASSERT_EQ(30, ranges.RowCount());
591590

592591
auto greater_finite = PredicateBuilder::GreaterThan(
593592
/*field_index=*/0, /*field_name=*/"value", field_type,
594593
field_type == FieldType::FLOAT ? Literal(2.0f) : Literal(2.0));
595-
ASSERT_OK_AND_ASSIGN(
596-
ranges, ColumnIndexFilter::CalculateRowRanges(
597-
greater_finite, page_index_reader, {{"value", 0}},
598-
/*row_group_index=*/0, reader->metadata()->RowGroup(0)->num_rows()));
594+
ASSERT_OK_AND_ASSIGN(ranges,
595+
ColumnIndexFilter::CalculateRowRanges(
596+
greater_finite, page_index_reader->RowGroup(0), {{"value", 0}},
597+
reader->metadata()->RowGroup(0)->num_rows()));
599598
ASSERT_TRUE(ranges.IsEmpty());
600599

601600
auto greater_between_pages = PredicateBuilder::GreaterThan(
602601
/*field_index=*/0, /*field_name=*/"value", field_type,
603602
field_type == FieldType::FLOAT ? Literal(0.5f) : Literal(0.5));
604603
ASSERT_OK_AND_ASSIGN(
605604
ranges, ColumnIndexFilter::CalculateRowRanges(
606-
greater_between_pages, page_index_reader, {{"value", 0}},
607-
/*row_group_index=*/0, reader->metadata()->RowGroup(0)->num_rows()));
605+
greater_between_pages, page_index_reader->RowGroup(0), {{"value", 0}},
606+
reader->metadata()->RowGroup(0)->num_rows()));
608607
ASSERT_EQ(10, ranges.RowCount());
609608
ASSERT_EQ(1, ranges.GetRanges().size());
610609
ASSERT_EQ(20, ranges.GetRanges()[0].from);
@@ -613,10 +612,10 @@ TEST_F(ColumnIndexFilterTest, SignedZeroUsesJavaOrderForFloatingPointPages) {
613612
auto equal_finite = PredicateBuilder::Equal(
614613
/*field_index=*/0, /*field_name=*/"value", field_type,
615614
field_type == FieldType::FLOAT ? Literal(2.0f) : Literal(2.0));
616-
ASSERT_OK_AND_ASSIGN(
617-
ranges, ColumnIndexFilter::CalculateRowRanges(
618-
equal_finite, page_index_reader, {{"value", 0}},
619-
/*row_group_index=*/0, reader->metadata()->RowGroup(0)->num_rows()));
615+
ASSERT_OK_AND_ASSIGN(ranges,
616+
ColumnIndexFilter::CalculateRowRanges(
617+
equal_finite, page_index_reader->RowGroup(0), {{"value", 0}},
618+
reader->metadata()->RowGroup(0)->num_rows()));
620619
ASSERT_TRUE(ranges.IsEmpty());
621620
}
622621
}

src/paimon/format/parquet/file_reader_wrapper.cpp

Lines changed: 31 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -232,15 +232,18 @@ Result<std::shared_ptr<arrow::RecordBatch>> FileReaderWrapper::NextPageFiltered(
232232
// Construct the per-RG streaming reader on demand.
233233
if (!current_page_filtered_reader_) {
234234
const auto& target_rg = target_row_groups_[current_row_group_idx_];
235+
auto row_group_page_index_reader = GetRowGroupPageIndexReader(rg_id);
235236
auto page_ranges = PageFilteredRowGroupReader::ComputePageRanges(
236-
target_rg, target_column_indices_, file_reader_->parquet_reader());
237+
target_rg, target_column_indices_, row_group_page_index_reader,
238+
file_reader_->parquet_reader());
237239
bool pre_buffered = !prebuffered_ranges_.empty();
238240
int64_t max_chunksize = batch_size_ > 0 ? batch_size_ : std::numeric_limits<int64_t>::max();
239241
PAIMON_ASSIGN_OR_RAISE(
240242
current_page_filtered_reader_,
241243
PageFilteredRowGroupReader::ReadFilteredRowGroup(
242244
target_rg, target_column_indices_, file_reader_->properties().cache_options(),
243-
pre_buffered, page_ranges, max_chunksize, pool_, file_reader_.get()));
245+
pre_buffered, page_ranges, max_chunksize, row_group_page_index_reader, pool_,
246+
file_reader_.get()));
244247
current_filtered_row_ranges_ = target_rg.GetRowRanges();
245248
current_filtered_rg_start_ = all_row_group_ranges_[rg_id].first;
246249
filtered_global_offset_ = 0;
@@ -298,6 +301,27 @@ Result<std::shared_ptr<arrow::RecordBatch>> FileReaderWrapper::NextFullyMatched(
298301
return record_batch;
299302
}
300303

304+
std::shared_ptr<::parquet::RowGroupPageIndexReader> FileReaderWrapper::GetRowGroupPageIndexReader(
305+
int32_t row_group_index) {
306+
auto cached = row_group_page_index_readers_.find(row_group_index);
307+
if (cached != row_group_page_index_readers_.end()) {
308+
return cached->second;
309+
}
310+
311+
std::shared_ptr<::parquet::RowGroupPageIndexReader> row_group_page_index_reader;
312+
auto page_index_reader = GetPageIndexReader();
313+
if (page_index_reader) {
314+
row_group_page_index_reader = page_index_reader->RowGroup(row_group_index);
315+
}
316+
317+
// To avoid OOM, limit the number of row group page index readers cached in memory.
318+
constexpr int32_t kMaxRowGroupPageIndexReaders = 1024;
319+
if (row_group_page_index_readers_.size() < kMaxRowGroupPageIndexReaders) {
320+
row_group_page_index_readers_.emplace(row_group_index, row_group_page_index_reader);
321+
}
322+
return row_group_page_index_reader;
323+
}
324+
301325
Result<std::shared_ptr<arrow::RecordBatch>> FileReaderWrapper::Next() {
302326
try {
303327
if (PAIMON_UNLIKELY(!reader_initialized_)) {
@@ -357,8 +381,9 @@ std::vector<::arrow::io::ReadRange> FileReaderWrapper::CollectPreBufferRanges(
357381

358382
if (trg.IsPartiallyMatched()) {
359383
// Page-filtered RGs: only matching page byte ranges.
384+
auto row_group_page_index_reader = GetRowGroupPageIndexReader(trg.GetRowGroupIndex());
360385
auto page_ranges = PageFilteredRowGroupReader::ComputePageRanges(
361-
trg, column_indices, file_reader_->parquet_reader());
386+
trg, column_indices, row_group_page_index_reader, file_reader_->parquet_reader());
362387
ranges.insert(ranges.end(), std::make_move_iterator(page_ranges.begin()),
363388
std::make_move_iterator(page_ranges.end()));
364389
} else {
@@ -500,13 +525,9 @@ Result<RowRanges> FileReaderWrapper::CalculateFilteredRowRanges(
500525
return RowRanges::CreateSingle(row_count);
501526
}
502527

503-
auto page_index_reader = GetPageIndexReader();
504-
if (!page_index_reader) {
505-
return RowRanges::CreateSingle(row_count);
506-
}
507-
508-
return ColumnIndexFilter::CalculateRowRanges(
509-
predicate, page_index_reader, column_name_to_index, row_group_index, row_count);
528+
return ColumnIndexFilter::CalculateRowRanges(predicate,
529+
GetRowGroupPageIndexReader(row_group_index),
530+
column_name_to_index, row_count);
510531
}
511532
PAIMON_PARQUET_CATCH_AND_RETURN_STATUS("FileReaderWrapper::CalculateFilteredRowRanges")
512533
}

src/paimon/format/parquet/file_reader_wrapper.h

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -143,6 +143,10 @@ class FileReaderWrapper {
143143
int32_t row_group_index, const std::shared_ptr<Predicate>& predicate,
144144
const std::map<std::string, int32_t>& column_name_to_index);
145145

146+
/// Get or create the page index reader for a row group.
147+
std::shared_ptr<::parquet::RowGroupPageIndexReader> GetRowGroupPageIndexReader(
148+
int32_t row_group_index);
149+
146150
private:
147151
FileReaderWrapper(std::unique_ptr<::parquet::arrow::FileReader>&& file_reader,
148152
const std::vector<std::pair<uint64_t, uint64_t>>& all_row_group_ranges,
@@ -196,6 +200,11 @@ class FileReaderWrapper {
196200

197201
// Track pre-buffered ranges so we can wait on destruction
198202
std::vector<::arrow::io::ReadRange> prebuffered_ranges_;
203+
204+
// Arrow caches the file-level PageIndexReader, but RowGroup() creates a new reader each time.
205+
// Keep one reader per row group so its page-index buffers are shared by all read stages.
206+
std::map<int32_t, std::shared_ptr<::parquet::RowGroupPageIndexReader>>
207+
row_group_page_index_readers_;
199208
};
200209

201210
} // namespace paimon::parquet

src/paimon/format/parquet/page_filtered_row_group_reader.cpp

Lines changed: 6 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -373,6 +373,7 @@ Result<std::unique_ptr<arrow::RecordBatchReader>> PageFilteredRowGroupReader::Re
373373
const TargetRowGroup& target_row_group, const std::vector<int32_t>& column_indices,
374374
const ::arrow::io::CacheOptions& cache_options, bool pre_buffered,
375375
const std::vector<::arrow::io::ReadRange>& page_ranges, int64_t max_chunksize,
376+
const std::shared_ptr<::parquet::RowGroupPageIndexReader>& row_group_page_index_reader,
376377
std::shared_ptr<::arrow::MemoryPool> pool, ::parquet::arrow::FileReader* arrow_file_reader) {
377378
auto parquet_reader = arrow_file_reader->parquet_reader();
378379
const auto& row_ranges = target_row_group.GetRowRanges();
@@ -388,14 +389,6 @@ Result<std::unique_ptr<arrow::RecordBatchReader>> PageFilteredRowGroupReader::Re
388389
auto rg_metadata = parquet_reader->metadata()->RowGroup(row_group_index);
389390
int64_t row_group_row_count = rg_metadata->num_rows();
390391

391-
// reuse RowGroupPageIndexReader for multiple columns in the same row group to avoid redundant
392-
// metadata reads
393-
std::shared_ptr<::parquet::RowGroupPageIndexReader> rg_page_index_reader;
394-
auto page_index_reader = parquet_reader->GetPageIndexReader();
395-
if (page_index_reader) {
396-
rg_page_index_reader = page_index_reader->RowGroup(row_group_index);
397-
}
398-
399392
const auto& manifest = arrow_file_reader->manifest();
400393
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(
401394
std::vector<int> field_indices,
@@ -407,8 +400,8 @@ Result<std::unique_ptr<arrow::RecordBatchReader>> PageFilteredRowGroupReader::Re
407400
for (int field_idx : field_indices) {
408401
PAIMON_ASSIGN_OR_RAISE(
409402
std::shared_ptr<arrow::ChunkedArray> chunked_array,
410-
ReadFilteredField(rg_page_index_reader, row_group_index, field_idx, column_indices,
411-
row_ranges, row_group_row_count, arrow_file_reader));
403+
ReadFilteredField(row_group_page_index_reader, row_group_index, field_idx,
404+
column_indices, row_ranges, row_group_row_count, arrow_file_reader));
412405

413406
if (chunked_array->length() != expected_rows) {
414407
return Status::Invalid(
@@ -434,6 +427,7 @@ Result<std::unique_ptr<arrow::RecordBatchReader>> PageFilteredRowGroupReader::Re
434427

435428
std::vector<::arrow::io::ReadRange> PageFilteredRowGroupReader::ComputePageRanges(
436429
const TargetRowGroup& target_row_group, const std::vector<int32_t>& column_indices,
430+
const std::shared_ptr<::parquet::RowGroupPageIndexReader>& row_group_page_index_reader,
437431
::parquet::ParquetFileReader* parquet_reader) {
438432
int32_t row_group_index = target_row_group.GetRowGroupIndex();
439433
const auto& row_ranges = target_row_group.GetRowRanges();
@@ -447,21 +441,15 @@ std::vector<::arrow::io::ReadRange> PageFilteredRowGroupReader::ComputePageRange
447441
auto rg_metadata = file_metadata->RowGroup(row_group_index);
448442
int64_t row_group_row_count = rg_metadata->num_rows();
449443

450-
auto page_index_reader = parquet_reader->GetPageIndexReader();
451-
std::shared_ptr<::parquet::RowGroupPageIndexReader> rg_page_index_reader;
452-
if (page_index_reader) {
453-
rg_page_index_reader = page_index_reader->RowGroup(row_group_index);
454-
}
455-
456444
for (int32_t col_idx : column_indices) {
457445
auto col_chunk = rg_metadata->ColumnChunk(col_idx);
458446
const int64_t column_chunk_offset = GetColumnChunkOffset(*col_chunk);
459447
const int64_t column_chunk_compressed_size = col_chunk->total_compressed_size();
460448

461449
// Try to get OffsetIndex for page-level ranges
462450
std::shared_ptr<::parquet::OffsetIndex> offset_index;
463-
if (rg_page_index_reader) {
464-
offset_index = rg_page_index_reader->GetOffsetIndex(col_idx);
451+
if (row_group_page_index_reader) {
452+
offset_index = row_group_page_index_reader->GetOffsetIndex(col_idx);
465453
}
466454

467455
if (!offset_index) {

src/paimon/format/parquet/page_filtered_row_group_reader.h

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -51,18 +51,20 @@ class PageFilteredRowGroupReader {
5151
/// Read a row group with page-level filtering.
5252
/// @param target_row_group Target row group with index and row ranges
5353
/// @param column_indices Leaf column indices to read
54-
/// @param pool Memory pool
5554
/// @param cache_options Cache options for PreBuffer
5655
/// @param pre_buffered If true, assumes PreBuffer was already called externally
5756
/// and only waits via WhenBuffered (no redundant PreBuffer).
5857
/// @param page_ranges If non-empty, wait via WhenBufferedRanges instead of WhenBuffered
5958
/// @param max_chunksize Per-batch row cap for the returned reader.
59+
/// @param row_group_page_index_reader Reusable page-index reader for the target row group
60+
/// @param pool Memory pool
6061
/// @param arrow_file_reader The Arrow FileReader for ColumnReader tree creation
6162
/// @return A RecordBatchReader streaming the filtered rows.
6263
static Result<std::unique_ptr<arrow::RecordBatchReader>> ReadFilteredRowGroup(
6364
const TargetRowGroup& target_row_group, const std::vector<int32_t>& column_indices,
6465
const ::arrow::io::CacheOptions& cache_options, bool pre_buffered,
6566
const std::vector<::arrow::io::ReadRange>& page_ranges, int64_t max_chunksize,
67+
const std::shared_ptr<::parquet::RowGroupPageIndexReader>& row_group_page_index_reader,
6668
std::shared_ptr<::arrow::MemoryPool> pool, ::parquet::arrow::FileReader* arrow_file_reader);
6769

6870
/// Compute the byte ranges of pages that overlap with the given RowRanges.
@@ -71,6 +73,7 @@ class PageFilteredRowGroupReader {
7173
/// Falls back to entire column chunk range if OffsetIndex is unavailable.
7274
static std::vector<::arrow::io::ReadRange> ComputePageRanges(
7375
const TargetRowGroup& target_row_group, const std::vector<int32_t>& column_indices,
76+
const std::shared_ptr<::parquet::RowGroupPageIndexReader>& row_group_page_index_reader,
7477
::parquet::ParquetFileReader* parquet_reader);
7578

7679
private:

0 commit comments

Comments
 (0)