Skip to content

Commit 0dafd41

Browse files
authored
feat(parquet): add metrics for parquet reader observability (#258)
1 parent 36e1e94 commit 0dafd41

3 files changed

Lines changed: 18 additions & 0 deletions

File tree

src/paimon/format/parquet/parquet_file_batch_reader.cpp

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -151,6 +151,9 @@ Status ParquetFileBatchReader::SetReadSchema(
151151
read_row_groups_ = row_groups;
152152
read_column_indices_ = column_indices;
153153

154+
metrics_->SetCounter(ParquetMetrics::READ_ROW_GROUPS_TOTAL, reader_->GetNumberOfRowGroups());
155+
metrics_->SetCounter(ParquetMetrics::READ_ROW_GROUPS_FILTERED, row_groups.size());
156+
154157
PAIMON_ASSIGN_OR_RAISE(std::set<int32_t> ordered_row_groups,
155158
reader_->FilterRowGroupsByReadRanges(read_ranges_, read_row_groups_));
156159
return reader_->PrepareForReadingLazy(ordered_row_groups, read_column_indices_);
@@ -243,6 +246,12 @@ Result<BatchReader::ReadBatch> ParquetFileBatchReader::NextBatch() {
243246
std::unique_ptr<ArrowArray> c_array = std::make_unique<ArrowArray>();
244247
std::unique_ptr<ArrowSchema> c_schema = std::make_unique<ArrowSchema>();
245248
PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportArray(*array, c_array.get(), c_schema.get()));
249+
250+
read_rows_ += array->length();
251+
read_batch_count_++;
252+
metrics_->SetCounter(ParquetMetrics::READ_ROWS, read_rows_);
253+
metrics_->SetCounter(ParquetMetrics::READ_BATCH_COUNT, read_batch_count_);
254+
246255
return make_pair(std::move(c_array), std::move(c_schema));
247256
}
248257

src/paimon/format/parquet/parquet_file_batch_reader.h

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -174,6 +174,9 @@ class ParquetFileBatchReader : public PrefetchFileBatchReader {
174174

175175
std::shared_ptr<Metrics> metrics_;
176176

177+
uint64_t read_rows_ = 0;
178+
uint64_t read_batch_count_ = 0;
179+
177180
// last time set read schema
178181
std::vector<int32_t> read_row_groups_;
179182
std::vector<int32_t> read_column_indices_;

src/paimon/format/parquet/parquet_format_defs.h

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,12 @@ static constexpr uint32_t DEFAULT_PARQUET_READ_PREDICATE_NODE_COUNT_LIMIT = 512;
6161
class ParquetMetrics {
6262
public:
6363
static inline const char WRITE_RECORD_COUNT[] = "parquet.write.record.count";
64+
65+
// read
66+
static inline const char READ_ROW_GROUPS_TOTAL[] = "parquet.read.row-groups.total";
67+
static inline const char READ_ROW_GROUPS_FILTERED[] = "parquet.read.row-groups.filtered";
68+
static inline const char READ_ROWS[] = "parquet.read.rows";
69+
static inline const char READ_BATCH_COUNT[] = "parquet.read.batch-count";
6470
};
6571

6672
} // namespace paimon::parquet

0 commit comments

Comments
 (0)