Coverage Report

Created: 2025-11-28 10:44

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/root/doris/be/src/olap/iterators.h
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
#pragma once
19
20
#include <cstddef>
21
#include <memory>
22
23
#include "common/status.h"
24
#include "io/io_common.h"
25
#include "olap/block_column_predicate.h"
26
#include "olap/column_predicate.h"
27
#include "olap/olap_common.h"
28
#include "olap/row_cursor.h"
29
#include "olap/rowset/segment_v2/ann_index/ann_topn_runtime.h"
30
#include "olap/rowset/segment_v2/row_ranges.h"
31
#include "olap/tablet_schema.h"
32
#include "runtime/runtime_state.h"
33
#include "vec/core/block.h"
34
#include "vec/exprs/score_runtime.h"
35
#include "vec/exprs/vexpr.h"
36
37
namespace doris {
38
39
class Schema;
40
class ColumnPredicate;
41
42
namespace vectorized {
43
struct IteratorRowRef;
44
};
45
46
namespace segment_v2 {
47
struct SubstreamIterator;
48
}
49
class StorageReadOptions {
50
public:
51
    struct KeyRange {
52
        KeyRange()
53
                : lower_key(nullptr),
54
                  include_lower(false),
55
                  upper_key(nullptr),
56
0
                  include_upper(false) {}
57
58
        KeyRange(const RowCursor* lower_key_, bool include_lower_, const RowCursor* upper_key_,
59
                 bool include_upper_)
60
0
                : lower_key(lower_key_),
61
0
                  include_lower(include_lower_),
62
0
                  upper_key(upper_key_),
63
0
                  include_upper(include_upper_) {}
64
65
        // the lower bound of the range, nullptr if not existed
66
        const RowCursor* lower_key = nullptr;
67
        // whether `lower_key` is included in the range
68
        bool include_lower;
69
        // the upper bound of the range, nullptr if not existed
70
        const RowCursor* upper_key = nullptr;
71
        // whether `upper_key` is included in the range
72
        bool include_upper;
73
74
0
        uint64_t get_digest(uint64_t seed) const {
75
0
            if (lower_key != nullptr) {
76
0
                auto key_str = lower_key->to_string();
77
0
                seed = HashUtil::hash64(key_str.c_str(), key_str.size(), seed);
78
0
                seed = HashUtil::hash64(&include_lower, sizeof(include_lower), seed);
79
0
            }
80
81
0
            if (upper_key != nullptr) {
82
0
                auto key_str = upper_key->to_string();
83
0
                seed = HashUtil::hash64(key_str.c_str(), key_str.size(), seed);
84
0
                seed = HashUtil::hash64(&include_upper, sizeof(include_upper), seed);
85
0
            }
86
87
0
            return seed;
88
0
        }
89
    };
90
91
    // reader's key ranges, empty if not existed.
92
    // used by short key index to filter row blocks
93
    std::vector<KeyRange> key_ranges;
94
95
    // For unique-key merge-on-write, the effect is similar to delete_conditions
96
    // that filters out rows that are deleted in realtime.
97
    // For a particular row, if delete_bitmap.contains(rowid) means that row is
98
    // marked deleted and invisible to user anymore.
99
    // segment_id -> roaring::Roaring*
100
    std::unordered_map<uint32_t, std::shared_ptr<roaring::Roaring>> delete_bitmap;
101
102
    std::shared_ptr<AndBlockColumnPredicate> delete_condition_predicates =
103
            AndBlockColumnPredicate::create_shared();
104
    // reader's column predicate, nullptr if not existed
105
    // used to fiter rows in row block
106
    std::vector<ColumnPredicate*> column_predicates;
107
    std::unordered_map<int32_t, std::shared_ptr<AndBlockColumnPredicate>> col_id_to_predicates;
108
    std::unordered_map<int32_t, std::vector<const ColumnPredicate*>> del_predicates_for_zone_map;
109
    TPushAggOp::type push_down_agg_type_opt = TPushAggOp::NONE;
110
111
    // REQUIRED (null is not allowed)
112
    OlapReaderStatistics* stats = nullptr;
113
    bool use_page_cache = false;
114
    uint32_t block_row_max = 4096 - 32; // see https://github.com/apache/doris/pull/11816
115
116
    TabletSchemaSPtr tablet_schema = nullptr;
117
    bool enable_unique_key_merge_on_write = false;
118
    bool record_rowids = false;
119
    std::vector<int> topn_filter_source_node_ids;
120
    int topn_filter_target_node_id = -1;
121
    // used for special optimization for query : ORDER BY key DESC LIMIT n
122
    bool read_orderby_key_reverse = false;
123
    // columns for orderby keys
124
    std::vector<uint32_t>* read_orderby_key_columns = nullptr;
125
    io::IOContext io_ctx;
126
    vectorized::VExpr* remaining_vconjunct_root = nullptr;
127
    std::vector<vectorized::VExprSPtr> remaining_conjunct_roots;
128
    vectorized::VExprContextSPtrs common_expr_ctxs_push_down;
129
    const std::set<int32_t>* output_columns = nullptr;
130
    // runtime state
131
    RuntimeState* runtime_state = nullptr;
132
    RowsetId rowset_id;
133
    Version version;
134
    int64_t tablet_id = 0;
135
    // slots that cast may be eliminated in storage layer
136
    std::map<std::string, vectorized::DataTypePtr> target_cast_type_for_variants;
137
    RowRanges row_ranges;
138
    size_t topn_limit = 0;
139
140
    std::map<ColumnId, vectorized::VExprContextSPtr> virtual_column_exprs;
141
    std::shared_ptr<segment_v2::AnnTopNRuntime> ann_topn_runtime;
142
    std::map<ColumnId, size_t> vir_cid_to_idx_in_block;
143
    std::map<size_t, vectorized::DataTypePtr> vir_col_idx_to_type;
144
145
    std::map<int32_t, TColumnAccessPaths> all_access_paths;
146
    std::map<int32_t, TColumnAccessPaths> predicate_access_paths;
147
148
    std::shared_ptr<vectorized::ScoreRuntime> score_runtime;
149
    CollectionStatisticsPtr collection_statistics;
150
151
    // Cache for sparse column data to avoid redundant reads
152
    // col_unique_id -> cached column_ptr
153
    std::unordered_map<int32_t, vectorized::ColumnPtr> sparse_column_cache;
154
155
    uint64_t condition_cache_digest = 0;
156
};
157
158
struct CompactionSampleInfo {
159
    int64_t bytes = 0;
160
    int64_t rows = 0;
161
    int64_t group_data_size = 0;
162
};
163
164
struct BlockWithSameBit {
165
    vectorized::Block* block;
166
    std::vector<bool>& same_bit;
167
168
201
    bool empty() const { return block->rows() == 0; }
169
};
170
171
class RowwiseIterator;
172
using RowwiseIteratorUPtr = std::unique_ptr<RowwiseIterator>;
173
class RowwiseIterator {
174
public:
175
11.7k
    RowwiseIterator() = default;
176
11.7k
    virtual ~RowwiseIterator() = default;
177
178
    // Initialize this iterator and make it ready to read with
179
    // input options.
180
    // Input options may contain scan range in which this scan.
181
    // Return Status::OK() if init successfully,
182
    // Return other error otherwise
183
0
    virtual Status init(const StorageReadOptions& opts) {
184
0
        return Status::InternalError("to be implemented, current class: " +
185
0
                                     demangle(typeid(*this).name()));
186
0
    }
187
188
0
    virtual Status init(const StorageReadOptions& opts, CompactionSampleInfo* sample_info) {
189
0
        return Status::InternalError("should not reach here, current class: " +
190
0
                                     demangle(typeid(*this).name()));
191
0
    }
192
193
    // If there is any valid data, this function will load data
194
    // into input batch with Status::OK() returned
195
    // If there is no data to read, will return Status::EndOfFile.
196
    // If other error happens, other error code will be returned.
197
0
    virtual Status next_batch(vectorized::Block* block) {
198
0
        return Status::InternalError("should not reach here, current class: " +
199
0
                                     demangle(typeid(*this).name()));
200
0
    }
201
202
0
    virtual Status next_batch(BlockWithSameBit* block_with_same_bit) {
203
0
        return Status::InternalError("should not reach here, current class: " +
204
0
                                     demangle(typeid(*this).name()));
205
0
    }
206
207
0
    virtual Status next_batch(vectorized::BlockView* block_view) {
208
0
        return Status::InternalError("should not reach here, current class: " +
209
0
                                     demangle(typeid(*this).name()));
210
0
    }
211
212
0
    virtual Status next_row(vectorized::IteratorRowRef* ref) {
213
0
        return Status::InternalError("should not reach here, current class: " +
214
0
                                     demangle(typeid(*this).name()));
215
0
    }
216
0
    virtual Status unique_key_next_row(vectorized::IteratorRowRef* ref) {
217
0
        return Status::InternalError("should not reach here, current class: " +
218
0
                                     demangle(typeid(*this).name()));
219
0
    }
220
221
0
    virtual bool is_merge_iterator() const { return false; }
222
223
0
    virtual Status current_block_row_locations(std::vector<RowLocation>* block_row_locations) {
224
0
        return Status::InternalError("should not reach here, current class: " +
225
0
                                     demangle(typeid(*this).name()));
226
0
    }
227
228
    // return schema for this Iterator
229
    virtual const Schema& schema() const = 0;
230
231
    // Return the data id such as segment id, used for keep the insert order when do
232
    // merge sort in priority queue
233
26.2k
    virtual uint64_t data_id() const { return 0; }
234
235
0
    virtual void update_profile(RuntimeProfile* profile) {}
236
    // return rows merged count by iterator
237
0
    virtual uint64_t merged_rows() const { return 0; }
238
239
    // return if it's an empty iterator
240
5.58k
    virtual bool empty() const { return false; }
241
};
242
243
} // namespace doris