Skip to content

Commit c5e453e

Browse files
committed
[fix](profile) sort out parquet reader profile (#58895)
Problem Summary: Refine some metrics in parquet reader profile. 1. Rename some `Statistics` class name to make it readable. (There are too many `Statistics` struct with same name) 2. Add `read page header timer` in parquet reader profile 3. fix issue of invalid check logic for `MergeRangeFileReader` when setting prefetch buffer size 4. fix issue that data cache profile is incorrect for external table can.
1 parent a2380c1 commit c5e453e

12 files changed

Lines changed: 129 additions & 159 deletions

be/src/io/fs/buffered_reader.cpp

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -819,12 +819,10 @@ Status BufferedFileStreamReader::read_bytes(const uint8_t** buf, uint64_t offset
819819
int64_t buf_remaining = _buf_end_offset - _buf_start_offset;
820820
int64_t to_read = std::min(_buf_size - buf_remaining, _file_end_offset - _buf_end_offset);
821821
int64_t has_read = 0;
822-
SCOPED_RAW_TIMER(&_statistics.read_time);
823822
while (has_read < to_read) {
824823
size_t loop_read = 0;
825824
Slice result(_buf.get() + buf_remaining + has_read, to_read - has_read);
826825
RETURN_IF_ERROR(_file->read_at(_buf_end_offset + has_read, result, &loop_read, io_ctx));
827-
_statistics.read_calls++;
828826
if (loop_read == 0) {
829827
break;
830828
}
@@ -833,7 +831,6 @@ Status BufferedFileStreamReader::read_bytes(const uint8_t** buf, uint64_t offset
833831
if (has_read != to_read) {
834832
return Status::Corruption("Try to read {} bytes, but received {} bytes", to_read, has_read);
835833
}
836-
_statistics.read_bytes += to_read;
837834
_buf_end_offset += to_read;
838835
*buf = _buf.get();
839836
return Status::OK();

be/src/io/fs/buffered_reader.h

Lines changed: 0 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -225,7 +225,6 @@ class MergeRangeFileReader : public io::FileReader {
225225
int64_t merged_io = 0;
226226
int64_t request_bytes = 0;
227227
int64_t merged_bytes = 0;
228-
int64_t apply_bytes = 0;
229228
};
230229

231230
struct RangeCachedData {
@@ -299,9 +298,6 @@ class MergeRangeFileReader : public io::FileReader {
299298
_merged_read_slice_size = READ_SLICE_SIZE;
300299
}
301300

302-
for (const PrefetchRange& range : _random_access_ranges) {
303-
_statistics.apply_bytes += range.end_offset - range.start_offset;
304-
}
305301
if (_profile != nullptr) {
306302
const char* random_profile = "MergedSmallIO";
307303
ADD_TIMER_WITH_LEVEL(_profile, random_profile, 1);
@@ -315,8 +311,6 @@ class MergeRangeFileReader : public io::FileReader {
315311
random_profile, 1);
316312
_merged_bytes = ADD_CHILD_COUNTER_WITH_LEVEL(_profile, "MergedBytes", TUnit::BYTES,
317313
random_profile, 1);
318-
_apply_bytes = ADD_CHILD_COUNTER_WITH_LEVEL(_profile, "ApplyBytes", TUnit::BYTES,
319-
random_profile, 1);
320314
}
321315
}
322316

@@ -359,7 +353,6 @@ class MergeRangeFileReader : public io::FileReader {
359353
COUNTER_UPDATE(_merged_io, _statistics.merged_io);
360354
COUNTER_UPDATE(_request_bytes, _statistics.request_bytes);
361355
COUNTER_UPDATE(_merged_bytes, _statistics.merged_bytes);
362-
COUNTER_UPDATE(_apply_bytes, _statistics.apply_bytes);
363356
if (_reader != nullptr) {
364357
_reader->collect_profile_before_close();
365358
}
@@ -373,7 +366,6 @@ class MergeRangeFileReader : public io::FileReader {
373366
RuntimeProfile::Counter* _merged_io = nullptr;
374367
RuntimeProfile::Counter* _request_bytes = nullptr;
375368
RuntimeProfile::Counter* _merged_bytes = nullptr;
376-
RuntimeProfile::Counter* _apply_bytes = nullptr;
377369

378370
int _search_read_range(size_t start_offset, size_t end_offset);
379371
void _clean_cached_data(RangeCachedData& cached_data);
@@ -619,12 +611,6 @@ class InMemoryFileReader final : public io::FileReader {
619611
*/
620612
class BufferedStreamReader {
621613
public:
622-
struct Statistics {
623-
int64_t read_time = 0;
624-
int64_t read_calls = 0;
625-
int64_t read_bytes = 0;
626-
};
627-
628614
/**
629615
* Return the address of underlying buffer that locates the start of data between [offset, offset + bytes_to_read)
630616
* @param buf the buffer address to save the start address of data
@@ -637,13 +623,9 @@ class BufferedStreamReader {
637623
* Save the data address to slice.data, and the slice.size is the bytes to read.
638624
*/
639625
virtual Status read_bytes(Slice& slice, uint64_t offset, const IOContext* io_ctx) = 0;
640-
Statistics& statistics() { return _statistics; }
641626
virtual ~BufferedStreamReader() = default;
642627
// return the file path
643628
virtual std::string path() = 0;
644-
645-
protected:
646-
Statistics _statistics;
647629
};
648630

649631
class BufferedFileStreamReader : public BufferedStreamReader, public ProfileCollector {

be/src/vec/exec/format/parquet/vparquet_column_chunk_reader.cpp

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -188,8 +188,8 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::load_page_data() {
188188
// check decompressed buffer size
189189
_reserve_decompress_buf(uncompressed_size);
190190
_page_data = Slice(_decompress_buf.get(), uncompressed_size);
191-
SCOPED_RAW_TIMER(&_statistics.decompress_time);
192-
_statistics.decompress_cnt++;
191+
SCOPED_RAW_TIMER(&_chunk_statistics.decompress_time);
192+
_chunk_statistics.decompress_cnt++;
193193
RETURN_IF_ERROR(_block_compress_codec->decompress(compressed_data, &_page_data));
194194
} else {
195195
// Don't need decompress
@@ -204,7 +204,7 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::load_page_data() {
204204

205205
// Initialize repetition level and definition level. Skip when level = 0, which means required field.
206206
if (_max_rep_level > 0) {
207-
SCOPED_RAW_TIMER(&_statistics.decode_level_time);
207+
SCOPED_RAW_TIMER(&_chunk_statistics.decode_level_time);
208208
if (header->__isset.data_page_header_v2) {
209209
RETURN_IF_ERROR(_rep_level_decoder.init_v2(_v2_rep_levels, _max_rep_level,
210210
_remaining_rep_nums));
@@ -215,7 +215,7 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::load_page_data() {
215215
}
216216
}
217217
if (_max_def_level > 0) {
218-
SCOPED_RAW_TIMER(&_statistics.decode_level_time);
218+
SCOPED_RAW_TIMER(&_chunk_statistics.decode_level_time);
219219
if (header->__isset.data_page_header_v2) {
220220
RETURN_IF_ERROR(_def_level_decoder.init_v2(_v2_def_levels, _max_def_level,
221221
_remaining_def_nums));
@@ -255,7 +255,7 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::_decode_dict_page() {
255255
const tparquet::PageHeader* header = nullptr;
256256
RETURN_IF_ERROR(_page_reader->get_page_header(header));
257257
DCHECK_EQ(tparquet::PageType::DICTIONARY_PAGE, header->type);
258-
SCOPED_RAW_TIMER(&_statistics.decode_dict_time);
258+
SCOPED_RAW_TIMER(&_chunk_statistics.decode_dict_time);
259259

260260
// Using the PLAIN_DICTIONARY enum value is deprecated in the Parquet 2.0 specification.
261261
// Prefer using RLE_DICTIONARY in a data page and PLAIN in a dictionary page for Parquet 2.0+ files.
@@ -314,7 +314,7 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::skip_values(size_t num_va
314314
}
315315
_remaining_num_values -= num_values;
316316
if (skip_data) {
317-
SCOPED_RAW_TIMER(&_statistics.decode_value_time);
317+
SCOPED_RAW_TIMER(&_chunk_statistics.decode_value_time);
318318
return _page_decoder->skip_values(num_values);
319319
} else {
320320
return Status::OK();
@@ -328,7 +328,7 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::decode_values(
328328
if (select_vector.num_values() == 0) {
329329
return Status::OK();
330330
}
331-
SCOPED_RAW_TIMER(&_statistics.decode_value_time);
331+
SCOPED_RAW_TIMER(&_chunk_statistics.decode_value_time);
332332
if (UNLIKELY((doris_column->is_column_dictionary() || is_dict_filter) && !_has_dict)) {
333333
return Status::IOError("Not dictionary coded");
334334
}

be/src/vec/exec/format/parquet/vparquet_column_chunk_reader.h

Lines changed: 11 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,7 @@ struct ColumnChunkReaderStatistics {
6060
int64_t decode_level_time = 0;
6161
int64_t skip_page_header_num = 0;
6262
int64_t parse_page_header_num = 0;
63+
int64_t read_page_header_time = 0;
6364
};
6465

6566
/**
@@ -158,12 +159,15 @@ class ColumnChunkReader {
158159
// Get page decoder
159160
Decoder* get_page_decoder() { return _page_decoder; }
160161

161-
ColumnChunkReaderStatistics& statistics() {
162-
_statistics.decode_header_time = _page_reader->statistics().decode_header_time;
163-
_statistics.skip_page_header_num = _page_reader->statistics().skip_page_header_num;
164-
_statistics.parse_page_header_num = _page_reader->statistics().parse_page_header_num;
165-
_statistics.read_page_header_time = _page_reader->statistics().read_page_header_time;
166-
return _statistics;
162+
ColumnChunkReaderStatistics& chunk_statistics() {
163+
_chunk_statistics.decode_header_time = _page_reader->page_statistics().decode_header_time;
164+
_chunk_statistics.skip_page_header_num =
165+
_page_reader->page_statistics().skip_page_header_num;
166+
_chunk_statistics.parse_page_header_num =
167+
_page_reader->page_statistics().parse_page_header_num;
168+
_chunk_statistics.read_page_header_time =
169+
_page_reader->page_statistics().read_page_header_time;
170+
return _chunk_statistics;
167171
}
168172

169173
Status read_dict_values_to_column(MutableColumnPtr& doris_column) {
@@ -251,7 +255,7 @@ class ColumnChunkReader {
251255
// Map: encoding -> Decoder
252256
// Plain or Dictionary encoding. If the dictionary grows too big, the encoding will fall back to the plain encoding
253257
std::unordered_map<int, std::unique_ptr<Decoder>> _decoders;
254-
ColumnChunkReaderStatistics _statistics;
258+
ColumnChunkReaderStatistics _chunk_statistics;
255259
};
256260
#include "common/compile_check_end.h"
257261

be/src/vec/exec/format/parquet/vparquet_column_reader.h

Lines changed: 29 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -53,12 +53,9 @@ using ColumnString = ColumnStr<UInt32>;
5353

5454
class ParquetColumnReader {
5555
public:
56-
struct Statistics {
57-
Statistics()
58-
: read_time(0),
59-
read_calls(0),
60-
page_index_read_calls(0),
61-
read_bytes(0),
56+
struct ColumnStatistics {
57+
ColumnStatistics()
58+
: page_index_read_calls(0),
6259
decompress_time(0),
6360
decompress_cnt(0),
6461
decode_header_time(0),
@@ -70,12 +67,8 @@ class ParquetColumnReader {
7067
parse_page_header_num(0),
7168
read_page_header_time(0) {}
7269

73-
Statistics(io::BufferedStreamReader::Statistics& fs, ColumnChunkReaderStatistics& cs,
74-
int64_t null_map_time)
75-
: read_time(fs.read_time),
76-
read_calls(fs.read_calls),
77-
page_index_read_calls(0),
78-
read_bytes(fs.read_bytes),
70+
ColumnStatistics(ColumnChunkReaderStatistics& cs, int64_t null_map_time)
71+
: page_index_read_calls(0),
7972
decompress_time(cs.decompress_time),
8073
decompress_cnt(cs.decompress_cnt),
8174
decode_header_time(cs.decode_header_time),
@@ -87,10 +80,7 @@ class ParquetColumnReader {
8780
parse_page_header_num(cs.parse_page_header_num),
8881
read_page_header_time(cs.read_page_header_time) {}
8982

90-
int64_t read_time;
91-
int64_t read_calls;
9283
int64_t page_index_read_calls;
93-
int64_t read_bytes;
9484
int64_t decompress_time;
9585
int64_t decompress_cnt;
9686
int64_t decode_header_time;
@@ -102,21 +92,18 @@ class ParquetColumnReader {
10292
int64_t parse_page_header_num;
10393
int64_t read_page_header_time;
10494

105-
void merge(Statistics& statistics) {
106-
read_time += statistics.read_time;
107-
read_calls += statistics.read_calls;
108-
read_bytes += statistics.read_bytes;
109-
page_index_read_calls += statistics.page_index_read_calls;
110-
decompress_time += statistics.decompress_time;
111-
decompress_cnt += statistics.decompress_cnt;
112-
decode_header_time += statistics.decode_header_time;
113-
decode_value_time += statistics.decode_value_time;
114-
decode_dict_time += statistics.decode_dict_time;
115-
decode_level_time += statistics.decode_level_time;
116-
decode_null_map_time += statistics.decode_null_map_time;
117-
skip_page_header_num += statistics.skip_page_header_num;
118-
parse_page_header_num += statistics.parse_page_header_num;
119-
read_page_header_time += statistics.read_page_header_time;
95+
void merge(ColumnStatistics& col_statistics) {
96+
page_index_read_calls += col_statistics.page_index_read_calls;
97+
decompress_time += col_statistics.decompress_time;
98+
decompress_cnt += col_statistics.decompress_cnt;
99+
decode_header_time += col_statistics.decode_header_time;
100+
decode_value_time += col_statistics.decode_value_time;
101+
decode_dict_time += col_statistics.decode_dict_time;
102+
decode_level_time += col_statistics.decode_level_time;
103+
decode_null_map_time += col_statistics.decode_null_map_time;
104+
skip_page_header_num += col_statistics.skip_page_header_num;
105+
parse_page_header_num += col_statistics.parse_page_header_num;
106+
read_page_header_time += col_statistics.read_page_header_time;
120107
}
121108
};
122109

@@ -148,7 +135,7 @@ class ParquetColumnReader {
148135
const std::set<uint64_t>& filter_column_ids = {});
149136
virtual const std::vector<level_t>& get_rep_level() const = 0;
150137
virtual const std::vector<level_t>& get_def_level() const = 0;
151-
virtual Statistics statistics() = 0;
138+
virtual ColumnStatistics column_statistics() = 0;
152139
virtual void close() = 0;
153140

154141
virtual void reset_filter_map_index() = 0;
@@ -191,9 +178,8 @@ class ScalarColumnReader : public ParquetColumnReader {
191178
MutableColumnPtr convert_dict_column_to_string_column(const ColumnInt32* dict_column) override;
192179
const std::vector<level_t>& get_rep_level() const override { return _rep_levels; }
193180
const std::vector<level_t>& get_def_level() const override { return _def_levels; }
194-
Statistics statistics() override {
195-
return Statistics(_stream_reader->statistics(), _chunk_reader->statistics(),
196-
_decode_null_map_time);
181+
ColumnStatistics column_statistics() override {
182+
return ColumnStatistics(_chunk_reader->chunk_statistics(), _decode_null_map_time);
197183
}
198184
void close() override {}
199185

@@ -307,7 +293,7 @@ class ArrayColumnReader : public ParquetColumnReader {
307293
const std::vector<level_t>& get_def_level() const override {
308294
return _element_reader->get_def_level();
309295
}
310-
Statistics statistics() override { return _element_reader->statistics(); }
296+
ColumnStatistics column_statistics() override { return _element_reader->column_statistics(); }
311297
void close() override {}
312298

313299
void reset_filter_map_index() override { _element_reader->reset_filter_map_index(); }
@@ -338,9 +324,9 @@ class MapColumnReader : public ParquetColumnReader {
338324
return _key_reader->get_def_level();
339325
}
340326

341-
Statistics statistics() override {
342-
Statistics kst = _key_reader->statistics();
343-
Statistics vst = _value_reader->statistics();
327+
ColumnStatistics column_statistics() override {
328+
ColumnStatistics kst = _key_reader->column_statistics();
329+
ColumnStatistics vst = _value_reader->column_statistics();
344330
kst.merge(vst);
345331
return kst;
346332
}
@@ -395,12 +381,12 @@ class StructColumnReader : public ParquetColumnReader {
395381
return _child_readers.begin()->second->get_def_level();
396382
}
397383

398-
Statistics statistics() override {
399-
Statistics st;
384+
ColumnStatistics column_statistics() override {
385+
ColumnStatistics st;
400386
for (const auto& column_name : _read_column_names) {
401387
auto reader = _child_readers.find(column_name);
402388
if (reader != _child_readers.end()) {
403-
Statistics cst = reader->second->statistics();
389+
ColumnStatistics cst = reader->second->column_statistics();
404390
st.merge(cst);
405391
}
406392
}
@@ -493,8 +479,8 @@ class SkipReadingReader : public ParquetColumnReader {
493479
}
494480

495481
// Implement required pure virtual methods from base class
496-
Statistics statistics() override {
497-
return Statistics(); // Return empty statistics
482+
ColumnStatistics column_statistics() override {
483+
return ColumnStatistics(); // Return empty statistics
498484
}
499485

500486
void close() override {

be/src/vec/exec/format/parquet/vparquet_group_reader.cpp

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1128,10 +1128,10 @@ void RowGroupReader::_convert_dict_cols_to_string_cols(Block* block) {
11281128
}
11291129
}
11301130

1131-
ParquetColumnReader::Statistics RowGroupReader::statistics() {
1132-
ParquetColumnReader::Statistics st;
1131+
ParquetColumnReader::ColumnStatistics RowGroupReader::merged_column_statistics() {
1132+
ParquetColumnReader::ColumnStatistics st;
11331133
for (auto& reader : _column_readers) {
1134-
auto ost = reader.second->statistics();
1134+
auto ost = reader.second->column_statistics();
11351135
st.merge(ost);
11361136
}
11371137
return st;

be/src/vec/exec/format/parquet/vparquet_group_reader.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -173,7 +173,7 @@ class RowGroupReader : public ProfileCollector {
173173
int64_t predicate_filter_time() const { return _predicate_filter_time; }
174174
int64_t dict_filter_rewrite_time() const { return _dict_filter_rewrite_time; }
175175

176-
ParquetColumnReader::Statistics statistics();
176+
ParquetColumnReader::ColumnStatistics merged_column_statistics();
177177
void set_remaining_rows(int64_t rows) { _remaining_rows = rows; }
178178
int64_t get_remaining_rows() { return _remaining_rows; }
179179

be/src/vec/exec/format/parquet/vparquet_page_reader.cpp

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -88,11 +88,11 @@ Status PageReader<IN_COLLECTION, OFFSET_INDEX>::parse_page_header() {
8888
}
8989
header_size = std::min(header_size, max_size);
9090
{
91-
SCOPED_RAW_TIMER(&_statistics.read_page_header_time);
91+
SCOPED_RAW_TIMER(&_page_statistics.read_page_header_time);
9292
RETURN_IF_ERROR(_reader->read_bytes(&page_header_buf, _offset, header_size, _io_ctx));
9393
}
9494
real_header_size = cast_set<uint32_t>(header_size);
95-
SCOPED_RAW_TIMER(&_statistics.decode_header_time);
95+
SCOPED_RAW_TIMER(&_page_statistics.decode_header_time);
9696
auto st =
9797
deserialize_thrift_msg(page_header_buf, &real_header_size, true, &_cur_page_header);
9898
if (st.ok()) {
@@ -115,7 +115,7 @@ Status PageReader<IN_COLLECTION, OFFSET_INDEX>::parse_page_header() {
115115
}
116116
}
117117

118-
_statistics.parse_page_header_num++;
118+
_page_statistics.parse_page_header_num++;
119119
_offset += real_header_size;
120120
_next_header_offset = _offset + _cur_page_header.compressed_page_size;
121121
_state = HEADER_PARSED;

0 commit comments

Comments
 (0)