be/src/exec/operator/scan_operator.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 <cstdint> |
21 | | #include <optional> |
22 | | #include <set> |
23 | | #include <string> |
24 | | |
25 | | #include "common/status.h" |
26 | | #include "common/thread_safety_annotations.h" |
27 | | #include "core/field.h" |
28 | | #include "exec/common/util.hpp" |
29 | | #include "exec/operator/operator.h" |
30 | | #include "exec/pipeline/dependency.h" |
31 | | #include "exec/runtime_filter/runtime_filter_consumer_helper.h" |
32 | | #include "exec/runtime_filter/runtime_filter_partition_pruner.h" |
33 | | #include "exec/scan/scan_node.h" |
34 | | #include "exec/scan/scanner_context.h" |
35 | | #include "exprs/function_filter.h" |
36 | | #include "exprs/vectorized_fn_call.h" |
37 | | #include "exprs/vin_predicate.h" |
38 | | #include "runtime/descriptors.h" |
39 | | #include "storage/predicate/filter_olap_param.h" |
40 | | |
41 | | namespace doris { |
42 | | class ScannerDelegate; |
43 | | class OlapScanner; |
44 | | } // namespace doris |
45 | | |
46 | | namespace doris { |
47 | | |
48 | | enum class PushDownType { |
49 | | // The predicate can not be pushed down to data source |
50 | | UNACCEPTABLE, |
51 | | // The predicate can be pushed down to data source |
52 | | // and the data source can fully evaludate it |
53 | | ACCEPTABLE, |
54 | | // The predicate can be pushed down to data source |
55 | | // but the data source can not fully evaluate it. |
56 | | PARTIAL_ACCEPTABLE |
57 | | }; |
58 | | |
59 | | class ScanLocalStateBase : public PipelineXLocalState<> { |
60 | | public: |
61 | | ScanLocalStateBase(RuntimeState* state, OperatorXBase* parent) |
62 | 132 | : PipelineXLocalState<>(state, parent), _helper(parent->runtime_filter_descs()) {} |
63 | 132 | ~ScanLocalStateBase() override = default; |
64 | | |
65 | | [[nodiscard]] virtual bool should_run_serial() const = 0; |
66 | | |
67 | | virtual RuntimeProfile* scanner_profile() = 0; |
68 | | |
69 | | [[nodiscard]] virtual const TupleDescriptor* input_tuple_desc() const = 0; |
70 | | [[nodiscard]] virtual const TupleDescriptor* output_tuple_desc() const = 0; |
71 | | |
72 | | virtual int64_t limit_per_scanner() = 0; |
73 | | virtual std::atomic<int64_t>* shared_scan_limit_ptr() = 0; |
74 | | |
75 | | virtual void set_scan_ranges(RuntimeState* state, |
76 | | const std::vector<TScanRangeParams>& scan_ranges) = 0; |
77 | | virtual TPushAggOp::type get_push_down_agg_type() = 0; |
78 | | virtual const std::optional<std::vector<int32_t>>& get_push_down_count_slot_ids() const = 0; |
79 | | |
80 | | static bool is_count_star_pushdown(TPushAggOp::type agg_type, |
81 | 4 | const std::optional<std::vector<int32_t>>& count_slot_ids) { |
82 | | // An absent argument field is an old plan with unknown semantics. Only an explicitly empty |
83 | | // argument list proves COUNT(*)/COUNT(1) and permits placeholder slots to be ignored. |
84 | 4 | return agg_type == TPushAggOp::type::COUNT && count_slot_ids.has_value() && |
85 | 4 | count_slot_ids->empty(); |
86 | 4 | } |
87 | | |
88 | 0 | bool is_count_star_pushdown() { |
89 | 0 | return is_count_star_pushdown(get_push_down_agg_type(), get_push_down_count_slot_ids()); |
90 | 0 | } |
91 | | |
92 | | // If scan operator is serial operator(like topn), its real parallelism is 1. |
93 | | // Otherwise, its real parallelism is query_parallel_instance_num. |
94 | | // query_parallel_instance_num of olap table is usually equal to session var parallel_pipeline_task_num. |
95 | | // for file scan operator, its real parallelism will be 1 if it is in batch mode. |
96 | | // Related pr: |
97 | | // https://github.com/apache/doris/pull/42460 |
98 | | // https://github.com/apache/doris/pull/44635 |
99 | | [[nodiscard]] virtual int max_scanners_concurrency(RuntimeState* state) const; |
100 | | [[nodiscard]] virtual int min_scanners_concurrency(RuntimeState* state) const; |
101 | | [[nodiscard]] virtual ScannerScheduler* scan_scheduler(RuntimeState* state) const; |
102 | | |
103 | | // Thread-safe check whether a partition has been pruned by runtime filter. |
104 | | // Callable from any scan type's scanner in scheduling threads. |
105 | | bool is_partition_pruned(int64_t partition_id) const; |
106 | | |
107 | 0 | [[nodiscard]] std::string get_name() { return _parent->get_name(); } |
108 | | |
109 | 17 | uint64_t get_condition_cache_digest() const { return _condition_cache_digest; } |
110 | | |
111 | | Status update_late_arrival_runtime_filter(RuntimeState* state, int& arrived_rf_num); |
112 | | |
113 | | Status clone_conjunct_ctxs(VExprContextSPtrs& scanner_conjuncts); |
114 | | |
115 | | protected: |
116 | | friend class ScannerContext; |
117 | | friend class Scanner; |
118 | | |
119 | | virtual Status _init_profile() = 0; |
120 | | |
121 | | // Hook for subclasses to react after new runtime filters are appended. |
122 | | // Called inside update_late_arrival_runtime_filter() while _conjuncts_lock is held. |
123 | | // Default implementation runs partition pruning on the newly appended RFs. |
124 | | virtual Status _on_runtime_filter_update(); |
125 | | |
126 | | Status _do_partition_pruning_by_rf(); |
127 | | |
128 | | std::atomic<bool> _opened {false}; |
129 | | |
130 | | DependencySPtr _scan_dependency = nullptr; |
131 | | |
132 | | std::shared_ptr<RuntimeProfile> _scanner_profile; |
133 | | RuntimeProfile::Counter* _scanner_wait_worker_timer = nullptr; |
134 | | // Num of newly created free blocks when running query |
135 | | RuntimeProfile::Counter* _newly_create_free_blocks_num = nullptr; |
136 | | // Max num of scanner thread |
137 | | RuntimeProfile::Counter* _max_scan_concurrency = nullptr; |
138 | | RuntimeProfile::Counter* _min_scan_concurrency = nullptr; |
139 | | RuntimeProfile::HighWaterMarkCounter* _peak_running_scanner = nullptr; |
140 | | // time of get block from scanner |
141 | | RuntimeProfile::Counter* _scan_timer = nullptr; |
142 | | RuntimeProfile::Counter* _scan_cpu_timer = nullptr; |
143 | | // time of filter output block from scanner |
144 | | RuntimeProfile::Counter* _filter_timer = nullptr; |
145 | | // rows read from the scanner (including those discarded by (pre)filters) |
146 | | RuntimeProfile::Counter* _rows_read_counter = nullptr; |
147 | | |
148 | | RuntimeProfile::Counter* _num_scanners = nullptr; |
149 | | |
150 | | RuntimeProfile::Counter* _wait_for_rf_timer = nullptr; |
151 | | |
152 | | RuntimeProfile::Counter* _scan_rows = nullptr; |
153 | | RuntimeProfile::Counter* _scan_bytes = nullptr; |
154 | | |
155 | | AnnotatedMutex _conjuncts_lock; |
156 | | RuntimeFilterConsumerHelper _helper; |
157 | | // magic number as seed to generate hash value for condition cache |
158 | | uint64_t _condition_cache_digest = 0; |
159 | | // condition cache filter stats |
160 | | RuntimeProfile::Counter* _condition_cache_hit_counter = nullptr; |
161 | | RuntimeProfile::Counter* _condition_cache_filtered_rows_counter = nullptr; |
162 | | |
163 | | // ---- Runtime-filter partition pruning (scan-agnostic) ---- |
164 | | RuntimeFilterPartitionPruner _rf_partition_pruner; |
165 | | RuntimeProfile::Counter* _partitions_pruned_by_rf_counter = nullptr; |
166 | | RuntimeProfile::Counter* _total_partitions_rf_counter = nullptr; |
167 | | |
168 | | // Moved from ScanLocalState<Derived> to avoid re-instantiation for each Derived type. |
169 | | std::atomic<bool> _eos = false; |
170 | | int _max_pushdown_conditions_per_column = 1024; |
171 | | // Save all function predicates which may be pushed down to data source. |
172 | | std::vector<FunctionFilter> _push_down_functions; |
173 | | |
174 | | // Virtual methods with default implementations; overridden by subclasses when supported. |
175 | | // Declared here so that the normalize methods below (non-Derived-template) can call them. |
176 | 0 | virtual bool _push_down_topn(const RuntimePredicate& predicate) { return false; } |
177 | 0 | virtual PushDownType _should_push_down_bloom_filter() const { |
178 | 0 | return PushDownType::UNACCEPTABLE; |
179 | 0 | } |
180 | 0 | virtual PushDownType _should_push_down_topn_filter() const { |
181 | 0 | return PushDownType::UNACCEPTABLE; |
182 | 0 | } |
183 | 0 | virtual PushDownType _should_push_down_is_null_predicate(VectorizedFnCall* fn_call) const { |
184 | 0 | return PushDownType::UNACCEPTABLE; |
185 | 0 | } |
186 | 0 | virtual PushDownType _should_push_down_in_predicate() const { |
187 | 0 | return PushDownType::UNACCEPTABLE; |
188 | 0 | } |
189 | | virtual PushDownType _should_push_down_binary_predicate( |
190 | | VectorizedFnCall* fn_call, VExprContext* expr_ctx, Field& constant_val, |
191 | 0 | const std::set<std::string> fn_name) const { |
192 | 0 | return PushDownType::UNACCEPTABLE; |
193 | 0 | } |
194 | | virtual Status _should_push_down_function_filter(VectorizedFnCall* fn_call, |
195 | | VExprContext* expr_ctx, |
196 | | StringRef* constant_str, |
197 | | doris::FunctionContext** fn_ctx, |
198 | 0 | PushDownType& pdt) { |
199 | 0 | pdt = PushDownType::UNACCEPTABLE; |
200 | 0 | return Status::OK(); |
201 | 0 | } |
202 | | |
203 | | // Non-templated normalize methods, moved here to avoid re-compilation per Derived type. |
204 | | Status _eval_const_conjuncts(VExprContext* expr_ctx, PushDownType* pdt); |
205 | | Status _normalize_bloom_filter(VExprContext* expr_ctx, const VExprSPtr& root, |
206 | | SlotDescriptor* slot, |
207 | | std::vector<std::shared_ptr<ColumnPredicate>>& predicates, |
208 | | PushDownType* pdt); |
209 | | Status _normalize_topn_filter(VExprContext* expr_ctx, const VExprSPtr& root, |
210 | | SlotDescriptor* slot, |
211 | | std::vector<std::shared_ptr<ColumnPredicate>>& predicates, |
212 | | PushDownType* pdt); |
213 | | Status _normalize_function_filters(VExprContext* expr_ctx, SlotDescriptor* slot, |
214 | | PushDownType* pdt); |
215 | | |
216 | | // Inner PrimitiveType-template methods. Moved to base to avoid N(Derived)×M(PrimitiveType) |
217 | | // instantiation blowup: now instantiated M times total instead of N×M times. |
218 | | template <PrimitiveType T> |
219 | | Status _normalize_in_predicate(VExprContext* expr_ctx, const VExprSPtr& root, |
220 | | SlotDescriptor* slot, |
221 | | std::vector<std::shared_ptr<ColumnPredicate>>& predicates, |
222 | | ColumnValueRange<T>& range, PushDownType* pdt); |
223 | | template <PrimitiveType T> |
224 | | Status _normalize_binary_predicate(VExprContext* expr_ctx, const VExprSPtr& root, |
225 | | SlotDescriptor* slot, |
226 | | std::vector<std::shared_ptr<ColumnPredicate>>& predicates, |
227 | | ColumnValueRange<T>& range, PushDownType* pdt); |
228 | | template <PrimitiveType T> |
229 | | Status _normalize_is_null_predicate(VExprContext* expr_ctx, const VExprSPtr& root, |
230 | | SlotDescriptor* slot, |
231 | | std::vector<std::shared_ptr<ColumnPredicate>>& predicates, |
232 | | ColumnValueRange<T>& range, PushDownType* pdt); |
233 | | template <PrimitiveType PrimitiveType, typename ChangeFixedValueRangeFunc> |
234 | | Status _change_value_range(bool is_equal_op, ColumnValueRange<PrimitiveType>& range, |
235 | | const Field& value, const ChangeFixedValueRangeFunc& func, |
236 | | const std::string& fn_name); |
237 | | }; |
238 | | |
239 | | template <typename LocalStateType> |
240 | | class ScanOperatorX; |
241 | | template <typename Derived> |
242 | | class ScanLocalState : public ScanLocalStateBase { |
243 | | ENABLE_FACTORY_CREATOR(ScanLocalState); |
244 | | ScanLocalState(RuntimeState* state, OperatorXBase* parent) |
245 | 132 | : ScanLocalStateBase(state, parent) {}_ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEEC2EPNS_12RuntimeStateEPNS_13OperatorXBaseE Line | Count | Source | 245 | 17 | : ScanLocalStateBase(state, parent) {} |
_ZN5doris14ScanLocalStateINS_18MockScanLocalStateEEC2EPNS_12RuntimeStateEPNS_13OperatorXBaseE Line | Count | Source | 245 | 110 | : ScanLocalStateBase(state, parent) {} |
_ZN5doris14ScanLocalStateINS_18FileScanLocalStateEEC2EPNS_12RuntimeStateEPNS_13OperatorXBaseE Line | Count | Source | 245 | 5 | : ScanLocalStateBase(state, parent) {} |
Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEEC2EPNS_12RuntimeStateEPNS_13OperatorXBaseE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEEC2EPNS_12RuntimeStateEPNS_13OperatorXBaseE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEEC2EPNS_12RuntimeStateEPNS_13OperatorXBaseE |
246 | 132 | ~ScanLocalState() override = default; _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEED2Ev Line | Count | Source | 246 | 17 | ~ScanLocalState() override = default; |
_ZN5doris14ScanLocalStateINS_18MockScanLocalStateEED2Ev Line | Count | Source | 246 | 110 | ~ScanLocalState() override = default; |
_ZN5doris14ScanLocalStateINS_18FileScanLocalStateEED2Ev Line | Count | Source | 246 | 5 | ~ScanLocalState() override = default; |
Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEED2Ev Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEED2Ev Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEED2Ev |
247 | | |
248 | | Status init(RuntimeState* state, LocalStateInfo& info) override; |
249 | | |
250 | | Status open(RuntimeState* state) override; |
251 | | |
252 | | Status close(RuntimeState* state) override; |
253 | | std::string debug_string(int indentation_level) const final; |
254 | | |
255 | | [[nodiscard]] bool should_run_serial() const override; |
256 | | |
257 | 37 | RuntimeProfile* scanner_profile() override { return _scanner_profile.get(); }_ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE15scanner_profileEv Line | Count | Source | 257 | 37 | RuntimeProfile* scanner_profile() override { return _scanner_profile.get(); } |
Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE15scanner_profileEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE15scanner_profileEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE15scanner_profileEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE15scanner_profileEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE15scanner_profileEv |
258 | | |
259 | | [[nodiscard]] const TupleDescriptor* input_tuple_desc() const override; |
260 | | [[nodiscard]] const TupleDescriptor* output_tuple_desc() const override; |
261 | | |
262 | | int64_t limit_per_scanner() override; |
263 | | std::atomic<int64_t>* shared_scan_limit_ptr() override; |
264 | | |
265 | | void set_scan_ranges(RuntimeState* state, |
266 | 1 | const std::vector<TScanRangeParams>& scan_ranges) override {}Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE15set_scan_rangesEPNS_12RuntimeStateERKSt6vectorINS_16TScanRangeParamsESaIS6_EE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE15set_scan_rangesEPNS_12RuntimeStateERKSt6vectorINS_16TScanRangeParamsESaIS6_EE _ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE15set_scan_rangesEPNS_12RuntimeStateERKSt6vectorINS_16TScanRangeParamsESaIS6_EE Line | Count | Source | 266 | 1 | const std::vector<TScanRangeParams>& scan_ranges) override {} |
Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE15set_scan_rangesEPNS_12RuntimeStateERKSt6vectorINS_16TScanRangeParamsESaIS6_EE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE15set_scan_rangesEPNS_12RuntimeStateERKSt6vectorINS_16TScanRangeParamsESaIS6_EE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE15set_scan_rangesEPNS_12RuntimeStateERKSt6vectorINS_16TScanRangeParamsESaIS6_EE |
267 | | |
268 | | TPushAggOp::type get_push_down_agg_type() override; |
269 | | const std::optional<std::vector<int32_t>>& get_push_down_count_slot_ids() const override; |
270 | | |
271 | 0 | std::vector<Dependency*> execution_dependencies() override { |
272 | 0 | if (_filter_dependencies.empty()) { |
273 | 0 | return {}; |
274 | 0 | } |
275 | 0 | std::vector<Dependency*> res(_filter_dependencies.size()); |
276 | 0 | std::transform(_filter_dependencies.begin(), _filter_dependencies.end(), res.begin(), |
277 | 0 | [](DependencySPtr dep) { return dep.get(); });Unexecuted instantiation: _ZZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE22execution_dependenciesEvENKUlSt10shared_ptrINS_10DependencyEEE_clES5_ Unexecuted instantiation: _ZZN5doris14ScanLocalStateINS_18MockScanLocalStateEE22execution_dependenciesEvENKUlSt10shared_ptrINS_10DependencyEEE_clES5_ Unexecuted instantiation: _ZZN5doris14ScanLocalStateINS_18FileScanLocalStateEE22execution_dependenciesEvENKUlSt10shared_ptrINS_10DependencyEEE_clES5_ Unexecuted instantiation: _ZZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE22execution_dependenciesEvENKUlSt10shared_ptrINS_10DependencyEEE_clES5_ Unexecuted instantiation: _ZZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE22execution_dependenciesEvENKUlSt10shared_ptrINS_10DependencyEEE_clES5_ Unexecuted instantiation: _ZZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE22execution_dependenciesEvENKUlSt10shared_ptrINS_10DependencyEEE_clES5_ |
278 | 0 | return res; |
279 | 0 | } Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE22execution_dependenciesEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE22execution_dependenciesEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE22execution_dependenciesEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE22execution_dependenciesEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE22execution_dependenciesEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE22execution_dependenciesEv |
280 | | |
281 | 0 | std::vector<Dependency*> dependencies() const override { return {_scan_dependency.get()}; }Unexecuted instantiation: _ZNK5doris14ScanLocalStateINS_18FileScanLocalStateEE12dependenciesEv Unexecuted instantiation: _ZNK5doris14ScanLocalStateINS_18OlapScanLocalStateEE12dependenciesEv Unexecuted instantiation: _ZNK5doris14ScanLocalStateINS_18MockScanLocalStateEE12dependenciesEv Unexecuted instantiation: _ZNK5doris14ScanLocalStateINS_21GroupCommitLocalStateEE12dependenciesEv Unexecuted instantiation: _ZNK5doris14ScanLocalStateINS_18JDBCScanLocalStateEE12dependenciesEv Unexecuted instantiation: _ZNK5doris14ScanLocalStateINS_18MetaScanLocalStateEE12dependenciesEv |
282 | | |
283 | 0 | std::vector<int> get_topn_filter_source_node_ids(RuntimeState* state, bool push_down) { |
284 | 0 | std::vector<int> result; |
285 | 0 | for (int id : _parent->cast<typename Derived::Parent>()._topn_filter_source_node_ids) { |
286 | 0 | const auto& pred = state->get_query_ctx()->get_runtime_predicate(id); |
287 | 0 | if (!pred.enable()) { |
288 | 0 | continue; |
289 | 0 | } |
290 | 0 | if (_push_down_topn(pred) == push_down) { |
291 | 0 | result.push_back(id); |
292 | 0 | } |
293 | 0 | } |
294 | 0 | return result; |
295 | 0 | } Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE31get_topn_filter_source_node_idsEPNS_12RuntimeStateEb Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE31get_topn_filter_source_node_idsEPNS_12RuntimeStateEb Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE31get_topn_filter_source_node_idsEPNS_12RuntimeStateEb Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE31get_topn_filter_source_node_idsEPNS_12RuntimeStateEb Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE31get_topn_filter_source_node_idsEPNS_12RuntimeStateEb Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE31get_topn_filter_source_node_idsEPNS_12RuntimeStateEb |
296 | | |
297 | | protected: |
298 | | template <typename LocalStateType> |
299 | | friend class ScanOperatorX; |
300 | | friend class ScannerContext; |
301 | | friend class Scanner; |
302 | | |
303 | | Status _init_profile() override; |
304 | 0 | virtual Status _process_conjuncts(RuntimeState* state) { |
305 | 0 | RETURN_IF_ERROR(_do_partition_pruning_by_rf()); |
306 | 0 | return _normalize_conjuncts(state); |
307 | 0 | } Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE18_process_conjunctsEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE18_process_conjunctsEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE18_process_conjunctsEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE18_process_conjunctsEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE18_process_conjunctsEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE18_process_conjunctsEPNS_12RuntimeStateE |
308 | 0 | virtual bool _should_push_down_common_expr(const VExprSPtr&) { return false; }Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE29_should_push_down_common_exprERKSt10shared_ptrINS_5VExprEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE29_should_push_down_common_exprERKSt10shared_ptrINS_5VExprEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE29_should_push_down_common_exprERKSt10shared_ptrINS_5VExprEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE29_should_push_down_common_exprERKSt10shared_ptrINS_5VExprEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE29_should_push_down_common_exprERKSt10shared_ptrINS_5VExprEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE29_should_push_down_common_exprERKSt10shared_ptrINS_5VExprEE |
309 | | |
310 | 103 | virtual bool can_push_down_column_predicate(const SlotDescriptor* slot) { |
311 | 103 | return _parent->cast<typename Derived::Parent>().can_push_down_column_predicate(slot); |
312 | 103 | } Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE Line | Count | Source | 310 | 4 | virtual bool can_push_down_column_predicate(const SlotDescriptor* slot) { | 311 | 4 | return _parent->cast<typename Derived::Parent>().can_push_down_column_predicate(slot); | 312 | 4 | } |
_ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE Line | Count | Source | 310 | 99 | virtual bool can_push_down_column_predicate(const SlotDescriptor* slot) { | 311 | 99 | return _parent->cast<typename Derived::Parent>().can_push_down_column_predicate(slot); | 312 | 99 | } |
Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE |
313 | | |
314 | 0 | virtual bool _storage_no_merge() { return false; }Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE17_storage_no_mergeEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE17_storage_no_mergeEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE17_storage_no_mergeEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE17_storage_no_mergeEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE17_storage_no_mergeEv Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE17_storage_no_mergeEv |
315 | 0 | virtual bool _is_key_column(const std::string& col_name) { return false; }Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE14_is_key_columnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE14_is_key_columnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE14_is_key_columnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE14_is_key_columnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE14_is_key_columnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE14_is_key_columnERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE |
316 | | |
317 | | // Create a list of scanners. |
318 | | // The number of scanners is related to the implementation of the data source, |
319 | | // predicate conditions, and scheduling strategy. |
320 | | // So this method needs to be implemented separately by the subclass of ScanNode. |
321 | | // Finally, a set of scanners that have been prepared are returned. |
322 | 0 | virtual Status _init_scanners(std::list<ScannerSPtr>* scanners) { return Status::OK(); }Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18FileScanLocalStateEE14_init_scannersEPNSt7__cxx114listISt10shared_ptrINS_7ScannerEESaIS7_EEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18OlapScanLocalStateEE14_init_scannersEPNSt7__cxx114listISt10shared_ptrINS_7ScannerEESaIS7_EEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MockScanLocalStateEE14_init_scannersEPNSt7__cxx114listISt10shared_ptrINS_7ScannerEESaIS7_EEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_21GroupCommitLocalStateEE14_init_scannersEPNSt7__cxx114listISt10shared_ptrINS_7ScannerEESaIS7_EEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18JDBCScanLocalStateEE14_init_scannersEPNSt7__cxx114listISt10shared_ptrINS_7ScannerEESaIS7_EEE Unexecuted instantiation: _ZN5doris14ScanLocalStateINS_18MetaScanLocalStateEE14_init_scannersEPNSt7__cxx114listISt10shared_ptrINS_7ScannerEESaIS7_EEE |
323 | | |
324 | | Status _normalize_conjuncts(RuntimeState* state); |
325 | | // Normalize a conjunct and try to convert it to column predicate recursively. |
326 | | Status _normalize_predicate(VExprContext* context, const VExprSPtr& root, |
327 | | VExprSPtr& output_expr); |
328 | | bool _is_predicate_acting_on_slot(const VExprSPtrs& children, SlotDescriptor** slot_desc, |
329 | | ColumnValueRangeType** range); |
330 | | Status _prepare_scanners(); |
331 | | |
332 | | // Submit the scanner to the thread pool and start execution |
333 | | Status _start_scanners(const std::list<std::shared_ptr<ScannerDelegate>>& scanners); |
334 | | |
335 | | // For some conjunct there is chance to elimate cast operator |
336 | | // Eg. Variant's sub column could eliminate cast in storage layer if |
337 | | // cast dst column type equals storage column type |
338 | | void get_cast_types_for_variants(); |
339 | | void _filter_and_collect_cast_type_for_variant( |
340 | | const VExpr* expr, |
341 | | std::unordered_map<std::string, std::vector<DataTypePtr>>& colname_to_cast_types); |
342 | | |
343 | | Status _get_topn_filters(RuntimeState* state); |
344 | | |
345 | | // Stores conjuncts that have been fully pushed down to the storage layer as predicate columns. |
346 | | // These expr contexts are kept alive to prevent their FunctionContext and constant strings |
347 | | // from being freed prematurely. |
348 | | VExprContextSPtrs _stale_expr_ctxs; |
349 | | VExprContextSPtrs _common_expr_ctxs_push_down; |
350 | | |
351 | | atomic_shared_ptr<ScannerContext> _scanner_ctx; |
352 | | |
353 | | // colname -> cast dst type |
354 | | std::map<std::string, DataTypePtr> _cast_types_for_variants; |
355 | | |
356 | | // slot id -> ColumnValueRange |
357 | | // Parsed from conjuncts |
358 | | phmap::flat_hash_map<int, ColumnValueRangeType> _slot_id_to_value_range; |
359 | | phmap::flat_hash_map<int, std::vector<std::shared_ptr<ColumnPredicate>>> _slot_id_to_predicates; |
360 | | std::vector<std::shared_ptr<MutilColumnBlockPredicate>> _or_predicates; |
361 | | |
362 | | std::vector<std::shared_ptr<Dependency>> _filter_dependencies; |
363 | | |
364 | | // ScanLocalState owns the ownership of scanner, scanner context only has its weakptr |
365 | | std::list<std::shared_ptr<ScannerDelegate>> _scanners; |
366 | | Arena _arena; |
367 | | int _instance_idx = 0; |
368 | | }; |
369 | | |
370 | | template <typename LocalStateType> |
371 | | class ScanOperatorX : public OperatorX<LocalStateType> { |
372 | | public: |
373 | | Status init(const TPlanNode& tnode, RuntimeState* state) override; |
374 | | Status prepare(RuntimeState* state) override; |
375 | | Status get_block_impl(RuntimeState* state, Block* block, bool* eos) override; |
376 | 1 | Status get_block_after_projects(RuntimeState* state, Block* block, bool* eos) override { |
377 | 1 | Status status = OperatorX<LocalStateType>::get_block(state, block, eos); |
378 | 1 | if (status.ok()) { |
379 | 1 | state->get_local_state(operator_id())->update_output_block_counters(*block); |
380 | 1 | } |
381 | 1 | return status; |
382 | 1 | } Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18FileScanLocalStateEE24get_block_after_projectsEPNS_12RuntimeStateEPNS_5BlockEPb Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18OlapScanLocalStateEE24get_block_after_projectsEPNS_12RuntimeStateEPNS_5BlockEPb _ZN5doris13ScanOperatorXINS_18MockScanLocalStateEE24get_block_after_projectsEPNS_12RuntimeStateEPNS_5BlockEPb Line | Count | Source | 376 | 1 | Status get_block_after_projects(RuntimeState* state, Block* block, bool* eos) override { | 377 | 1 | Status status = OperatorX<LocalStateType>::get_block(state, block, eos); | 378 | 1 | if (status.ok()) { | 379 | 1 | state->get_local_state(operator_id())->update_output_block_counters(*block); | 380 | 1 | } | 381 | 1 | return status; | 382 | 1 | } |
Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18JDBCScanLocalStateEE24get_block_after_projectsEPNS_12RuntimeStateEPNS_5BlockEPb Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18MetaScanLocalStateEE24get_block_after_projectsEPNS_12RuntimeStateEPNS_5BlockEPb Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_21GroupCommitLocalStateEE24get_block_after_projectsEPNS_12RuntimeStateEPNS_5BlockEPb |
383 | 0 | [[nodiscard]] bool is_source() const override { return true; }Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18FileScanLocalStateEE9is_sourceEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18OlapScanLocalStateEE9is_sourceEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18MockScanLocalStateEE9is_sourceEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18JDBCScanLocalStateEE9is_sourceEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18MetaScanLocalStateEE9is_sourceEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_21GroupCommitLocalStateEE9is_sourceEv |
384 | | |
385 | | [[nodiscard]] size_t get_reserve_mem_size(RuntimeState* state) override; |
386 | | |
387 | 132 | const std::vector<TRuntimeFilterDesc>& runtime_filter_descs() override { |
388 | 132 | return _runtime_filter_descs; |
389 | 132 | } _ZN5doris13ScanOperatorXINS_18FileScanLocalStateEE20runtime_filter_descsEv Line | Count | Source | 387 | 5 | const std::vector<TRuntimeFilterDesc>& runtime_filter_descs() override { | 388 | 5 | return _runtime_filter_descs; | 389 | 5 | } |
_ZN5doris13ScanOperatorXINS_18OlapScanLocalStateEE20runtime_filter_descsEv Line | Count | Source | 387 | 17 | const std::vector<TRuntimeFilterDesc>& runtime_filter_descs() override { | 388 | 17 | return _runtime_filter_descs; | 389 | 17 | } |
_ZN5doris13ScanOperatorXINS_18MockScanLocalStateEE20runtime_filter_descsEv Line | Count | Source | 387 | 110 | const std::vector<TRuntimeFilterDesc>& runtime_filter_descs() override { | 388 | 110 | return _runtime_filter_descs; | 389 | 110 | } |
Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18JDBCScanLocalStateEE20runtime_filter_descsEv Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18MetaScanLocalStateEE20runtime_filter_descsEv Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_21GroupCommitLocalStateEE20runtime_filter_descsEv |
390 | | |
391 | | // Expose this operator's per-fragment shared partition-boundary parse |
392 | | // result to the non-templated ScanLocalStateBase so it can drive runtime |
393 | | // filter partition pruning without down-casting to a specific scan type. |
394 | | // Subclasses are expected to populate `_parsed_partition_boundaries` from |
395 | | // their own partition-boundary thrift field inside their `prepare()` |
396 | | // override before any LocalState observes the result. |
397 | 0 | const ParsedPartitionBoundaries* parsed_partition_boundaries() const override { |
398 | 0 | return &_parsed_partition_boundaries; |
399 | 0 | } Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18FileScanLocalStateEE27parsed_partition_boundariesEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18OlapScanLocalStateEE27parsed_partition_boundariesEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18MockScanLocalStateEE27parsed_partition_boundariesEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18JDBCScanLocalStateEE27parsed_partition_boundariesEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18MetaScanLocalStateEE27parsed_partition_boundariesEv Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_21GroupCommitLocalStateEE27parsed_partition_boundariesEv |
400 | | |
401 | 0 | [[nodiscard]] virtual int get_column_id(const std::string& col_name) const { return -1; }Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18FileScanLocalStateEE13get_column_idERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18OlapScanLocalStateEE13get_column_idERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18MockScanLocalStateEE13get_column_idERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18JDBCScanLocalStateEE13get_column_idERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18MetaScanLocalStateEE13get_column_idERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_21GroupCommitLocalStateEE13get_column_idERKNSt7__cxx1112basic_stringIcSt11char_traitsIcESaIcEEE |
402 | | |
403 | 103 | [[nodiscard]] virtual bool can_push_down_column_predicate(const SlotDescriptor*) const { |
404 | 103 | return true; |
405 | 103 | } Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18FileScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE _ZNK5doris13ScanOperatorXINS_18OlapScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE Line | Count | Source | 403 | 4 | [[nodiscard]] virtual bool can_push_down_column_predicate(const SlotDescriptor*) const { | 404 | 4 | return true; | 405 | 4 | } |
_ZNK5doris13ScanOperatorXINS_18MockScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE Line | Count | Source | 403 | 99 | [[nodiscard]] virtual bool can_push_down_column_predicate(const SlotDescriptor*) const { | 404 | 99 | return true; | 405 | 99 | } |
Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_21GroupCommitLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18JDBCScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18MetaScanLocalStateEE30can_push_down_column_predicateEPKNS_14SlotDescriptorE |
406 | | |
407 | 0 | TPushAggOp::type get_push_down_agg_type() { return _push_down_agg_type; }Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18OlapScanLocalStateEE22get_push_down_agg_typeEv Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18JDBCScanLocalStateEE22get_push_down_agg_typeEv Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18FileScanLocalStateEE22get_push_down_agg_typeEv Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18MetaScanLocalStateEE22get_push_down_agg_typeEv Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_21GroupCommitLocalStateEE22get_push_down_agg_typeEv Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18MockScanLocalStateEE22get_push_down_agg_typeEv |
408 | | |
409 | 0 | DataDistribution required_data_distribution(RuntimeState* /*state*/) const override { |
410 | 0 | if (OperatorX<LocalStateType>::is_serial_operator()) { |
411 | | // `is_serial_operator()` returns true means we ignore the distribution. |
412 | 0 | return {TLocalPartitionType::NOOP}; |
413 | 0 | } |
414 | 0 | return {TLocalPartitionType::BUCKET_HASH_SHUFFLE}; |
415 | 0 | } Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18FileScanLocalStateEE26required_data_distributionEPNS_12RuntimeStateE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18OlapScanLocalStateEE26required_data_distributionEPNS_12RuntimeStateE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18MockScanLocalStateEE26required_data_distributionEPNS_12RuntimeStateE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18JDBCScanLocalStateEE26required_data_distributionEPNS_12RuntimeStateE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_18MetaScanLocalStateEE26required_data_distributionEPNS_12RuntimeStateE Unexecuted instantiation: _ZNK5doris13ScanOperatorXINS_21GroupCommitLocalStateEE26required_data_distributionEPNS_12RuntimeStateE |
416 | | |
417 | 0 | void set_low_memory_mode(RuntimeState* state) override { |
418 | 0 | auto& local_state = get_local_state(state); |
419 | |
|
420 | 0 | if (auto ctx = local_state._scanner_ctx.load()) { |
421 | 0 | ctx->clear_free_blocks(); |
422 | 0 | } |
423 | 0 | } Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18FileScanLocalStateEE19set_low_memory_modeEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18OlapScanLocalStateEE19set_low_memory_modeEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18MockScanLocalStateEE19set_low_memory_modeEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18JDBCScanLocalStateEE19set_low_memory_modeEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18MetaScanLocalStateEE19set_low_memory_modeEPNS_12RuntimeStateE Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_21GroupCommitLocalStateEE19set_low_memory_modeEPNS_12RuntimeStateE |
424 | | |
425 | | using OperatorX<LocalStateType>::node_id; |
426 | | using OperatorX<LocalStateType>::operator_id; |
427 | | using OperatorX<LocalStateType>::get_local_state; |
428 | | |
429 | | #ifdef BE_TEST |
430 | 26 | ScanOperatorX() = default; |
431 | | #endif |
432 | | |
433 | | protected: |
434 | | using LocalState = LocalStateType; |
435 | | friend class OlapScanner; |
436 | | ScanOperatorX(ObjectPool* pool, const TPlanNode& tnode, int operator_id, |
437 | | const DescriptorTbl& descs, int parallel_tasks = 0); |
438 | 50 | virtual ~ScanOperatorX() = default; _ZN5doris13ScanOperatorXINS_18OlapScanLocalStateEED2Ev Line | Count | Source | 438 | 19 | virtual ~ScanOperatorX() = default; |
_ZN5doris13ScanOperatorXINS_18MockScanLocalStateEED2Ev Line | Count | Source | 438 | 26 | virtual ~ScanOperatorX() = default; |
_ZN5doris13ScanOperatorXINS_18FileScanLocalStateEED2Ev Line | Count | Source | 438 | 5 | virtual ~ScanOperatorX() = default; |
Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_21GroupCommitLocalStateEED2Ev Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18JDBCScanLocalStateEED2Ev Unexecuted instantiation: _ZN5doris13ScanOperatorXINS_18MetaScanLocalStateEED2Ev |
439 | | template <typename Derived> |
440 | | friend class ScanLocalState; |
441 | | friend class OlapScanLocalState; |
442 | | |
443 | | // For load scan node, there should be both input and output tuple descriptor. |
444 | | // For query scan node, there is only output_tuple_desc. |
445 | | TupleId _input_tuple_id = -1; |
446 | | TupleId _output_tuple_id = -1; |
447 | | const TupleDescriptor* _input_tuple_desc = nullptr; |
448 | | const TupleDescriptor* _output_tuple_desc = nullptr; |
449 | | |
450 | | phmap::flat_hash_map<int, SlotDescriptor*> _slot_id_to_slot_desc; |
451 | | std::unordered_map<std::string, int> _colname_to_slot_id; |
452 | | |
453 | | // These two values are from query_options |
454 | | int _max_scan_key_num = 48; |
455 | | int _max_pushdown_conditions_per_column = 1024; |
456 | | |
457 | | // If the query like select * from table limit 10; then the query should run in |
458 | | // single scanner to avoid too many scanners which will cause lots of useless read. |
459 | | bool _should_run_serial = false; |
460 | | |
461 | | VExprContextSPtrs _common_expr_ctxs_push_down; |
462 | | |
463 | | // If sort info is set, push limit to each scanner; |
464 | | int64_t _limit_per_scanner = -1; |
465 | | |
466 | | // Shared remaining limit across all parallel instances and their scanners. |
467 | | // Initialized to _limit (SQL LIMIT); -1 means no limit. |
468 | | std::atomic<int64_t> _shared_scan_limit {-1}; |
469 | | |
470 | | std::vector<TRuntimeFilterDesc> _runtime_filter_descs; |
471 | | |
472 | | TPushAggOp::type _push_down_agg_type; |
473 | | |
474 | | // Semantic arguments of a pushed-down COUNT. This is deliberately optional because absence |
475 | | // and an empty list have different meanings during a BE-first rolling upgrade: |
476 | | // |
477 | | // - nullopt: an old FE did not send the field, so the new BE must use the normal scan; |
478 | | // - empty: the new FE explicitly planned COUNT(*)/COUNT(1); |
479 | | // - non-empty: the new FE explicitly planned COUNT(col). |
480 | | // |
481 | | // Treating nullopt as empty would silently reinterpret an old plan as COUNT(*). |
482 | | std::optional<std::vector<int32_t>> _push_down_count_slot_ids; |
483 | | |
484 | | // Record the value of the aggregate function 'count' from doris's be |
485 | | int64_t _push_down_count = -1; |
486 | | const int _parallel_tasks = 0; |
487 | | |
488 | | std::vector<int> _topn_filter_source_node_ids; |
489 | | |
490 | | std::shared_ptr<MemShareArbitrator> _mem_arb = nullptr; |
491 | | std::shared_ptr<MemLimiter> _mem_limiter = nullptr; |
492 | | |
493 | | // Shared parse result of partition boundaries for runtime-filter partition |
494 | | // pruning. Lives here (rather than on the Olap-specific subclass) so any |
495 | | // future scan type can populate it in its `prepare()` override and reuse |
496 | | // the generic pruning machinery in ScanLocalStateBase. |
497 | | ParsedPartitionBoundaries _parsed_partition_boundaries; |
498 | | }; |
499 | | |
500 | | } // namespace doris |