Skip to content

Commit a77b747

Browse files
authored
fix: small fixes for AI reviewing feedback (alibaba#365)
1 parent 300228c commit a77b747

4 files changed

Lines changed: 23 additions & 15 deletions

File tree

src/paimon/format/parquet/file_reader_wrapper.cpp

Lines changed: 15 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -105,7 +105,7 @@ Result<std::unique_ptr<FileReaderWrapper>> FileReaderWrapper::Create(
105105
std::move(file_reader), all_row_group_ranges, num_rows, batch_size, pool));
106106
std::vector<TargetRowGroup> all_target_row_groups;
107107
for (int32_t i = 0; i < file_reader_wrapper->GetNumberOfRowGroups(); i++) {
108-
all_target_row_groups.emplace_back(/*rg_index=*/i, /*page_filtered=*/false,
108+
all_target_row_groups.emplace_back(/*rg_index=*/i, /*is_partially_matched=*/false,
109109
/*ranges=*/RowRanges());
110110
}
111111
PAIMON_RETURN_NOT_OK(
@@ -144,7 +144,7 @@ FileReaderWrapper::FileReaderWrapper(
144144
int64_t batch_size, std::shared_ptr<::arrow::MemoryPool> pool)
145145
: file_reader_(std::move(file_reader)),
146146
all_row_group_ranges_(all_row_group_ranges),
147-
pool_(pool),
147+
pool_(std::move(pool)),
148148
batch_size_(batch_size),
149149
num_rows_(num_rows) {}
150150

@@ -184,7 +184,7 @@ Status FileReaderWrapper::SeekToRow(uint64_t row_number) {
184184
if (target_row_groups_[i].excluded_by_read_range) {
185185
continue;
186186
}
187-
uint32_t rg_id = target_row_groups_[i].row_group_index;
187+
int32_t rg_id = target_row_groups_[i].row_group_index;
188188
uint64_t rg_start = all_row_group_ranges_[rg_id].first;
189189
uint64_t rg_end = all_row_group_ranges_[rg_id].second;
190190
if (row_number > rg_start && row_number < rg_end) {
@@ -297,12 +297,13 @@ Result<std::shared_ptr<arrow::RecordBatch>> FileReaderWrapper::Next() {
297297
}
298298

299299
while (current_row_group_idx_ < target_row_groups_.size()) {
300-
bool is_page_filtered = target_row_groups_[current_row_group_idx_].is_partially_matched;
300+
bool is_partially_matched =
301+
target_row_groups_[current_row_group_idx_].is_partially_matched;
301302
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::RecordBatch> batch,
302-
is_page_filtered ? NextPageFiltered() : NextFullyMatched());
303+
is_partially_matched ? NextPageFiltered() : NextFullyMatched());
303304
if (batch) {
304305
return batch;
305-
} else if (!is_page_filtered) {
306+
} else if (!is_partially_matched) {
306307
// Null from fully-matched path means batch_reader_ is globally exhausted.
307308
break;
308309
}
@@ -424,13 +425,18 @@ Status FileReaderWrapper::PrepareForReading(const std::vector<TargetRowGroup>& t
424425
}
425426
}
426427

427-
bool has_page_filtered = fully_matched_row_groups.size() != active_count;
428-
if (has_page_filtered) {
428+
bool has_partially_matched = fully_matched_row_groups.size() != active_count;
429+
if (has_partially_matched) {
429430
PAIMON_RETURN_NOT_OK(BuildPageFilteredSchema(column_indices));
430431
}
431432

432433
WaitForPendingPreBuffer();
433434

435+
// TODO(Yonghao Fang): Neither Paimon nor Arrow manage the size and lifecycle of prebuffered
436+
// caches. So when a lot of row is needed, there is possibility of OOM due to too much
437+
// prebuffering. Also, DispatchPreBuffer will drop previous prebuffered ranges by
438+
// GetRecordBatchReader, which cause IO wastes.
439+
434440
// Create standard reader for fully-matched row groups.
435441
if (!fully_matched_row_groups.empty()) {
436442
PAIMON_RETURN_NOT_OK_FROM_ARROW(file_reader_->GetRecordBatchReader(
@@ -441,7 +447,7 @@ Status FileReaderWrapper::PrepareForReading(const std::vector<TargetRowGroup>& t
441447

442448
// When page-filtered RGs exist, issue a single PreBuffer covering both kinds.
443449
// Otherwise GetRecordBatchReader already issued PreBuffer internally.
444-
if (has_page_filtered) {
450+
if (has_partially_matched) {
445451
auto all_ranges = CollectPreBufferRanges(column_indices);
446452
DispatchPreBuffer(std::move(all_ranges));
447453
}

src/paimon/format/parquet/page_filtered_row_group_reader.cpp

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -123,8 +123,9 @@ std::pair<RowRanges, int64_t> PageFilteredRowGroupReader::ComputeCompressedRowRa
123123
}
124124

125125
Status PageFilteredRowGroupReader::ExecuteSkipReadPattern(
126-
std::shared_ptr<::parquet::internal::RecordReader> record_reader, const RowRanges& ranges,
127-
int64_t total_row_count, int32_t row_group_index, int32_t column_index) {
126+
const std::shared_ptr<::parquet::internal::RecordReader>& record_reader,
127+
const RowRanges& ranges, int64_t total_row_count, int32_t row_group_index,
128+
int32_t column_index) {
128129
int64_t current_row = 0;
129130
for (const auto& range : ranges.GetRanges()) {
130131
if (range.from > current_row) {

src/paimon/format/parquet/page_filtered_row_group_reader.h

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -88,8 +88,9 @@ class PageFilteredRowGroupReader {
8888

8989
/// Execute the skip/read pattern on a RecordReader based on RowRanges.
9090
static Status ExecuteSkipReadPattern(
91-
std::shared_ptr<::parquet::internal::RecordReader> record_reader, const RowRanges& ranges,
92-
int64_t total_row_count, int32_t row_group_index, int32_t column_index);
91+
const std::shared_ptr<::parquet::internal::RecordReader>& record_reader,
92+
const RowRanges& ranges, int64_t total_row_count, int32_t row_group_index,
93+
int32_t column_index);
9394

9495
/// Create a data_page_filter callback for a column based on RowRanges + OffsetIndex.
9596
static std::function<bool(const ::parquet::DataPageStats&)> MakePageFilter(

src/paimon/format/parquet/parquet_file_batch_reader.cpp

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -199,11 +199,11 @@ Status ParquetFileBatchReader::SetReadSchema(
199199
for (int32_t rg_id : row_groups) {
200200
auto it = row_group_row_ranges.find(rg_id);
201201
if (it != row_group_row_ranges.end()) {
202-
target_row_groups.emplace_back(/*rg_index=*/rg_id, /*page_filtered=*/true,
202+
target_row_groups.emplace_back(/*rg_index=*/rg_id, /*is_partially_matched=*/true,
203203
/*ranges=*/it->second);
204204
} else {
205205
target_row_groups.emplace_back(/*rg_index=*/rg_id,
206-
/*page_filtered=*/false,
206+
/*is_partially_matched=*/false,
207207
/*ranges=*/RowRanges());
208208
}
209209
}

0 commit comments

Comments
 (0)