Coverage Report

Created: 2026-07-30 07:20

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/parquet/vparquet_page_reader.cpp
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  You may obtain a copy of the License at
8
//
9
//   http://www.apache.org/licenses/LICENSE-2.0
10
//
11
// Unless required by applicable law or agreed to in writing,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
18
#include "format/parquet/vparquet_page_reader.h"
19
20
#include <fmt/format.h>
21
#include <gen_cpp/parquet_types.h>
22
#include <stddef.h>
23
#include <stdint.h>
24
25
#include <algorithm>
26
27
#include "common/compiler_util.h" // IWYU pragma: keep
28
#include "common/config.h"
29
#include "format/parquet/parquet_common.h"
30
#include "io/fs/buffered_reader.h"
31
#include "runtime/runtime_profile.h"
32
#include "storage/cache/page_cache.h"
33
#include "util/slice.h"
34
#include "util/thrift_util.h"
35
36
namespace doris {
37
namespace io {
38
struct IOContext;
39
} // namespace io
40
} // namespace doris
41
42
namespace doris {
43
static constexpr size_t INIT_PAGE_HEADER_SIZE = 128;
44
45
193
void ParquetPageCacheKeyBuilder::init(const std::string& path, int64_t mtime) {
46
193
    _file_key_prefix = fmt::format("{}::{}", path, mtime);
47
193
}
48
49
template <bool IN_COLLECTION, bool OFFSET_INDEX>
50
PageReader<IN_COLLECTION, OFFSET_INDEX>::PageReader(io::BufferedStreamReader* reader,
51
                                                    io::IOContext* io_ctx, uint64_t offset,
52
                                                    uint64_t length, size_t total_rows,
53
                                                    const tparquet::ColumnMetaData& metadata,
54
                                                    const ParquetPageReadContext& page_read_ctx,
55
                                                    const tparquet::OffsetIndex* offset_index)
56
193
        : _reader(reader),
57
193
          _io_ctx(io_ctx),
58
193
          _offset(offset),
59
193
          _start_offset(offset),
60
193
          _end_offset(offset + length),
61
193
          _total_rows(total_rows),
62
193
          _metadata(metadata),
63
193
          _page_read_ctx(page_read_ctx),
64
193
          _offset_index(offset_index) {
65
193
    _next_header_offset = _offset;
66
193
    _state = INITIALIZED;
67
193
    _page_cache_key_builder.init(_reader->path(), _reader->mtime());
68
69
193
    if constexpr (OFFSET_INDEX) {
70
3
        _end_row = _offset_index->page_locations.size() >= 2
71
3
                           ? _offset_index->page_locations[1].first_row_index
72
3
                           : _total_rows;
73
3
    }
74
193
}
Unexecuted instantiation: _ZN5doris10PageReaderILb1ELb1EEC2EPNS_2io20BufferedStreamReaderEPNS2_9IOContextEmmmRKN8tparquet14ColumnMetaDataERKNS_22ParquetPageReadContextEPKNS7_11OffsetIndexE
_ZN5doris10PageReaderILb1ELb0EEC2EPNS_2io20BufferedStreamReaderEPNS2_9IOContextEmmmRKN8tparquet14ColumnMetaDataERKNS_22ParquetPageReadContextEPKNS7_11OffsetIndexE
Line
Count
Source
56
3
        : _reader(reader),
57
3
          _io_ctx(io_ctx),
58
3
          _offset(offset),
59
3
          _start_offset(offset),
60
3
          _end_offset(offset + length),
61
3
          _total_rows(total_rows),
62
3
          _metadata(metadata),
63
3
          _page_read_ctx(page_read_ctx),
64
3
          _offset_index(offset_index) {
65
3
    _next_header_offset = _offset;
66
3
    _state = INITIALIZED;
67
3
    _page_cache_key_builder.init(_reader->path(), _reader->mtime());
68
69
    if constexpr (OFFSET_INDEX) {
70
        _end_row = _offset_index->page_locations.size() >= 2
71
                           ? _offset_index->page_locations[1].first_row_index
72
                           : _total_rows;
73
    }
74
3
}
_ZN5doris10PageReaderILb0ELb1EEC2EPNS_2io20BufferedStreamReaderEPNS2_9IOContextEmmmRKN8tparquet14ColumnMetaDataERKNS_22ParquetPageReadContextEPKNS7_11OffsetIndexE
Line
Count
Source
56
3
        : _reader(reader),
57
3
          _io_ctx(io_ctx),
58
3
          _offset(offset),
59
3
          _start_offset(offset),
60
3
          _end_offset(offset + length),
61
3
          _total_rows(total_rows),
62
3
          _metadata(metadata),
63
3
          _page_read_ctx(page_read_ctx),
64
3
          _offset_index(offset_index) {
65
3
    _next_header_offset = _offset;
66
3
    _state = INITIALIZED;
67
3
    _page_cache_key_builder.init(_reader->path(), _reader->mtime());
68
69
3
    if constexpr (OFFSET_INDEX) {
70
3
        _end_row = _offset_index->page_locations.size() >= 2
71
3
                           ? _offset_index->page_locations[1].first_row_index
72
3
                           : _total_rows;
73
3
    }
74
3
}
_ZN5doris10PageReaderILb0ELb0EEC2EPNS_2io20BufferedStreamReaderEPNS2_9IOContextEmmmRKN8tparquet14ColumnMetaDataERKNS_22ParquetPageReadContextEPKNS7_11OffsetIndexE
Line
Count
Source
56
187
        : _reader(reader),
57
187
          _io_ctx(io_ctx),
58
187
          _offset(offset),
59
187
          _start_offset(offset),
60
187
          _end_offset(offset + length),
61
187
          _total_rows(total_rows),
62
187
          _metadata(metadata),
63
187
          _page_read_ctx(page_read_ctx),
64
187
          _offset_index(offset_index) {
65
187
    _next_header_offset = _offset;
66
187
    _state = INITIALIZED;
67
187
    _page_cache_key_builder.init(_reader->path(), _reader->mtime());
68
69
    if constexpr (OFFSET_INDEX) {
70
        _end_row = _offset_index->page_locations.size() >= 2
71
                           ? _offset_index->page_locations[1].first_row_index
72
                           : _total_rows;
73
    }
74
187
}
75
76
template <bool IN_COLLECTION, bool OFFSET_INDEX>
77
384
Status PageReader<IN_COLLECTION, OFFSET_INDEX>::parse_page_header() {
78
384
    if (_state == HEADER_PARSED) {
79
140
        return Status::OK();
80
140
    }
81
244
    if (UNLIKELY(_offset < _start_offset || _offset >= _end_offset)) {
82
0
        return Status::IOError("Out-of-bounds Access");
83
0
    }
84
244
    if (UNLIKELY(_offset != _next_header_offset)) {
85
0
        return Status::IOError("Wrong header position, should seek to a page header first");
86
0
    }
87
244
    if (UNLIKELY(_state != INITIALIZED)) {
88
0
        return Status::IOError("Should skip or load current page to get next page");
89
0
    }
90
91
244
    _page_statistics.page_read_counter += 1;
92
93
    // Parse page header from file; header bytes are saved for possible cache insertion
94
244
    const uint8_t* page_header_buf = nullptr;
95
244
    size_t max_size = _end_offset - _offset;
96
244
    size_t header_size = std::min(INIT_PAGE_HEADER_SIZE, max_size);
97
244
    const size_t MAX_PAGE_HEADER_SIZE = config::parquet_header_max_size_mb << 20;
98
244
    uint32_t real_header_size = 0;
99
100
    // Try a header-only lookup in the page cache. Cached pages store
101
    // header + optional v2 levels + uncompressed payload, so we can
102
    // parse the page header directly from the cached bytes and avoid
103
    // a file read for the header.
104
244
    if (_page_read_ctx.enable_parquet_file_page_cache && !config::disable_storage_page_cache &&
105
244
        StoragePageCache::instance() != nullptr) {
106
226
        PageCacheHandle handle;
107
226
        StoragePageCache::CacheKey key = make_page_cache_key(static_cast<int64_t>(_offset));
108
226
        if (StoragePageCache::instance()->lookup(key, &handle, segment_v2::DATA_PAGE)) {
109
            // Parse header directly from cached data
110
118
            _page_cache_handle = std::move(handle);
111
118
            Slice s = _page_cache_handle.data();
112
118
            real_header_size = cast_set<uint32_t>(s.size);
113
118
            SCOPED_RAW_TIMER(&_page_statistics.decode_header_time);
114
118
            auto st = deserialize_thrift_msg(reinterpret_cast<const uint8_t*>(s.data),
115
118
                                             &real_header_size, true, &_cur_page_header);
116
118
            if (!st.ok()) return st;
117
            // Increment page cache counters for a true cache hit on header+payload
118
118
            _page_statistics.page_cache_hit_counter += 1;
119
            // Detect whether the cached payload is compressed or decompressed and record
120
118
            bool is_cache_payload_decompressed =
121
118
                    should_cache_decompressed(&_cur_page_header, _metadata);
122
123
118
            if (is_cache_payload_decompressed) {
124
118
                _page_statistics.page_cache_decompressed_hit_counter += 1;
125
118
            } else {
126
0
                _page_statistics.page_cache_compressed_hit_counter += 1;
127
0
            }
128
129
118
            _is_cache_payload_decompressed = is_cache_payload_decompressed;
130
131
118
            if constexpr (OFFSET_INDEX == false) {
132
118
                if (is_header_v2()) {
133
11
                    _end_row = _start_row + _cur_page_header.data_page_header_v2.num_rows;
134
106
                } else if constexpr (!IN_COLLECTION) {
135
106
                    _end_row = _start_row + _cur_page_header.data_page_header.num_values;
136
106
                }
137
118
            }
138
139
            // Save header bytes for later use (e.g., to insert updated cache entries)
140
118
            _header_buf.assign(s.data, s.data + real_header_size);
141
118
            _last_header_size = real_header_size;
142
118
            _page_statistics.parse_page_header_num++;
143
118
            _offset += real_header_size;
144
118
            _next_header_offset = _offset + _cur_page_header.compressed_page_size;
145
118
            _state = HEADER_PARSED;
146
118
            return Status::OK();
147
118
        } else {
148
108
            _page_statistics.page_cache_missing_counter += 1;
149
            // Clear any existing cache handle on miss to avoid holding stale handle
150
108
            _page_cache_handle = PageCacheHandle();
151
108
        }
152
226
    }
153
    // NOTE: page cache lookup for *decompressed* page data is handled in
154
    // ColumnChunkReader::load_page_data(). PageReader should only be
155
    // responsible for parsing the header bytes from the file and saving
156
    // them in `_header_buf` for possible later insertion into the cache.
157
128
    while (true) {
158
128
        if (UNLIKELY(_io_ctx && _io_ctx->should_stop)) {
159
0
            return Status::EndOfFile("stop");
160
0
        }
161
128
        header_size = std::min(header_size, max_size);
162
128
        {
163
128
            SCOPED_RAW_TIMER(&_page_statistics.read_page_header_time);
164
128
            RETURN_IF_ERROR(_reader->read_bytes(&page_header_buf, _offset, header_size, _io_ctx));
165
128
        }
166
127
        real_header_size = cast_set<uint32_t>(header_size);
167
127
        SCOPED_RAW_TIMER(&_page_statistics.decode_header_time);
168
127
        auto st =
169
127
                deserialize_thrift_msg(page_header_buf, &real_header_size, true, &_cur_page_header);
170
127
        if (st.ok()) {
171
125
            break;
172
125
        }
173
2
        if (_offset + header_size >= _end_offset || real_header_size > MAX_PAGE_HEADER_SIZE) {
174
0
            return Status::IOError(
175
0
                    "Failed to deserialize parquet page header. offset: {}, "
176
0
                    "header size: {}, end offset: {}, real header size: {}",
177
0
                    _offset, header_size, _end_offset, real_header_size);
178
0
        }
179
2
        header_size <<= 2;
180
2
    }
181
182
125
    if constexpr (OFFSET_INDEX == false) {
183
119
        if (is_header_v2()) {
184
9
            _end_row = _start_row + _cur_page_header.data_page_header_v2.num_rows;
185
108
        } else if constexpr (!IN_COLLECTION) {
186
108
            _end_row = _start_row + _cur_page_header.data_page_header.num_values;
187
108
        }
188
119
    }
189
190
    // Save header bytes for possible cache insertion later
191
125
    _header_buf.assign(page_header_buf, page_header_buf + real_header_size);
192
125
    _last_header_size = real_header_size;
193
125
    _page_statistics.parse_page_header_num++;
194
125
    _offset += real_header_size;
195
125
    _next_header_offset = _offset + _cur_page_header.compressed_page_size;
196
125
    _state = HEADER_PARSED;
197
125
    return Status::OK();
198
126
}
Unexecuted instantiation: _ZN5doris10PageReaderILb1ELb1EE17parse_page_headerEv
_ZN5doris10PageReaderILb1ELb0EE17parse_page_headerEv
Line
Count
Source
77
6
Status PageReader<IN_COLLECTION, OFFSET_INDEX>::parse_page_header() {
78
6
    if (_state == HEADER_PARSED) {
79
3
        return Status::OK();
80
3
    }
81
3
    if (UNLIKELY(_offset < _start_offset || _offset >= _end_offset)) {
82
0
        return Status::IOError("Out-of-bounds Access");
83
0
    }
84
3
    if (UNLIKELY(_offset != _next_header_offset)) {
85
0
        return Status::IOError("Wrong header position, should seek to a page header first");
86
0
    }
87
3
    if (UNLIKELY(_state != INITIALIZED)) {
88
0
        return Status::IOError("Should skip or load current page to get next page");
89
0
    }
90
91
3
    _page_statistics.page_read_counter += 1;
92
93
    // Parse page header from file; header bytes are saved for possible cache insertion
94
3
    const uint8_t* page_header_buf = nullptr;
95
3
    size_t max_size = _end_offset - _offset;
96
3
    size_t header_size = std::min(INIT_PAGE_HEADER_SIZE, max_size);
97
3
    const size_t MAX_PAGE_HEADER_SIZE = config::parquet_header_max_size_mb << 20;
98
3
    uint32_t real_header_size = 0;
99
100
    // Try a header-only lookup in the page cache. Cached pages store
101
    // header + optional v2 levels + uncompressed payload, so we can
102
    // parse the page header directly from the cached bytes and avoid
103
    // a file read for the header.
104
3
    if (_page_read_ctx.enable_parquet_file_page_cache && !config::disable_storage_page_cache &&
105
3
        StoragePageCache::instance() != nullptr) {
106
3
        PageCacheHandle handle;
107
3
        StoragePageCache::CacheKey key = make_page_cache_key(static_cast<int64_t>(_offset));
108
3
        if (StoragePageCache::instance()->lookup(key, &handle, segment_v2::DATA_PAGE)) {
109
            // Parse header directly from cached data
110
1
            _page_cache_handle = std::move(handle);
111
1
            Slice s = _page_cache_handle.data();
112
1
            real_header_size = cast_set<uint32_t>(s.size);
113
1
            SCOPED_RAW_TIMER(&_page_statistics.decode_header_time);
114
1
            auto st = deserialize_thrift_msg(reinterpret_cast<const uint8_t*>(s.data),
115
1
                                             &real_header_size, true, &_cur_page_header);
116
1
            if (!st.ok()) return st;
117
            // Increment page cache counters for a true cache hit on header+payload
118
1
            _page_statistics.page_cache_hit_counter += 1;
119
            // Detect whether the cached payload is compressed or decompressed and record
120
1
            bool is_cache_payload_decompressed =
121
1
                    should_cache_decompressed(&_cur_page_header, _metadata);
122
123
1
            if (is_cache_payload_decompressed) {
124
1
                _page_statistics.page_cache_decompressed_hit_counter += 1;
125
1
            } else {
126
0
                _page_statistics.page_cache_compressed_hit_counter += 1;
127
0
            }
128
129
1
            _is_cache_payload_decompressed = is_cache_payload_decompressed;
130
131
1
            if constexpr (OFFSET_INDEX == false) {
132
1
                if (is_header_v2()) {
133
0
                    _end_row = _start_row + _cur_page_header.data_page_header_v2.num_rows;
134
1
                } else if constexpr (!IN_COLLECTION) {
135
1
                    _end_row = _start_row + _cur_page_header.data_page_header.num_values;
136
1
                }
137
1
            }
138
139
            // Save header bytes for later use (e.g., to insert updated cache entries)
140
1
            _header_buf.assign(s.data, s.data + real_header_size);
141
1
            _last_header_size = real_header_size;
142
1
            _page_statistics.parse_page_header_num++;
143
1
            _offset += real_header_size;
144
1
            _next_header_offset = _offset + _cur_page_header.compressed_page_size;
145
1
            _state = HEADER_PARSED;
146
1
            return Status::OK();
147
2
        } else {
148
2
            _page_statistics.page_cache_missing_counter += 1;
149
            // Clear any existing cache handle on miss to avoid holding stale handle
150
2
            _page_cache_handle = PageCacheHandle();
151
2
        }
152
3
    }
153
    // NOTE: page cache lookup for *decompressed* page data is handled in
154
    // ColumnChunkReader::load_page_data(). PageReader should only be
155
    // responsible for parsing the header bytes from the file and saving
156
    // them in `_header_buf` for possible later insertion into the cache.
157
2
    while (true) {
158
2
        if (UNLIKELY(_io_ctx && _io_ctx->should_stop)) {
159
0
            return Status::EndOfFile("stop");
160
0
        }
161
2
        header_size = std::min(header_size, max_size);
162
2
        {
163
2
            SCOPED_RAW_TIMER(&_page_statistics.read_page_header_time);
164
2
            RETURN_IF_ERROR(_reader->read_bytes(&page_header_buf, _offset, header_size, _io_ctx));
165
2
        }
166
2
        real_header_size = cast_set<uint32_t>(header_size);
167
2
        SCOPED_RAW_TIMER(&_page_statistics.decode_header_time);
168
2
        auto st =
169
2
                deserialize_thrift_msg(page_header_buf, &real_header_size, true, &_cur_page_header);
170
2
        if (st.ok()) {
171
2
            break;
172
2
        }
173
0
        if (_offset + header_size >= _end_offset || real_header_size > MAX_PAGE_HEADER_SIZE) {
174
0
            return Status::IOError(
175
0
                    "Failed to deserialize parquet page header. offset: {}, "
176
0
                    "header size: {}, end offset: {}, real header size: {}",
177
0
                    _offset, header_size, _end_offset, real_header_size);
178
0
        }
179
0
        header_size <<= 2;
180
0
    }
181
182
2
    if constexpr (OFFSET_INDEX == false) {
183
2
        if (is_header_v2()) {
184
0
            _end_row = _start_row + _cur_page_header.data_page_header_v2.num_rows;
185
2
        } else if constexpr (!IN_COLLECTION) {
186
2
            _end_row = _start_row + _cur_page_header.data_page_header.num_values;
187
2
        }
188
2
    }
189
190
    // Save header bytes for possible cache insertion later
191
2
    _header_buf.assign(page_header_buf, page_header_buf + real_header_size);
192
2
    _last_header_size = real_header_size;
193
2
    _page_statistics.parse_page_header_num++;
194
2
    _offset += real_header_size;
195
2
    _next_header_offset = _offset + _cur_page_header.compressed_page_size;
196
2
    _state = HEADER_PARSED;
197
2
    return Status::OK();
198
2
}
_ZN5doris10PageReaderILb0ELb1EE17parse_page_headerEv
Line
Count
Source
77
6
Status PageReader<IN_COLLECTION, OFFSET_INDEX>::parse_page_header() {
78
6
    if (_state == HEADER_PARSED) {
79
0
        return Status::OK();
80
0
    }
81
6
    if (UNLIKELY(_offset < _start_offset || _offset >= _end_offset)) {
82
0
        return Status::IOError("Out-of-bounds Access");
83
0
    }
84
6
    if (UNLIKELY(_offset != _next_header_offset)) {
85
0
        return Status::IOError("Wrong header position, should seek to a page header first");
86
0
    }
87
6
    if (UNLIKELY(_state != INITIALIZED)) {
88
0
        return Status::IOError("Should skip or load current page to get next page");
89
0
    }
90
91
6
    _page_statistics.page_read_counter += 1;
92
93
    // Parse page header from file; header bytes are saved for possible cache insertion
94
6
    const uint8_t* page_header_buf = nullptr;
95
6
    size_t max_size = _end_offset - _offset;
96
6
    size_t header_size = std::min(INIT_PAGE_HEADER_SIZE, max_size);
97
6
    const size_t MAX_PAGE_HEADER_SIZE = config::parquet_header_max_size_mb << 20;
98
6
    uint32_t real_header_size = 0;
99
100
    // Try a header-only lookup in the page cache. Cached pages store
101
    // header + optional v2 levels + uncompressed payload, so we can
102
    // parse the page header directly from the cached bytes and avoid
103
    // a file read for the header.
104
6
    if (_page_read_ctx.enable_parquet_file_page_cache && !config::disable_storage_page_cache &&
105
6
        StoragePageCache::instance() != nullptr) {
106
0
        PageCacheHandle handle;
107
0
        StoragePageCache::CacheKey key = make_page_cache_key(static_cast<int64_t>(_offset));
108
0
        if (StoragePageCache::instance()->lookup(key, &handle, segment_v2::DATA_PAGE)) {
109
            // Parse header directly from cached data
110
0
            _page_cache_handle = std::move(handle);
111
0
            Slice s = _page_cache_handle.data();
112
0
            real_header_size = cast_set<uint32_t>(s.size);
113
0
            SCOPED_RAW_TIMER(&_page_statistics.decode_header_time);
114
0
            auto st = deserialize_thrift_msg(reinterpret_cast<const uint8_t*>(s.data),
115
0
                                             &real_header_size, true, &_cur_page_header);
116
0
            if (!st.ok()) return st;
117
            // Increment page cache counters for a true cache hit on header+payload
118
0
            _page_statistics.page_cache_hit_counter += 1;
119
            // Detect whether the cached payload is compressed or decompressed and record
120
0
            bool is_cache_payload_decompressed =
121
0
                    should_cache_decompressed(&_cur_page_header, _metadata);
122
123
0
            if (is_cache_payload_decompressed) {
124
0
                _page_statistics.page_cache_decompressed_hit_counter += 1;
125
0
            } else {
126
0
                _page_statistics.page_cache_compressed_hit_counter += 1;
127
0
            }
128
129
0
            _is_cache_payload_decompressed = is_cache_payload_decompressed;
130
131
            if constexpr (OFFSET_INDEX == false) {
132
                if (is_header_v2()) {
133
                    _end_row = _start_row + _cur_page_header.data_page_header_v2.num_rows;
134
                } else if constexpr (!IN_COLLECTION) {
135
                    _end_row = _start_row + _cur_page_header.data_page_header.num_values;
136
                }
137
            }
138
139
            // Save header bytes for later use (e.g., to insert updated cache entries)
140
0
            _header_buf.assign(s.data, s.data + real_header_size);
141
0
            _last_header_size = real_header_size;
142
0
            _page_statistics.parse_page_header_num++;
143
0
            _offset += real_header_size;
144
0
            _next_header_offset = _offset + _cur_page_header.compressed_page_size;
145
0
            _state = HEADER_PARSED;
146
0
            return Status::OK();
147
0
        } else {
148
0
            _page_statistics.page_cache_missing_counter += 1;
149
            // Clear any existing cache handle on miss to avoid holding stale handle
150
0
            _page_cache_handle = PageCacheHandle();
151
0
        }
152
0
    }
153
    // NOTE: page cache lookup for *decompressed* page data is handled in
154
    // ColumnChunkReader::load_page_data(). PageReader should only be
155
    // responsible for parsing the header bytes from the file and saving
156
    // them in `_header_buf` for possible later insertion into the cache.
157
6
    while (true) {
158
6
        if (UNLIKELY(_io_ctx && _io_ctx->should_stop)) {
159
0
            return Status::EndOfFile("stop");
160
0
        }
161
6
        header_size = std::min(header_size, max_size);
162
6
        {
163
6
            SCOPED_RAW_TIMER(&_page_statistics.read_page_header_time);
164
6
            RETURN_IF_ERROR(_reader->read_bytes(&page_header_buf, _offset, header_size, _io_ctx));
165
6
        }
166
6
        real_header_size = cast_set<uint32_t>(header_size);
167
6
        SCOPED_RAW_TIMER(&_page_statistics.decode_header_time);
168
6
        auto st =
169
6
                deserialize_thrift_msg(page_header_buf, &real_header_size, true, &_cur_page_header);
170
6
        if (st.ok()) {
171
6
            break;
172
6
        }
173
0
        if (_offset + header_size >= _end_offset || real_header_size > MAX_PAGE_HEADER_SIZE) {
174
0
            return Status::IOError(
175
0
                    "Failed to deserialize parquet page header. offset: {}, "
176
0
                    "header size: {}, end offset: {}, real header size: {}",
177
0
                    _offset, header_size, _end_offset, real_header_size);
178
0
        }
179
0
        header_size <<= 2;
180
0
    }
181
182
    if constexpr (OFFSET_INDEX == false) {
183
        if (is_header_v2()) {
184
            _end_row = _start_row + _cur_page_header.data_page_header_v2.num_rows;
185
        } else if constexpr (!IN_COLLECTION) {
186
            _end_row = _start_row + _cur_page_header.data_page_header.num_values;
187
        }
188
    }
189
190
    // Save header bytes for possible cache insertion later
191
6
    _header_buf.assign(page_header_buf, page_header_buf + real_header_size);
192
6
    _last_header_size = real_header_size;
193
6
    _page_statistics.parse_page_header_num++;
194
6
    _offset += real_header_size;
195
6
    _next_header_offset = _offset + _cur_page_header.compressed_page_size;
196
6
    _state = HEADER_PARSED;
197
6
    return Status::OK();
198
6
}
_ZN5doris10PageReaderILb0ELb0EE17parse_page_headerEv
Line
Count
Source
77
372
Status PageReader<IN_COLLECTION, OFFSET_INDEX>::parse_page_header() {
78
372
    if (_state == HEADER_PARSED) {
79
137
        return Status::OK();
80
137
    }
81
235
    if (UNLIKELY(_offset < _start_offset || _offset >= _end_offset)) {
82
0
        return Status::IOError("Out-of-bounds Access");
83
0
    }
84
235
    if (UNLIKELY(_offset != _next_header_offset)) {
85
0
        return Status::IOError("Wrong header position, should seek to a page header first");
86
0
    }
87
235
    if (UNLIKELY(_state != INITIALIZED)) {
88
0
        return Status::IOError("Should skip or load current page to get next page");
89
0
    }
90
91
235
    _page_statistics.page_read_counter += 1;
92
93
    // Parse page header from file; header bytes are saved for possible cache insertion
94
235
    const uint8_t* page_header_buf = nullptr;
95
235
    size_t max_size = _end_offset - _offset;
96
235
    size_t header_size = std::min(INIT_PAGE_HEADER_SIZE, max_size);
97
235
    const size_t MAX_PAGE_HEADER_SIZE = config::parquet_header_max_size_mb << 20;
98
235
    uint32_t real_header_size = 0;
99
100
    // Try a header-only lookup in the page cache. Cached pages store
101
    // header + optional v2 levels + uncompressed payload, so we can
102
    // parse the page header directly from the cached bytes and avoid
103
    // a file read for the header.
104
235
    if (_page_read_ctx.enable_parquet_file_page_cache && !config::disable_storage_page_cache &&
105
235
        StoragePageCache::instance() != nullptr) {
106
223
        PageCacheHandle handle;
107
223
        StoragePageCache::CacheKey key = make_page_cache_key(static_cast<int64_t>(_offset));
108
223
        if (StoragePageCache::instance()->lookup(key, &handle, segment_v2::DATA_PAGE)) {
109
            // Parse header directly from cached data
110
117
            _page_cache_handle = std::move(handle);
111
117
            Slice s = _page_cache_handle.data();
112
117
            real_header_size = cast_set<uint32_t>(s.size);
113
117
            SCOPED_RAW_TIMER(&_page_statistics.decode_header_time);
114
117
            auto st = deserialize_thrift_msg(reinterpret_cast<const uint8_t*>(s.data),
115
117
                                             &real_header_size, true, &_cur_page_header);
116
117
            if (!st.ok()) return st;
117
            // Increment page cache counters for a true cache hit on header+payload
118
117
            _page_statistics.page_cache_hit_counter += 1;
119
            // Detect whether the cached payload is compressed or decompressed and record
120
117
            bool is_cache_payload_decompressed =
121
117
                    should_cache_decompressed(&_cur_page_header, _metadata);
122
123
117
            if (is_cache_payload_decompressed) {
124
117
                _page_statistics.page_cache_decompressed_hit_counter += 1;
125
117
            } else {
126
0
                _page_statistics.page_cache_compressed_hit_counter += 1;
127
0
            }
128
129
117
            _is_cache_payload_decompressed = is_cache_payload_decompressed;
130
131
117
            if constexpr (OFFSET_INDEX == false) {
132
117
                if (is_header_v2()) {
133
11
                    _end_row = _start_row + _cur_page_header.data_page_header_v2.num_rows;
134
106
                } else if constexpr (!IN_COLLECTION) {
135
106
                    _end_row = _start_row + _cur_page_header.data_page_header.num_values;
136
106
                }
137
117
            }
138
139
            // Save header bytes for later use (e.g., to insert updated cache entries)
140
117
            _header_buf.assign(s.data, s.data + real_header_size);
141
117
            _last_header_size = real_header_size;
142
117
            _page_statistics.parse_page_header_num++;
143
117
            _offset += real_header_size;
144
117
            _next_header_offset = _offset + _cur_page_header.compressed_page_size;
145
117
            _state = HEADER_PARSED;
146
117
            return Status::OK();
147
117
        } else {
148
106
            _page_statistics.page_cache_missing_counter += 1;
149
            // Clear any existing cache handle on miss to avoid holding stale handle
150
106
            _page_cache_handle = PageCacheHandle();
151
106
        }
152
223
    }
153
    // NOTE: page cache lookup for *decompressed* page data is handled in
154
    // ColumnChunkReader::load_page_data(). PageReader should only be
155
    // responsible for parsing the header bytes from the file and saving
156
    // them in `_header_buf` for possible later insertion into the cache.
157
120
    while (true) {
158
120
        if (UNLIKELY(_io_ctx && _io_ctx->should_stop)) {
159
0
            return Status::EndOfFile("stop");
160
0
        }
161
120
        header_size = std::min(header_size, max_size);
162
120
        {
163
120
            SCOPED_RAW_TIMER(&_page_statistics.read_page_header_time);
164
120
            RETURN_IF_ERROR(_reader->read_bytes(&page_header_buf, _offset, header_size, _io_ctx));
165
120
        }
166
119
        real_header_size = cast_set<uint32_t>(header_size);
167
119
        SCOPED_RAW_TIMER(&_page_statistics.decode_header_time);
168
119
        auto st =
169
119
                deserialize_thrift_msg(page_header_buf, &real_header_size, true, &_cur_page_header);
170
119
        if (st.ok()) {
171
117
            break;
172
117
        }
173
2
        if (_offset + header_size >= _end_offset || real_header_size > MAX_PAGE_HEADER_SIZE) {
174
0
            return Status::IOError(
175
0
                    "Failed to deserialize parquet page header. offset: {}, "
176
0
                    "header size: {}, end offset: {}, real header size: {}",
177
0
                    _offset, header_size, _end_offset, real_header_size);
178
0
        }
179
2
        header_size <<= 2;
180
2
    }
181
182
117
    if constexpr (OFFSET_INDEX == false) {
183
117
        if (is_header_v2()) {
184
9
            _end_row = _start_row + _cur_page_header.data_page_header_v2.num_rows;
185
108
        } else if constexpr (!IN_COLLECTION) {
186
108
            _end_row = _start_row + _cur_page_header.data_page_header.num_values;
187
108
        }
188
117
    }
189
190
    // Save header bytes for possible cache insertion later
191
117
    _header_buf.assign(page_header_buf, page_header_buf + real_header_size);
192
117
    _last_header_size = real_header_size;
193
117
    _page_statistics.parse_page_header_num++;
194
117
    _offset += real_header_size;
195
117
    _next_header_offset = _offset + _cur_page_header.compressed_page_size;
196
117
    _state = HEADER_PARSED;
197
117
    return Status::OK();
198
118
}
199
200
template <bool IN_COLLECTION, bool OFFSET_INDEX>
201
116
Status PageReader<IN_COLLECTION, OFFSET_INDEX>::get_page_data(Slice& slice) {
202
116
    if (UNLIKELY(_state != HEADER_PARSED)) {
203
0
        return Status::IOError("Should generate page header first to load current page data");
204
0
    }
205
116
    if (UNLIKELY(_io_ctx && _io_ctx->should_stop)) {
206
0
        return Status::EndOfFile("stop");
207
0
    }
208
116
    slice.size = _cur_page_header.compressed_page_size;
209
116
    RETURN_IF_ERROR(_reader->read_bytes(slice, _offset, _io_ctx));
210
115
    _offset += slice.size;
211
115
    _state = DATA_LOADED;
212
115
    return Status::OK();
213
116
}
Unexecuted instantiation: _ZN5doris10PageReaderILb1ELb1EE13get_page_dataERNS_5SliceE
_ZN5doris10PageReaderILb1ELb0EE13get_page_dataERNS_5SliceE
Line
Count
Source
201
2
Status PageReader<IN_COLLECTION, OFFSET_INDEX>::get_page_data(Slice& slice) {
202
2
    if (UNLIKELY(_state != HEADER_PARSED)) {
203
0
        return Status::IOError("Should generate page header first to load current page data");
204
0
    }
205
2
    if (UNLIKELY(_io_ctx && _io_ctx->should_stop)) {
206
0
        return Status::EndOfFile("stop");
207
0
    }
208
2
    slice.size = _cur_page_header.compressed_page_size;
209
2
    RETURN_IF_ERROR(_reader->read_bytes(slice, _offset, _io_ctx));
210
2
    _offset += slice.size;
211
2
    _state = DATA_LOADED;
212
2
    return Status::OK();
213
2
}
_ZN5doris10PageReaderILb0ELb1EE13get_page_dataERNS_5SliceE
Line
Count
Source
201
2
Status PageReader<IN_COLLECTION, OFFSET_INDEX>::get_page_data(Slice& slice) {
202
2
    if (UNLIKELY(_state != HEADER_PARSED)) {
203
0
        return Status::IOError("Should generate page header first to load current page data");
204
0
    }
205
2
    if (UNLIKELY(_io_ctx && _io_ctx->should_stop)) {
206
0
        return Status::EndOfFile("stop");
207
0
    }
208
2
    slice.size = _cur_page_header.compressed_page_size;
209
2
    RETURN_IF_ERROR(_reader->read_bytes(slice, _offset, _io_ctx));
210
2
    _offset += slice.size;
211
2
    _state = DATA_LOADED;
212
2
    return Status::OK();
213
2
}
_ZN5doris10PageReaderILb0ELb0EE13get_page_dataERNS_5SliceE
Line
Count
Source
201
112
Status PageReader<IN_COLLECTION, OFFSET_INDEX>::get_page_data(Slice& slice) {
202
112
    if (UNLIKELY(_state != HEADER_PARSED)) {
203
0
        return Status::IOError("Should generate page header first to load current page data");
204
0
    }
205
112
    if (UNLIKELY(_io_ctx && _io_ctx->should_stop)) {
206
0
        return Status::EndOfFile("stop");
207
0
    }
208
112
    slice.size = _cur_page_header.compressed_page_size;
209
112
    RETURN_IF_ERROR(_reader->read_bytes(slice, _offset, _io_ctx));
210
111
    _offset += slice.size;
211
111
    _state = DATA_LOADED;
212
111
    return Status::OK();
213
112
}
214
215
template class PageReader<true, true>;
216
template class PageReader<true, false>;
217
template class PageReader<false, true>;
218
template class PageReader<false, false>;
219
220
} // namespace doris