Coverage Report

Created: 2026-08-07 19:51

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/storage/index/indexed_column_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 "storage/index/indexed_column_reader.h"
19
20
#include <gen_cpp/segment_v2.pb.h>
21
22
#include <algorithm>
23
24
#include "common/status.h"
25
#include "io/fs/file_reader.h"
26
#include "io/io_common.h"
27
#include "storage/key_coder.h"
28
#include "storage/olap_common.h"
29
#include "storage/segment/encoding_info.h" // for EncodingInfo
30
#include "storage/segment/options.h"
31
#include "storage/segment/page_decoder.h"
32
#include "storage/segment/page_io.h"
33
#include "storage/types.h"
34
#include "util/block_compression.h"
35
#include "util/bvar_helper.h"
36
37
namespace doris {
38
using namespace ErrorCode;
39
namespace segment_v2 {
40
41
static bvar::Adder<uint64_t> g_index_reader_bytes("doris_pk", "index_reader_bytes");
42
static bvar::Adder<uint64_t> g_index_reader_compressed_bytes("doris_pk",
43
                                                             "index_reader_compressed_bytes");
44
static bvar::PerSecond<bvar::Adder<uint64_t>> g_index_reader_bytes_per_second(
45
        "doris_pk", "index_reader_bytes_per_second", &g_index_reader_bytes, 60);
46
static bvar::Adder<uint64_t> g_index_reader_pages("doris_pk", "index_reader_pages");
47
static bvar::PerSecond<bvar::Adder<uint64_t>> g_index_reader_pages_per_second(
48
        "doris_pk", "index_reader_pages_per_second", &g_index_reader_pages, 60);
49
static bvar::Adder<uint64_t> g_index_reader_cached_pages("doris_pk", "index_reader_cached_pages");
50
static bvar::PerSecond<bvar::Adder<uint64_t>> g_index_reader_cached_pages_per_second(
51
        "doris_pk", "index_reader_cached_pages_per_second", &g_index_reader_cached_pages, 60);
52
static bvar::Adder<uint64_t> g_index_reader_seek_count("doris_pk", "index_reader_seek_count");
53
static bvar::PerSecond<bvar::Adder<uint64_t>> g_index_reader_seek_per_second(
54
        "doris_pk", "index_reader_seek_per_second", &g_index_reader_seek_count, 60);
55
static bvar::Adder<uint64_t> g_index_reader_pk_pages("doris_pk", "index_reader_pk_pages");
56
static bvar::PerSecond<bvar::Adder<uint64_t>> g_index_reader_pk_bytes_per_second(
57
        "doris_pk", "index_reader_pk_pages_per_second", &g_index_reader_pk_pages, 60);
58
59
777
int64_t IndexedColumnReader::get_metadata_size() const {
60
777
    return sizeof(IndexedColumnReader) + _meta.ByteSizeLong();
61
777
}
62
63
Status IndexedColumnReader::load(bool use_page_cache, bool kept_in_memory,
64
                                 OlapReaderStatistics* index_load_stats,
65
777
                                 const io::IOContext* io_ctx) {
66
777
    _use_page_cache = use_page_cache;
67
777
    _kept_in_memory = kept_in_memory;
68
69
777
    _type = (FieldType)_meta.data_type();
70
777
    if (!is_scalar_type(_type)) {
71
0
        return Status::NotSupported("unsupported typeinfo, type={}", _meta.data_type());
72
0
    }
73
777
    RETURN_IF_ERROR(EncodingInfo::get(_type, _meta.encoding(), &_encoding_info));
74
777
    _value_key_coder = get_key_coder(_type);
75
76
    // read and parse ordinal index page when exists
77
777
    if (_meta.has_ordinal_index_meta()) {
78
777
        if (_meta.ordinal_index_meta().is_root_data_page()) {
79
756
            _sole_data_page = PagePointer(_meta.ordinal_index_meta().root_page());
80
756
        } else {
81
21
            RETURN_IF_ERROR(load_index_page(_meta.ordinal_index_meta().root_page(),
82
21
                                            &_ordinal_index_page_handle,
83
21
                                            _ordinal_index_reader.get(), index_load_stats, io_ctx));
84
21
            _has_index_page = true;
85
21
        }
86
777
    }
87
88
    // read and parse value index page when exists
89
777
    if (_meta.has_value_index_meta()) {
90
522
        if (_meta.value_index_meta().is_root_data_page()) {
91
501
            _sole_data_page = PagePointer(_meta.value_index_meta().root_page());
92
501
        } else {
93
21
            RETURN_IF_ERROR(load_index_page(_meta.value_index_meta().root_page(),
94
21
                                            &_value_index_page_handle, _value_index_reader.get(),
95
21
                                            index_load_stats, io_ctx));
96
21
            _has_index_page = true;
97
21
        }
98
522
    }
99
777
    _num_values = _meta.num_values();
100
101
777
    update_metadata_size();
102
777
    return Status::OK();
103
777
}
104
105
Status IndexedColumnReader::load_index_page(const PagePointerPB& pp, PageHandle* handle,
106
                                            IndexPageReader* reader,
107
                                            OlapReaderStatistics* index_load_stats,
108
42
                                            const io::IOContext* io_ctx) {
109
42
    Slice body;
110
42
    PageFooterPB footer;
111
42
    BlockCompressionCodec* local_compress_codec;
112
42
    RETURN_IF_ERROR(get_block_compression_codec(_meta.compression(), &local_compress_codec));
113
42
    RETURN_IF_ERROR(read_page(PagePointer(pp), handle, &body, &footer, INDEX_PAGE,
114
42
                              local_compress_codec, false, index_load_stats, io_ctx));
115
42
    RETURN_IF_ERROR(reader->parse(body, footer.index_page_footer()));
116
42
    _mem_size += body.get_size();
117
42
    return Status::OK();
118
42
}
119
120
Status IndexedColumnReader::read_page(const PagePointer& pp, PageHandle* handle, Slice* body,
121
                                      PageFooterPB* footer, PageTypePB type,
122
                                      BlockCompressionCodec* codec, bool pre_decode,
123
                                      OlapReaderStatistics* stats,
124
2.02k
                                      const io::IOContext* io_ctx) const {
125
2.02k
    OlapReaderStatistics tmp_stats;
126
2.02k
    OlapReaderStatistics* stats_ptr = stats != nullptr ? stats : &tmp_stats;
127
2.02k
    io::IOContext page_io_ctx = io_ctx != nullptr ? *io_ctx : io::IOContext {};
128
2.02k
    page_io_ctx.is_index_data = true;
129
2.02k
    page_io_ctx.file_cache_stats = &stats_ptr->file_cache_stats;
130
2.02k
    PageReadOptions opts(page_io_ctx);
131
2.02k
    opts.use_page_cache = _use_page_cache;
132
2.02k
    opts.kept_in_memory = _kept_in_memory;
133
2.02k
    opts.pre_decode = pre_decode;
134
2.02k
    opts.type = type;
135
2.02k
    opts.file_reader = _file_reader.get();
136
2.02k
    opts.page_pointer = pp;
137
2.02k
    opts.codec = codec;
138
2.02k
    opts.stats = stats_ptr;
139
2.02k
    opts.encoding_info = _encoding_info;
140
141
2.02k
    if (_is_pk_index) {
142
489
        opts.type = PRIMARY_KEY_INDEX_PAGE;
143
489
    }
144
2.02k
    auto st = PageIO::read_and_decompress_page(opts, handle, body, footer);
145
2.02k
    g_index_reader_compressed_bytes << pp.size;
146
2.02k
    g_index_reader_bytes << footer->uncompressed_size();
147
2.02k
    g_index_reader_pages << 1;
148
2.02k
    g_index_reader_cached_pages << tmp_stats.cached_pages_num;
149
2.02k
    return st;
150
2.02k
}
151
152
777
IndexedColumnReader::~IndexedColumnReader() = default;
153
154
///////////////////////////////////////////////////////////////////////////////
155
156
1.98k
Status IndexedColumnIterator::_read_data_page(const PagePointer& pp) {
157
1.98k
    Status status;
158
    // there is not init() for IndexedColumnIterator, so do it here
159
1.98k
    if (!_compress_codec) {
160
1.91k
        RETURN_IF_ERROR(get_block_compression_codec(_reader->get_compression(), &_compress_codec));
161
1.91k
    }
162
163
1.98k
    PageHandle handle;
164
1.98k
    Slice body;
165
1.98k
    PageFooterPB footer;
166
1.98k
    RETURN_IF_ERROR(_reader->read_page(pp, &handle, &body, &footer, DATA_PAGE, _compress_codec,
167
1.98k
                                       true, _stats, _io_ctx));
168
    // parse data page
169
    // note that page_index is not used in IndexedColumnIterator, so we pass 0
170
1.98k
    PageDecoderOptions opts;
171
1.98k
    opts.need_check_bitmap = false;
172
1.98k
    status = ParsedPage::create(std::move(handle), body, footer.data_page_footer(),
173
1.98k
                                _reader->encoding_info(), pp, 0, &_data_page, opts);
174
1.98k
    if (!status.ok()) {
175
0
        LOG(WARNING) << "failed to create ParsedPage in IndexedColumnIterator, file="
176
0
                     << _reader->_file_reader->path().native() << ", page_offset=" << pp.offset
177
0
                     << ", page_size=" << pp.size << ", error=" << status;
178
0
    }
179
1.98k
    DCHECK(_reader->_meta.ordinal_index_meta().is_root_data_page()
180
1.98k
                   ? _reader->_meta.num_values() == _data_page.num_rows
181
1.98k
                   : true);
182
1.98k
    return status;
183
1.98k
}
184
185
886
Status IndexedColumnIterator::seek_to_ordinal(ordinal_t idx) {
186
886
    DCHECK(idx <= _reader->num_values());
187
188
886
    if (!_reader->support_ordinal_seek()) {
189
0
        return Status::NotSupported("no ordinal index");
190
0
    }
191
192
    // it's ok to seek past the last value
193
886
    if (idx == _reader->num_values()) {
194
1
        _current_ordinal = idx;
195
1
        _seeked = true;
196
1
        return Status::OK();
197
1
    }
198
199
885
    if (!_data_page || !_data_page.contains(idx)) {
200
        // need to read the data page containing row at idx
201
715
        if (_reader->_has_index_page) {
202
38
            std::string key;
203
38
            KeyCoderTraits<FieldType::OLAP_FIELD_TYPE_UNSIGNED_BIGINT>::full_encode_ascending(&idx,
204
38
                                                                                              &key);
205
38
            RETURN_IF_ERROR(_ordinal_iter.seek_at_or_before(key));
206
38
            RETURN_IF_ERROR(_read_data_page(_ordinal_iter.current_page_pointer()));
207
38
            _current_iter = &_ordinal_iter;
208
677
        } else {
209
677
            RETURN_IF_ERROR(_read_data_page(_reader->_sole_data_page));
210
677
        }
211
715
    }
212
213
885
    ordinal_t offset_in_page = idx - _data_page.first_ordinal;
214
885
    RETURN_IF_ERROR(_data_page.data_decoder->seek_to_position_in_page(offset_in_page));
215
885
    DCHECK(offset_in_page == _data_page.data_decoder->current_index());
216
885
    _data_page.offset_in_page = offset_in_page;
217
885
    _current_ordinal = idx;
218
885
    _seeked = true;
219
885
    return Status::OK();
220
885
}
221
222
6.43k
Status IndexedColumnIterator::seek_at_or_after(const void* key, bool* exact_match) {
223
6.43k
    if (!_reader->support_value_seek()) {
224
0
        return Status::NotSupported("no value index");
225
0
    }
226
227
6.43k
    if (_reader->num_values() == 0) {
228
0
        return Status::Error<ErrorCode::ENTRY_NOT_FOUND>("value index is empty ");
229
0
    }
230
231
6.43k
    g_index_reader_seek_count << 1;
232
233
6.43k
    bool load_data_page = false;
234
6.43k
    PagePointer data_page_pp;
235
6.43k
    if (_reader->_has_index_page) {
236
        // seek index to determine the data page to seek
237
102
        std::string encoded_key;
238
102
        _reader->_value_key_coder->full_encode_ascending(key, &encoded_key);
239
102
        Status st = _value_iter.seek_at_or_before(encoded_key);
240
102
        if (st.is<ENTRY_NOT_FOUND>()) {
241
            // all keys in page is greater than `encoded_key`, point to the first page.
242
            // otherwise, we may missing some pages.
243
            // For example, the predicate is `col1 > 2`, and the index page is [3,5,7].
244
            // so the `seek_at_or_before(2)` will return Status::Error<ENTRY_NOT_FOUND>().
245
            // But actually, we expect it to point to page `3`.
246
0
            _value_iter.seek_to_first();
247
102
        } else if (!st.ok()) {
248
0
            return st;
249
0
        }
250
102
        data_page_pp = _value_iter.current_page_pointer();
251
102
        _current_iter = &_value_iter;
252
102
        if (!_data_page || _data_page.page_pointer != data_page_pp) {
253
            // load when it's not the same with the current
254
24
            load_data_page = true;
255
24
        }
256
6.33k
    } else if (!_data_page) {
257
        // no index page, load data page for the first time
258
1.21k
        load_data_page = true;
259
1.21k
        data_page_pp = PagePointer(_reader->_sole_data_page);
260
1.21k
    }
261
262
6.43k
    if (load_data_page) {
263
1.23k
        RETURN_IF_ERROR(_read_data_page(data_page_pp));
264
1.23k
    }
265
266
    // seek inside data page
267
6.43k
    Status st = _data_page.data_decoder->seek_at_or_after_value(key, exact_match);
268
    // return the first row of next page when not found
269
6.43k
    if (st.is<ENTRY_NOT_FOUND>() && _reader->_has_index_page) {
270
4
        if (_value_iter.has_next()) {
271
3
            _seeked = true;
272
3
            *exact_match = false;
273
3
            _current_ordinal = _data_page.first_ordinal + _data_page.num_rows;
274
            // move offset to the end of the page
275
3
            _data_page.offset_in_page = _data_page.num_rows;
276
3
            return Status::OK();
277
3
        }
278
4
    }
279
6.43k
    RETURN_IF_ERROR(st);
280
5.89k
    _data_page.offset_in_page = _data_page.data_decoder->current_index();
281
5.89k
    _current_ordinal = _data_page.first_ordinal + _data_page.offset_in_page;
282
5.89k
    DCHECK(_data_page.contains(_current_ordinal));
283
5.89k
    _seeked = true;
284
5.89k
    return Status::OK();
285
6.43k
}
286
287
1.55k
Status IndexedColumnIterator::next_batch(size_t* n, MutableColumnPtr& dst) {
288
1.55k
    DCHECK(_seeked);
289
1.55k
    if (_current_ordinal == _reader->num_values()) {
290
0
        *n = 0;
291
0
        return Status::OK();
292
0
    }
293
294
1.55k
    size_t remaining = *n;
295
3.14k
    while (remaining > 0) {
296
1.58k
        if (!_data_page.has_remaining()) {
297
            // trying to read next data page
298
30
            if (!_reader->_has_index_page) {
299
0
                break; // no more data page
300
0
            }
301
30
            bool has_next = _current_iter->move_next();
302
30
            if (!has_next) {
303
0
                break; // no more data page
304
0
            }
305
30
            RETURN_IF_ERROR(_read_data_page(_current_iter->current_page_pointer()));
306
30
        }
307
308
1.58k
        size_t rows_to_read = std::min(_data_page.remaining(), remaining);
309
1.58k
        size_t rows_read = rows_to_read;
310
1.58k
        RETURN_IF_ERROR(_data_page.data_decoder->next_batch(&rows_read, dst));
311
1.58k
        DCHECK(rows_to_read == rows_read);
312
313
1.58k
        _data_page.offset_in_page += rows_read;
314
1.58k
        _current_ordinal += rows_read;
315
1.58k
        remaining -= rows_read;
316
1.58k
    }
317
1.55k
    *n -= remaining;
318
1.55k
    _seeked = false;
319
1.55k
    return Status::OK();
320
1.55k
}
321
322
} // namespace segment_v2
323
} // namespace doris