SteNicholas commented on code in PR #209:
URL: https://github.com/apache/paimon-cpp/pull/209#discussion_r3810802931
##########
src/paimon/format/parquet/file_reader_wrapper.cpp:
##########
@@ -376,39 +393,68 @@ Status FileReaderWrapper::PrepareForReadingLazy(
target_row_groups_ = target_row_groups;
target_column_indices_ = column_indices;
reader_initialized_ = false;
+ pending_start_idx_.reset();
return Status::OK();
}
-std::vector<::arrow::io::ReadRange> FileReaderWrapper::CollectPreBufferRanges(
- const std::vector<int32_t>& column_indices) {
- std::vector<::arrow::io::ReadRange> ranges;
- auto file_metadata = file_reader_->parquet_reader()->metadata();
-
- for (const auto& trg : target_row_groups_) {
- if (trg.IsExcludedByReadRange()) continue;
-
- if (trg.IsPartiallyMatched()) {
- // Page-filtered RGs: only matching page byte ranges.
- auto row_group_page_index_reader =
GetRowGroupPageIndexReader(trg.GetRowGroupIndex());
- auto page_ranges = PageFilteredRowGroupReader::ComputePageRanges(
- trg, column_indices, row_group_page_index_reader,
file_reader_->parquet_reader());
- ranges.insert(ranges.end(),
std::make_move_iterator(page_ranges.begin()),
- std::make_move_iterator(page_ranges.end()));
- } else {
- // Fully-matched RGs: entire column chunk ranges.
- auto rg_metadata = file_metadata->RowGroup(trg.GetRowGroupIndex());
- for (int32_t col_idx : column_indices) {
- auto col_chunk = rg_metadata->ColumnChunk(col_idx);
- int64_t offset = col_chunk->data_page_offset();
- if (col_chunk->has_dictionary_page() &&
col_chunk->dictionary_page_offset() > 0 &&
- offset > col_chunk->dictionary_page_offset()) {
- offset = col_chunk->dictionary_page_offset();
+Result<std::vector<::arrow::io::ReadRange>>
FileReaderWrapper::CollectPreBufferRanges(
+ const std::vector<int32_t>& column_indices, uint64_t start_idx) {
+ return DoCollectPreBufferRanges(column_indices,
/*skip_read_range_excluded=*/true, start_idx);
+}
+
+Result<std::vector<::arrow::io::ReadRange>>
FileReaderWrapper::DoCollectPreBufferRanges(
+ const std::vector<int32_t>& column_indices, bool skip_read_range_excluded,
uint64_t start_idx) {
+ try {
+ std::vector<::arrow::io::ReadRange> ranges;
+ auto file_metadata = file_reader_->parquet_reader()->metadata();
+
+ for (uint64_t idx = start_idx; idx < target_row_groups_.size(); idx++)
{
+ const auto& trg = target_row_groups_[idx];
+ if (skip_read_range_excluded && trg.IsExcludedByReadRange()) {
+ continue;
+ }
+
+ if (trg.IsPartiallyMatched()) {
+ // Page-filtered RGs: only matching page byte ranges.
+ auto row_group_page_index_reader =
+ GetRowGroupPageIndexReader(trg.GetRowGroupIndex());
+ auto page_ranges =
PageFilteredRowGroupReader::ComputePageRanges(
+ trg, column_indices, row_group_page_index_reader,
+ file_reader_->parquet_reader());
+ ranges.insert(ranges.end(),
std::make_move_iterator(page_ranges.begin()),
+ std::make_move_iterator(page_ranges.end()));
+ } else {
+ // Fully-matched RGs: entire column chunk ranges.
+ auto rg_metadata =
file_metadata->RowGroup(trg.GetRowGroupIndex());
+ for (int32_t col_idx : column_indices) {
+ auto col_chunk = rg_metadata->ColumnChunk(col_idx);
+ int64_t offset = col_chunk->data_page_offset();
+ if (col_chunk->has_dictionary_page() &&
+ col_chunk->dictionary_page_offset() > 0 &&
+ offset > col_chunk->dictionary_page_offset()) {
+ offset = col_chunk->dictionary_page_offset();
+ }
+ ranges.push_back({offset,
col_chunk->total_compressed_size()});
}
- ranges.push_back({offset, col_chunk->total_compressed_size()});
}
}
+ return ranges;
}
- return ranges;
+
PAIMON_PARQUET_CATCH_AND_RETURN_STATUS("FileReaderWrapper::DoCollectPreBufferRanges")
+}
+
+Result<std::vector<std::pair<uint64_t, uint64_t>>>
FileReaderWrapper::GetPreBufferRanges() {
+ PAIMON_ASSIGN_OR_RAISE(std::vector<::arrow::io::ReadRange> ranges,
+ DoCollectPreBufferRanges(target_column_indices_,
+
/*skip_read_range_excluded=*/false,
+ /*start_idx=*/0));
+ std::vector<std::pair<uint64_t, uint64_t>> pre_buffer_ranges;
+ pre_buffer_ranges.reserve(ranges.size());
+ for (const auto& range : ranges) {
+ pre_buffer_ranges.emplace_back(static_cast<uint64_t>(range.offset),
+ static_cast<uint64_t>(range.length));
Review Comment:
Validate signed Parquet ranges before conversion. These Arrow ranges contain
signed metadata values, but negative offsets or `total_compressed_size` values
are cast directly to very large `uint64_t`s. `ReadAheadCache::Init()` then
calls `CoalesceByteRanges()` before its bounds checks; the combiner performs
unchecked `offset + length` arithmetic and can enter a pathological splitting
loop, exhausting memory on a corrupt footer. Please validate that offset and
length are non-negative and that their sum does not overflow before converting
them.
--
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]