Coverage Report

Created: 2026-06-04 06:55

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/storage/tablet/tablet_reader.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 <gen_cpp/Descriptors_types.h>
21
#include <gen_cpp/PaloInternalService_types.h>
22
#include <gen_cpp/PlanNodes_types.h>
23
#include <stddef.h>
24
#include <stdint.h>
25
26
#include <memory>
27
#include <set>
28
#include <string>
29
#include <unordered_set>
30
#include <utility>
31
#include <vector>
32
33
#include "agent/be_exec_version_manager.h"
34
#include "common/status.h"
35
#include "exprs/function_filter.h"
36
#include "io/io_common.h"
37
#include "storage/delete/delete_handler.h"
38
#include "storage/iterators.h"
39
#include "storage/olap_common.h"
40
#include "storage/olap_tuple.h"
41
#include "storage/predicate/filter_olap_param.h"
42
#include "storage/row_cursor.h"
43
#include "storage/rowid_conversion.h"
44
#include "storage/rowset/rowset.h"
45
#include "storage/rowset/rowset_meta.h"
46
#include "storage/rowset/rowset_reader.h"
47
#include "storage/rowset/rowset_reader_context.h"
48
#include "storage/tablet/base_tablet.h"
49
#include "storage/tablet/tablet_fwd.h"
50
51
namespace doris {
52
53
class RuntimeState;
54
class BitmapFilterFuncBase;
55
class BloomFilterFuncBase;
56
class ColumnPredicate;
57
class DeleteBitmap;
58
class HybridSetBase;
59
class RuntimeProfile;
60
61
class VCollectIterator;
62
class Block;
63
class VExpr;
64
class Arena;
65
class VExprContext;
66
67
// Used to compare row with input scan key. Scan key only contains key columns,
68
// row contains all key columns, which is superset of key columns.
69
// So we should compare the common prefix columns of lhs and rhs.
70
//
71
// NOTE: if you are not sure if you can use it, please don't use this function.
72
1.43M
inline int compare_row_key(const RowCursor& lhs, const RowCursor& rhs) {
73
1.43M
    auto cmp_cids = std::min(lhs.field_count(), rhs.field_count());
74
3.12M
    for (uint32_t cid = 0; cid < cmp_cids; ++cid) {
75
3.08M
        const auto& lf = lhs.field(cid);
76
3.08M
        const auto& rf = rhs.field(cid);
77
        // Handle nulls: null < non-null
78
3.08M
        if (lf.is_null() != rf.is_null()) {
79
153k
            return lf.is_null() ? -1 : 1;
80
153k
        }
81
2.93M
        if (lf.is_null()) {
82
3.77k
            continue; // both null
83
3.77k
        }
84
2.93M
        auto cmp = lf <=> rf;
85
2.93M
        if (cmp < 0) return -1;
86
1.68M
        if (cmp > 0) return 1;
87
1.68M
    }
88
35.3k
    return 0;
89
1.43M
}
90
91
class TabletReader {
92
    struct KeysParam {
93
        std::vector<RowCursor> start_keys;
94
        std::vector<RowCursor> end_keys;
95
        bool start_key_include = false;
96
        bool end_key_include = false;
97
    };
98
99
public:
100
    // Params for Reader,
101
    // mainly include tablet, data version and fetch range.
102
    struct ReaderParams {
103
1.01M
        bool has_single_version() const {
104
1.01M
            return (rs_splits.size() == 1 &&
105
1.01M
                    rs_splits[0].rs_reader->rowset()->start_version() == 0 &&
106
1.01M
                    !rs_splits[0].rs_reader->rowset()->rowset_meta()->is_segments_overlapping()) ||
107
1.01M
                   (rs_splits.size() == 2 &&
108
1.01M
                    rs_splits[0].rs_reader->rowset()->rowset_meta()->num_rows() == 0 &&
109
1.01M
                    rs_splits[1].rs_reader->rowset()->start_version() == 2 &&
110
1.01M
                    !rs_splits[1].rs_reader->rowset()->rowset_meta()->is_segments_overlapping());
111
1.01M
        }
112
113
14.1k
        int get_be_exec_version() const {
114
14.1k
            if (runtime_state) {
115
10.5k
                return runtime_state->be_exec_version();
116
10.5k
            }
117
3.63k
            return BeExecVersionManager::get_newest_version();
118
14.1k
        }
119
120
1.04M
        void set_read_source(TabletReadSource read_source, bool skip_delete_bitmap = false) {
121
1.04M
            rs_splits = std::move(read_source.rs_splits);
122
1.04M
            delete_predicates = std::move(read_source.delete_predicates);
123
1.04M
#ifndef BE_TEST
124
1.04M
            if (tablet->enable_unique_key_merge_on_write() && !skip_delete_bitmap) {
125
636k
                delete_bitmap = std::move(read_source.delete_bitmap);
126
636k
            }
127
1.04M
#endif
128
1.04M
        }
129
130
        BaseTabletSPtr tablet;
131
        TabletSchemaSPtr tablet_schema;
132
        ReaderType reader_type = ReaderType::READER_QUERY;
133
        bool direct_mode = false;
134
        bool aggregation = false;
135
        // for compaction, schema_change, check_sum: we don't use page cache
136
        // for query, when the BE config disable_storage_page_cache is false, we use page cache
137
        bool use_page_cache = false;
138
        Version version = Version(-1, 0);
139
140
        std::vector<OlapTuple> start_key;
141
        std::vector<OlapTuple> end_key;
142
        bool start_key_include = false;
143
        bool end_key_include = false;
144
145
        std::vector<std::shared_ptr<ColumnPredicate>> predicates;
146
        std::vector<FunctionFilter> function_filters;
147
        std::vector<RowsetMetaSharedPtr> delete_predicates;
148
        // slots that cast may be eliminated in storage layer
149
        std::map<std::string, DataTypePtr> target_cast_type_for_variants;
150
151
        std::map<int32_t, TColumnAccessPaths> all_access_paths;
152
        std::map<int32_t, TColumnAccessPaths> predicate_access_paths;
153
154
        std::vector<RowSetSplits> rs_splits;
155
        // For unique key table with merge-on-write
156
        DeleteBitmapPtr delete_bitmap = nullptr;
157
158
        // return_columns is init from query schema
159
        std::vector<ColumnId> return_columns;
160
        // output_columns only contain columns in OrderByExprs and outputExprs
161
        std::set<int32_t> output_columns;
162
        RuntimeProfile* profile = nullptr;
163
        RuntimeState* runtime_state = nullptr;
164
165
        // use only in vec exec engine
166
        std::vector<ColumnId>* origin_return_columns = nullptr;
167
        std::unordered_set<uint32_t>* tablet_columns_convert_to_null_set = nullptr;
168
        TPushAggOp::type push_down_agg_type_opt = TPushAggOp::NONE;
169
        VExprContextSPtrs common_expr_ctxs_push_down;
170
171
        // used for compaction to record row ids
172
        bool record_rowids = false;
173
        RowIdConversion* rowid_conversion = nullptr;
174
        std::vector<int> topn_filter_source_node_ids;
175
        int topn_filter_target_node_id = -1;
176
        // used for special optimization for query : ORDER BY key LIMIT n
177
        bool read_orderby_key = false;
178
        // used for special optimization for query : ORDER BY key DESC LIMIT n
179
        bool read_orderby_key_reverse = false;
180
        // For rows with the same key, use ascending order (small-to-large) for tie-breakers.
181
        // For example, use lower rowset version / segment id first.
182
        bool use_insert_order_when_same = false;
183
        // num of columns for orderby key
184
        size_t read_orderby_key_num_prefix_columns = 0;
185
        // limit of rows for read_orderby_key
186
        size_t read_orderby_key_limit = 0;
187
        // for vertical compaction
188
        bool is_key_column_group = false;
189
        std::vector<uint32_t> key_group_cluster_key_idxes;
190
191
        // For sparse column compaction optimization
192
        // When true, use optimized path for sparse wide tables
193
        bool enable_sparse_optimization = false;
194
195
        bool is_segcompaction = false;
196
197
        // Enable value predicate pushdown for MOR tables
198
        bool enable_mor_value_predicate_pushdown = false;
199
200
        std::vector<RowwiseIteratorUPtr>* segment_iters_ptr = nullptr;
201
202
        void check_validation() const;
203
204
        int64_t batch_size = -1;
205
206
        std::map<ColumnId, VExprContextSPtr> virtual_column_exprs;
207
        std::map<ColumnId, size_t> vir_cid_to_idx_in_block;
208
        std::map<size_t, DataTypePtr> vir_col_idx_to_type;
209
210
        std::shared_ptr<ScoreRuntime> score_runtime;
211
        CollectionStatisticsPtr collection_statistics;
212
        std::shared_ptr<segment_v2::AnnTopNRuntime> ann_topn_runtime;
213
214
        uint64_t condition_cache_digest = 0;
215
216
        // General LIMIT budget forwarded to SegmentIterator. -1 means no limit.
217
        int64_t general_read_limit = -1;
218
    };
219
220
1.04M
    TabletReader() = default;
221
222
1.03M
    virtual ~TabletReader() = default;
223
224
    TabletReader(const TabletReader&) = delete;
225
    void operator=(const TabletReader&) = delete;
226
227
    // Initialize TabletReader with tablet, data version and fetch range.
228
    virtual Status init(const ReaderParams& read_params);
229
230
    // Read next block with aggregation.
231
    // Return OK and set `*eof` to false when next block is read
232
    // Return OK and set `*eof` to true when no more rows can be read.
233
    // Return others when unexpected error happens.
234
0
    virtual Status next_block_with_aggregation(Block* block, bool* eof) {
235
0
        return Status::Error<ErrorCode::READER_INITIALIZE_ERROR>(
236
0
                "TabletReader not support next_block_with_aggregation");
237
0
    }
238
239
1.44k
    virtual uint64_t merged_rows() const { return _merged_rows; }
240
241
9.97k
    uint64_t filtered_rows() const {
242
9.97k
        return _stats.rows_del_filtered + _stats.rows_del_by_bitmap +
243
9.97k
               _stats.rows_conditions_filtered + _stats.rows_vec_del_cond_filtered +
244
9.97k
               _stats.rows_vec_cond_filtered + _stats.rows_short_circuit_cond_filtered;
245
9.97k
    }
246
247
1.01M
    void set_batch_size(int batch_size) { _reader_context.batch_size = batch_size; }
248
249
0
    int batch_size() const { return _reader_context.batch_size; }
250
251
10.4M
    size_t batch_max_rows() const { return _reader_context.batch_size; }
252
253
1.01M
    void set_preferred_block_size_bytes(size_t bytes) {
254
1.01M
        _reader_context.preferred_block_size_bytes = bytes;
255
1.01M
    }
256
257
    // Returns the preferred output block byte budget. Subclasses that support adaptive batch size
258
    // should override this; the base returns 0 (disabled) so VCollectIterator degrades safely
259
    // when called through a TabletReader* that has not been configured.
260
0
    virtual size_t preferred_block_size_bytes() const { return 0; }
261
262
3.48M
    const OlapReaderStatistics& stats() const { return _stats; }
263
4.23M
    OlapReaderStatistics* mutable_stats() { return &_stats; }
264
265
0
    virtual void update_profile(RuntimeProfile* profile) {}
266
    static Status init_reader_params_and_create_block(
267
            TabletSharedPtr tablet, ReaderType reader_type,
268
            const std::vector<RowsetSharedPtr>& input_rowsets,
269
            TabletReader::ReaderParams* reader_params, Block* block);
270
271
protected:
272
    friend class VCollectIterator;
273
    friend class DeleteHandler;
274
275
    Status _init_params(const ReaderParams& read_params);
276
277
    Status _capture_rs_readers(const ReaderParams& read_params);
278
279
    Status _init_keys_param(const ReaderParams& read_params);
280
281
    Status _init_orderby_keys_param(const ReaderParams& read_params);
282
283
    Status _init_conditions_param(const ReaderParams& read_params);
284
285
    virtual std::shared_ptr<ColumnPredicate> _parse_to_predicate(
286
            const FunctionFilter& function_filter);
287
288
    Status _init_delete_condition(const ReaderParams& read_params);
289
290
    Status _init_return_columns(const ReaderParams& read_params);
291
292
2.06M
    const BaseTabletSPtr& tablet() { return _tablet; }
293
    // If original column is a variant type column, and it's predicate is normalized
294
    // so in order to get the real type of column predicate, we need to reset type
295
    // according to the related type in `target_cast_type_for_variants`.Since variant is not
296
    // an predicate applicable type.Otherwise return the original tablet column.
297
    // Eg. `where cast(v:a as bigint) > 1` will elimate cast, and materialize this variant column
298
    // to type bigint
299
    TabletColumn materialize_column(const TabletColumn& orig);
300
301
4.30M
    const TabletSchema& tablet_schema() { return *_tablet_schema; }
302
303
    Arena _predicate_arena;
304
    std::vector<ColumnId> _return_columns;
305
306
    // used for special optimization for query : ORDER BY key [ASC|DESC] LIMIT n
307
    // columns for orderby keys
308
    std::vector<uint32_t> _orderby_key_columns;
309
    // only use in outer join which change the column nullable which must keep same in
310
    // vec query engine
311
    std::unordered_set<uint32_t>* _tablet_columns_convert_to_null_set = nullptr;
312
313
    BaseTabletSPtr _tablet;
314
    RowsetReaderContext _reader_context;
315
    TabletSchemaSPtr _tablet_schema;
316
    KeysParam _keys_param;
317
    std::vector<bool> _is_lower_keys_included;
318
    std::vector<bool> _is_upper_keys_included;
319
    std::vector<std::shared_ptr<ColumnPredicate>> _col_predicates;
320
    std::vector<std::shared_ptr<ColumnPredicate>> _value_col_predicates;
321
    DeleteHandler _delete_handler;
322
323
    // Indicates whether the tablets has do a aggregation in storage engine.
324
    bool _aggregation = false;
325
    // for agg query, we don't need to finalize when scan agg object data
326
    ReaderType _reader_type = ReaderType::READER_QUERY;
327
    bool _next_delete_flag = false;
328
    bool _delete_sign_available = false;
329
    bool _filter_delete = false;
330
    int32_t _sequence_col_idx = -1;
331
    bool _direct_mode = false;
332
333
    std::vector<uint32_t> _key_cids;
334
    std::vector<uint32_t> _value_cids;
335
336
    uint64_t _merged_rows = 0;
337
    OlapReaderStatistics _stats;
338
};
339
340
} // namespace doris