be/src/service/point_query_executor.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 "service/point_query_executor.h" |
19 | | |
20 | | #include <fmt/format.h> |
21 | | #include <gen_cpp/Descriptors_types.h> |
22 | | #include <gen_cpp/Exprs_types.h> |
23 | | #include <gen_cpp/PaloInternalService_types.h> |
24 | | #include <gen_cpp/internal_service.pb.h> |
25 | | #include <glog/logging.h> |
26 | | #include <google/protobuf/extension_set.h> |
27 | | #include <stdlib.h> |
28 | | |
29 | | #include <memory> |
30 | | #include <unordered_map> |
31 | | #include <vector> |
32 | | |
33 | | #include "cloud/cloud_tablet.h" |
34 | | #include "cloud/config.h" |
35 | | #include "common/cast_set.h" |
36 | | #include "common/consts.h" |
37 | | #include "common/status.h" |
38 | | #include "core/data_type/data_type_factory.hpp" |
39 | | #include "core/data_type_serde/data_type_serde.h" |
40 | | #include "exec/sink/writer/vmysql_result_writer.h" |
41 | | #include "exprs/vexpr.h" |
42 | | #include "exprs/vexpr_context.h" |
43 | | #include "exprs/vexpr_fwd.h" |
44 | | #include "exprs/vslot_ref.h" |
45 | | #include "io/cache/remote_scan_cache_write_limiter.h" |
46 | | #include "io/io_common.h" |
47 | | #include "runtime/descriptors.h" |
48 | | #include "runtime/exec_env.h" |
49 | | #include "runtime/result_block_buffer.h" |
50 | | #include "runtime/runtime_profile.h" |
51 | | #include "runtime/runtime_state.h" |
52 | | #include "runtime/thread_context.h" |
53 | | #include "storage/read_time_hidden_column.h" |
54 | | #include "storage/row_cursor.h" |
55 | | #include "storage/rowset/beta_rowset.h" |
56 | | #include "storage/rowset/rowset_fwd.h" |
57 | | #include "storage/segment/column_reader.h" |
58 | | #include "storage/tablet/tablet_schema.h" |
59 | | #include "storage/utils.h" |
60 | | #include "util/defer_op.h" |
61 | | #include "util/jsonb/serialize.h" |
62 | | #include "util/lru_cache.h" |
63 | | #include "util/simd/bits.h" |
64 | | #include "util/thrift_util.h" |
65 | | |
66 | | namespace doris { |
67 | | |
68 | 0 | static BetaRowsetSharedPtr get_beta_rowset(RowsetSharedPtr* rowset_ptr) { |
69 | 0 | DORIS_CHECK(rowset_ptr != nullptr); |
70 | 0 | DORIS_CHECK(*rowset_ptr != nullptr); |
71 | 0 | return std::static_pointer_cast<BetaRowset>(*rowset_ptr); |
72 | 0 | } |
73 | | |
74 | | static void replace_point_query_read_time_hidden_columns( |
75 | | const std::vector<std::pair<int32_t, uint32_t>>& hidden_columns, const TabletSchema& schema, |
76 | 0 | const BetaRowset& rowset, MutableColumns& result_columns) { |
77 | 0 | for (const auto& [column_uid, position] : hidden_columns) { |
78 | 0 | replace_suffix_with_read_time_hidden_column( |
79 | 0 | get_read_time_hidden_column_type(schema, column_uid), rowset.version(), |
80 | 0 | rowset.commit_tso(), 1, *result_columns[position]); |
81 | 0 | } |
82 | 0 | } |
83 | | |
84 | | class PointQueryResultBlockBuffer final : public MySQLResultBlockBuffer { |
85 | | public: |
86 | 0 | PointQueryResultBlockBuffer(RuntimeState* state) : MySQLResultBlockBuffer(state) {} |
87 | 0 | ~PointQueryResultBlockBuffer() override = default; |
88 | 0 | std::shared_ptr<TFetchDataResult> get_block() { |
89 | 0 | std::lock_guard<std::mutex> l(_lock); |
90 | 0 | DCHECK_EQ(_result_batch_queue.size(), 1); |
91 | 0 | auto result = std::move(_result_batch_queue.front()); |
92 | 0 | _result_batch_queue.pop_front(); |
93 | 0 | return result; |
94 | 0 | } |
95 | | }; |
96 | | |
97 | 4.09k | Reusable::~Reusable() = default; |
98 | | |
99 | | // get missing and include column ids |
100 | | // input include_cids : the output expr slots columns unique ids |
101 | | // missing_cids : the output expr columns that not in row columns cids |
102 | | static void get_missing_and_include_cids(const TabletSchema& schema, |
103 | | const std::vector<SlotDescriptor*>& slots, |
104 | | int target_rs_column_id, bool has_delete_sign, |
105 | | std::unordered_set<int>& missing_cids, |
106 | 4.04k | std::unordered_set<int>& include_cids) { |
107 | 4.04k | missing_cids.clear(); |
108 | 4.04k | include_cids.clear(); |
109 | 4.04k | for (auto* slot : slots) { |
110 | 4.04k | missing_cids.insert(slot->col_unique_id()); |
111 | 4.04k | } |
112 | 4.04k | if (has_delete_sign) { |
113 | 0 | missing_cids.insert(schema.columns()[schema.delete_sign_idx()]->unique_id()); |
114 | 0 | } |
115 | 4.04k | if (target_rs_column_id == -1) { |
116 | | // no row store columns |
117 | 4.04k | return; |
118 | 4.04k | } |
119 | 1 | const TabletColumn& target_rs_column = schema.column_by_uid(target_rs_column_id); |
120 | 1 | DCHECK(target_rs_column.is_row_store_column()); |
121 | | // The full column group is considered a full match, thus no missing cids |
122 | 1 | if (schema.row_columns_uids().empty()) { |
123 | 0 | missing_cids.clear(); |
124 | 0 | return; |
125 | 0 | } |
126 | 1 | for (int cid : schema.row_columns_uids()) { |
127 | 0 | missing_cids.erase(cid); |
128 | 0 | include_cids.insert(cid); |
129 | 0 | } |
130 | 1 | } |
131 | | |
132 | | constexpr static int s_preallocted_blocks_num = 32; |
133 | | |
134 | | static void extract_slot_ref(const VExprSPtr& expr, TupleDescriptor* tuple_desc, |
135 | 4.04k | std::vector<SlotDescriptor*>& slots) { |
136 | 4.04k | const auto& children = expr->children(); |
137 | 4.04k | for (const auto& i : children) { |
138 | 0 | extract_slot_ref(i, tuple_desc, slots); |
139 | 0 | } |
140 | | |
141 | 4.04k | auto node_type = expr->node_type(); |
142 | 4.04k | if (node_type == TExprNodeType::SLOT_REF) { |
143 | 4.04k | int column_id = static_cast<const VSlotRef*>(expr.get())->column_id(); |
144 | 4.04k | auto* slot_desc = tuple_desc->slots()[column_id]; |
145 | 4.04k | slots.push_back(slot_desc); |
146 | 4.04k | } |
147 | 4.04k | } |
148 | | |
149 | | Status Reusable::init(const TDescriptorTable& t_desc_tbl, const std::vector<TExpr>& output_exprs, |
150 | | const TQueryOptions& query_options, const TabletSchema& schema, |
151 | 4.04k | size_t block_size) { |
152 | 4.04k | _runtime_state = RuntimeState::create_unique(); |
153 | 4.04k | _runtime_state->set_query_options(query_options); |
154 | 4.04k | RETURN_IF_ERROR(DescriptorTbl::create(_runtime_state->obj_pool(), t_desc_tbl, &_desc_tbl)); |
155 | 4.04k | _runtime_state->set_desc_tbl(_desc_tbl); |
156 | 88.9k | for (const auto* slot : tuple_desc()->slots()) { |
157 | 89.0k | if (!slot->all_access_paths().empty() || !slot->predicate_access_paths().empty()) { |
158 | 0 | return Status::InternalError( |
159 | 0 | "Short-circuit point query does not support nested column access paths, " |
160 | 0 | "slot: {}. Please upgrade FE to disable nested column pruning for " |
161 | 0 | "short-circuit point queries.", |
162 | 0 | slot->col_name()); |
163 | 0 | } |
164 | 88.9k | } |
165 | 4.04k | _block_pool.resize(block_size); |
166 | 8.09k | for (auto& i : _block_pool) { |
167 | 8.09k | i = Block::create_unique(tuple_desc()->slots(), 2); |
168 | | // Name is useless but cost space |
169 | 8.09k | i->clear_names(); |
170 | 8.09k | } |
171 | | |
172 | 4.04k | RETURN_IF_ERROR(VExpr::create_expr_trees(output_exprs, _output_exprs_ctxs)); |
173 | 4.04k | RowDescriptor row_desc(tuple_desc()); |
174 | | // Prepare the exprs to run. |
175 | 4.04k | RETURN_IF_ERROR(VExpr::prepare(_output_exprs_ctxs, _runtime_state.get(), row_desc)); |
176 | 4.04k | RETURN_IF_ERROR(VExpr::open(_output_exprs_ctxs, _runtime_state.get())); |
177 | 4.04k | _create_timestamp = butil::gettimeofday_ms(); |
178 | 4.04k | _data_type_serdes = create_data_type_serdes(tuple_desc()->slots()); |
179 | 4.04k | _col_default_values.resize(tuple_desc()->slots().size()); |
180 | 4.04k | bool has_delete_sign = false; |
181 | 93.0k | for (int i = 0; i < tuple_desc()->slots().size(); ++i) { |
182 | 88.9k | auto* slot = tuple_desc()->slots()[i]; |
183 | 88.9k | if (slot->col_name() == DELETE_SIGN) { |
184 | 0 | has_delete_sign = true; |
185 | 0 | } |
186 | 88.9k | _col_uid_to_idx[slot->col_unique_id()] = i; |
187 | 88.9k | _col_default_values[i] = slot->col_default_value(); |
188 | 88.9k | } |
189 | | |
190 | | // Get the output slot descriptors |
191 | 4.04k | std::vector<SlotDescriptor*> output_slot_descs; |
192 | 4.04k | for (const auto& expr : _output_exprs_ctxs) { |
193 | 4.04k | extract_slot_ref(expr->root(), tuple_desc(), output_slot_descs); |
194 | 4.04k | } |
195 | | |
196 | 4.04k | std::unordered_set<int32_t> read_time_hidden_column_uids; |
197 | 4.04k | for (const auto* slot : output_slot_descs) { |
198 | 4.04k | const int32_t column_uid = slot->col_unique_id(); |
199 | 4.04k | if (get_read_time_hidden_column_type(schema, column_uid) != |
200 | 4.04k | ReadTimeHiddenColumnType::NONE && |
201 | 4.04k | read_time_hidden_column_uids.insert(column_uid).second) { |
202 | 0 | _read_time_hidden_columns.emplace_back(column_uid, _col_uid_to_idx.at(column_uid)); |
203 | 0 | } |
204 | 4.04k | } |
205 | | |
206 | | // get the delete sign idx in block |
207 | 4.04k | if (has_delete_sign) { |
208 | 0 | _delete_sign_idx = _col_uid_to_idx[schema.columns()[schema.delete_sign_idx()]->unique_id()]; |
209 | 0 | } |
210 | | |
211 | 4.04k | if (schema.have_column(BeConsts::ROW_STORE_COL)) { |
212 | 0 | const auto& column = *DORIS_TRY(schema.column(BeConsts::ROW_STORE_COL)); |
213 | 0 | _row_store_column_ids = column.unique_id(); |
214 | 0 | } |
215 | 4.04k | get_missing_and_include_cids(schema, output_slot_descs, _row_store_column_ids, has_delete_sign, |
216 | 4.04k | _missing_col_uids, _include_col_uids); |
217 | | |
218 | 4.04k | return Status::OK(); |
219 | 4.04k | } |
220 | | |
221 | 0 | std::unique_ptr<Block> Reusable::get_block() { |
222 | 0 | std::lock_guard lock(_block_mutex); |
223 | 0 | if (_block_pool.empty()) { |
224 | 0 | auto block = Block::create_unique(tuple_desc()->slots(), 2); |
225 | | // Name is useless but cost space |
226 | 0 | block->clear_names(); |
227 | 0 | return block; |
228 | 0 | } |
229 | 0 | auto block = std::move(_block_pool.back()); |
230 | 0 | CHECK(block != nullptr); |
231 | 0 | _block_pool.pop_back(); |
232 | 0 | return block; |
233 | 0 | } |
234 | | |
235 | 0 | void Reusable::return_block(std::unique_ptr<Block>& block) { |
236 | 0 | std::lock_guard lock(_block_mutex); |
237 | 0 | if (block == nullptr) { |
238 | 0 | return; |
239 | 0 | } |
240 | 0 | block->clear_column_data(); |
241 | 0 | _block_pool.push_back(std::move(block)); |
242 | 0 | if (_block_pool.size() > s_preallocted_blocks_num) { |
243 | 0 | _block_pool.resize(s_preallocted_blocks_num); |
244 | 0 | } |
245 | 0 | } |
246 | | |
247 | 0 | LookupConnectionCache* LookupConnectionCache::create_global_instance(size_t capacity) { |
248 | 0 | DCHECK(ExecEnv::GetInstance()->get_lookup_connection_cache() == nullptr); |
249 | 0 | auto* res = new LookupConnectionCache(capacity); |
250 | 0 | return res; |
251 | 0 | } |
252 | | |
253 | | RowCache::RowCache(int64_t capacity, int num_shards) |
254 | 5 | : LRUCachePolicy(CachePolicy::CacheType::POINT_QUERY_ROW_CACHE, capacity, |
255 | 5 | LRUCacheType::SIZE, config::point_query_row_cache_stale_sweep_time_sec, |
256 | 5 | num_shards, /*element count capacity */ 0, |
257 | 5 | /*enable prune*/ true, /*is lru-k*/ true) {} |
258 | | |
259 | | // Create global instance of this class |
260 | 0 | RowCache* RowCache::create_global_cache(int64_t capacity, uint32_t num_shards) { |
261 | 0 | DCHECK(ExecEnv::GetInstance()->get_row_cache() == nullptr); |
262 | 0 | auto* res = new RowCache(capacity, num_shards); |
263 | 0 | return res; |
264 | 0 | } |
265 | | |
266 | 12 | RowCache* RowCache::instance() { |
267 | 12 | return ExecEnv::GetInstance()->get_row_cache(); |
268 | 12 | } |
269 | | |
270 | 10 | bool RowCache::lookup(const RowCacheKey& key, CacheHandle* handle) { |
271 | 10 | const std::string& encoded_key = key.encode(); |
272 | 10 | auto* lru_handle = LRUCachePolicy::lookup(encoded_key); |
273 | 10 | if (!lru_handle) { |
274 | | // cache miss |
275 | 5 | return false; |
276 | 5 | } |
277 | 5 | *handle = CacheHandle(this, lru_handle); |
278 | 5 | return true; |
279 | 10 | } |
280 | | |
281 | 15 | void RowCache::insert(const RowCacheKey& key, const Slice& value) { |
282 | 15 | char* cache_value = static_cast<char*>(malloc(value.size)); |
283 | 15 | memcpy(cache_value, value.data, value.size); |
284 | 15 | auto* row_cache_value = new RowCacheValue; |
285 | 15 | row_cache_value->cache_value = cache_value; |
286 | 15 | const std::string& encoded_key = key.encode(); |
287 | 15 | auto* handle = LRUCachePolicy::insert(encoded_key, row_cache_value, value.size, value.size, |
288 | 15 | CachePriority::NORMAL); |
289 | | // handle will released |
290 | 15 | auto tmp = CacheHandle {this, handle}; |
291 | 15 | } |
292 | | |
293 | 3 | void RowCache::erase(const RowCacheKey& key) { |
294 | 3 | const std::string& encoded_key = key.encode(); |
295 | 3 | LRUCachePolicy::erase(encoded_key); |
296 | 3 | } |
297 | | |
298 | 2.09k | LookupConnectionCache::CacheValue::~CacheValue() { |
299 | 2.09k | SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER( |
300 | 2.09k | ExecEnv::GetInstance()->point_query_executor_mem_tracker()); |
301 | 2.09k | item.reset(); |
302 | 2.09k | } |
303 | | |
304 | 0 | PointQueryExecutor::~PointQueryExecutor() { |
305 | 0 | SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER( |
306 | 0 | ExecEnv::GetInstance()->point_query_executor_mem_tracker()); |
307 | 0 | _tablet.reset(); |
308 | 0 | _reusable.reset(); |
309 | 0 | _result_block.reset(); |
310 | 0 | _row_read_ctxs.clear(); |
311 | 0 | } |
312 | | |
313 | 0 | void PointQueryExecutor::_init_remote_scan_cache_write_limiter() { |
314 | 0 | const auto& query_options = _reusable->runtime_state()->query_options(); |
315 | 0 | const bool initialize_remote_scan_cache_write_limiter = |
316 | 0 | config::is_cloud_mode() && config::enable_file_cache && |
317 | 0 | query_options.__isset.file_cache_query_limit_bytes && |
318 | 0 | query_options.file_cache_query_limit_bytes >= 0 && |
319 | 0 | query_options.query_type == TQueryType::SELECT; |
320 | 0 | if (initialize_remote_scan_cache_write_limiter) { |
321 | 0 | _remote_scan_cache_write_limiter = std::make_unique<io::RemoteScanCacheWriteLimiter>( |
322 | 0 | _reusable->runtime_state()->query_id(), query_options.file_cache_query_limit_bytes); |
323 | 0 | } |
324 | 0 | } |
325 | | |
326 | | Status PointQueryExecutor::init(const PTabletKeyLookupRequest* request, |
327 | 0 | PTabletKeyLookupResponse* response) { |
328 | 0 | SCOPED_TIMER(&_profile_metrics.init_ns); |
329 | 0 | _response = response; |
330 | | // using cache |
331 | 0 | __int128_t uuid = |
332 | 0 | static_cast<__int128_t>(request->uuid().uuid_high()) << 64 | request->uuid().uuid_low(); |
333 | 0 | SCOPED_ATTACH_TASK(ExecEnv::GetInstance()->point_query_executor_mem_tracker()); |
334 | 0 | auto cache_handle = LookupConnectionCache::instance()->get(uuid); |
335 | 0 | _binary_row_format = request->is_binary_row(); |
336 | 0 | _tablet = DORIS_TRY(ExecEnv::get_tablet(request->tablet_id())); |
337 | 0 | if (cache_handle != nullptr) { |
338 | 0 | _reusable = cache_handle; |
339 | 0 | _profile_metrics.hit_lookup_cache = true; |
340 | 0 | } else { |
341 | | // Lightweight request: FE may omit reusable query context and rely on uuid cache. |
342 | | // If cache miss and required parameters are absent, ask FE to resend a full request. |
343 | 0 | if (uuid != 0 && (!request->has_desc_tbl() || request->desc_tbl().empty() || |
344 | 0 | !request->has_output_expr() || request->output_expr().empty() || |
345 | 0 | !request->has_query_options() || request->query_options().empty())) { |
346 | 0 | if (VLOG_DEBUG_IS_ON) { |
347 | 0 | VLOG_DEBUG << "lookup connection cache miss, ask FE to resend query context" |
348 | 0 | << ", tablet_id=" << request->tablet_id() |
349 | 0 | << ", uuid_high=" << request->uuid().uuid_high() |
350 | 0 | << ", uuid_low=" << request->uuid().uuid_low(); |
351 | 0 | } |
352 | 0 | response->set_need_resend_query_context(true); |
353 | 0 | return Status::OK(); |
354 | 0 | } |
355 | 0 | if (uuid == 0 && (!request->has_desc_tbl() || request->desc_tbl().empty() || |
356 | 0 | !request->has_output_expr() || request->output_expr().empty())) { |
357 | 0 | return Status::InvalidArgument( |
358 | 0 | "tablet_fetch_data requires desc_tbl/output_expr when uuid is not set"); |
359 | 0 | } |
360 | | // init handle |
361 | 0 | auto reusable_ptr = std::make_shared<Reusable>(); |
362 | 0 | TDescriptorTable t_desc_tbl; |
363 | 0 | TExprList t_output_exprs; |
364 | 0 | auto len = cast_set<uint32_t>(request->desc_tbl().size()); |
365 | 0 | RETURN_IF_ERROR( |
366 | 0 | deserialize_thrift_msg(reinterpret_cast<const uint8_t*>(request->desc_tbl().data()), |
367 | 0 | &len, false, &t_desc_tbl)); |
368 | 0 | len = cast_set<uint32_t>(request->output_expr().size()); |
369 | 0 | RETURN_IF_ERROR(deserialize_thrift_msg( |
370 | 0 | reinterpret_cast<const uint8_t*>(request->output_expr().data()), &len, false, |
371 | 0 | &t_output_exprs)); |
372 | 0 | _reusable = reusable_ptr; |
373 | 0 | TQueryOptions t_query_options; |
374 | 0 | len = cast_set<uint32_t>(request->query_options().size()); |
375 | 0 | if (request->has_query_options()) { |
376 | 0 | RETURN_IF_ERROR(deserialize_thrift_msg( |
377 | 0 | reinterpret_cast<const uint8_t*>(request->query_options().data()), &len, false, |
378 | 0 | &t_query_options)); |
379 | 0 | } |
380 | 0 | if (uuid != 0) { |
381 | | // could be reused by requests after, pre allocte more blocks |
382 | 0 | RETURN_IF_ERROR(reusable_ptr->init(t_desc_tbl, t_output_exprs.exprs, t_query_options, |
383 | 0 | *_tablet->tablet_schema(), |
384 | 0 | s_preallocted_blocks_num)); |
385 | 0 | LookupConnectionCache::instance()->add(uuid, reusable_ptr); |
386 | 0 | } else { |
387 | 0 | RETURN_IF_ERROR(reusable_ptr->init(t_desc_tbl, t_output_exprs.exprs, t_query_options, |
388 | 0 | *_tablet->tablet_schema(), 1)); |
389 | 0 | } |
390 | 0 | } |
391 | 0 | _init_remote_scan_cache_write_limiter(); |
392 | | // Set timezone from request for functions like from_unixtime() |
393 | 0 | if (request->has_time_zone() && !request->time_zone().empty()) { |
394 | 0 | _reusable->runtime_state()->set_timezone(request->time_zone()); |
395 | 0 | } |
396 | 0 | if (request->has_version() && request->version() >= 0) { |
397 | 0 | _version = request->version(); |
398 | 0 | } |
399 | 0 | RETURN_IF_ERROR(_init_keys(request)); |
400 | 0 | _result_block = _reusable->get_block(); |
401 | 0 | CHECK(_result_block != nullptr); |
402 | |
|
403 | 0 | return Status::OK(); |
404 | 0 | } |
405 | | |
406 | 0 | Status PointQueryExecutor::lookup_up() { |
407 | 0 | SCOPED_ATTACH_TASK(ExecEnv::GetInstance()->point_query_executor_mem_tracker()); |
408 | 0 | RETURN_IF_ERROR(_lookup_row_key()); |
409 | 0 | RETURN_IF_ERROR(_lookup_row_data()); |
410 | 0 | RETURN_IF_ERROR(_output_data()); |
411 | 0 | return Status::OK(); |
412 | 0 | } |
413 | | |
414 | 0 | void PointQueryExecutor::print_profile() { |
415 | 0 | auto init_us = _profile_metrics.init_ns.value() / 1000; |
416 | 0 | auto init_key_us = _profile_metrics.init_key_ns.value() / 1000; |
417 | 0 | auto lookup_key_us = _profile_metrics.lookup_key_ns.value() / 1000; |
418 | 0 | auto lookup_data_us = _profile_metrics.lookup_data_ns.value() / 1000; |
419 | 0 | auto output_data_us = _profile_metrics.output_data_ns.value() / 1000; |
420 | 0 | auto load_segments_key_us = _profile_metrics.load_segment_key_stage_ns.value() / 1000; |
421 | 0 | auto load_segments_data_us = _profile_metrics.load_segment_data_stage_ns.value() / 1000; |
422 | 0 | auto total_us = init_us + lookup_key_us + lookup_data_us + output_data_us; |
423 | 0 | auto read_stats = _profile_metrics.read_stats; |
424 | 0 | const std::string stats_str = fmt::format( |
425 | 0 | "[lookup profile:{}us] init:{}us, init_key:{}us," |
426 | 0 | " lookup_key:{}us, load_segments_key:{}us, lookup_data:{}us, load_segments_data:{}us," |
427 | 0 | " output_data:{}us, " |
428 | 0 | "hit_lookup_cache:{}" |
429 | 0 | ", is_binary_row:{}, output_columns:{}, total_keys:{}, row_cache_hits:{}" |
430 | 0 | ", hit_cached_pages:{}, total_pages_read:{}, compressed_bytes_read:{}, " |
431 | 0 | "io_latency:{}ns, " |
432 | 0 | "uncompressed_bytes_read:{}, result_data_bytes:{}, row_hits:{}" |
433 | 0 | ", rs_column_uid:{}, bytes_read_from_local:{}, bytes_read_from_remote:{}, " |
434 | 0 | "local_io_timer:{}, remote_io_timer:{}, local_write_timer:{}", |
435 | 0 | total_us, init_us, init_key_us, lookup_key_us, load_segments_key_us, lookup_data_us, |
436 | 0 | load_segments_data_us, output_data_us, _profile_metrics.hit_lookup_cache, |
437 | 0 | _binary_row_format, _reusable->output_exprs().size(), _row_read_ctxs.size(), |
438 | 0 | _profile_metrics.row_cache_hits, read_stats.cached_pages_num, |
439 | 0 | read_stats.total_pages_num, read_stats.compressed_bytes_read, read_stats.io_ns, |
440 | 0 | read_stats.uncompressed_bytes_read, _profile_metrics.result_data_bytes, _row_hits, |
441 | 0 | _reusable->rs_column_uid(), |
442 | 0 | _profile_metrics.read_stats.file_cache_stats.bytes_read_from_local, |
443 | 0 | _profile_metrics.read_stats.file_cache_stats.bytes_read_from_remote, |
444 | 0 | _profile_metrics.read_stats.file_cache_stats.local_io_timer, |
445 | 0 | _profile_metrics.read_stats.file_cache_stats.remote_io_timer, |
446 | 0 | _profile_metrics.read_stats.file_cache_stats.write_cache_io_timer); |
447 | |
|
448 | 0 | constexpr static int kSlowThreholdUs = 50 * 1000; // 50ms |
449 | 0 | if (total_us > kSlowThreholdUs) { |
450 | 0 | LOG(WARNING) << "slow query, " << stats_str; |
451 | 0 | } else if (VLOG_DEBUG_IS_ON) { |
452 | 0 | VLOG_DEBUG << stats_str; |
453 | 0 | } else { |
454 | 0 | LOG_EVERY_N(INFO, 1000) << stats_str; |
455 | 0 | } |
456 | 0 | } |
457 | | |
458 | 0 | Status PointQueryExecutor::_init_keys(const PTabletKeyLookupRequest* request) { |
459 | 0 | SCOPED_TIMER(&_profile_metrics.init_key_ns); |
460 | 0 | const auto& schema = _tablet->tablet_schema(); |
461 | | // Point query is only supported on merge-on-write unique key tables. |
462 | 0 | DCHECK(schema->keys_type() == UNIQUE_KEYS && _tablet->enable_unique_key_merge_on_write()); |
463 | 0 | if (schema->keys_type() != UNIQUE_KEYS || !_tablet->enable_unique_key_merge_on_write()) { |
464 | 0 | return Status::InvalidArgument( |
465 | 0 | "Point query is only supported on merge-on-write unique key tables, " |
466 | 0 | "tablet_id={}", |
467 | 0 | _tablet->tablet_id()); |
468 | 0 | } |
469 | | // 1. get primary key from conditions |
470 | 0 | _row_read_ctxs.resize(request->key_tuples().size()); |
471 | | // get row cursor and encode keys |
472 | 0 | for (int i = 0; i < request->key_tuples().size(); ++i) { |
473 | 0 | const KeyTuple& key_tuple = request->key_tuples(i); |
474 | 0 | if (UNLIKELY(cast_set<size_t>(key_tuple.key_column_literals_size()) != |
475 | 0 | schema->num_key_columns())) { |
476 | 0 | return Status::InvalidArgument( |
477 | 0 | "Key column count mismatch. expected={}, actual={}, tablet_id={}", |
478 | 0 | schema->num_key_columns(), key_tuple.key_column_literals_size(), |
479 | 0 | _tablet->tablet_id()); |
480 | 0 | } |
481 | 0 | RowCursor cursor; |
482 | 0 | std::vector<Field> key_fields; |
483 | 0 | key_fields.reserve(key_tuple.key_column_literals_size()); |
484 | 0 | for (int j = 0; j < key_tuple.key_column_literals_size(); ++j) { |
485 | 0 | const auto& literal_bytes = key_tuple.key_column_literals(j); |
486 | 0 | TExprNode expr_node; |
487 | 0 | auto len = cast_set<uint32_t>(literal_bytes.size()); |
488 | 0 | RETURN_IF_ERROR( |
489 | 0 | deserialize_thrift_msg(reinterpret_cast<const uint8_t*>(literal_bytes.data()), |
490 | 0 | &len, false, &expr_node)); |
491 | 0 | const auto& col = schema->column(j); |
492 | 0 | auto data_type = DataTypeFactory::instance().create_data_type( |
493 | 0 | col.type(), col.precision(), col.frac(), col.length()); |
494 | 0 | key_fields.push_back(data_type->get_field(expr_node)); |
495 | 0 | } |
496 | 0 | RETURN_IF_ERROR(cursor.init_scan_key(_tablet->tablet_schema(), std::move(key_fields))); |
497 | 0 | cursor.encode_key_with_padding<true>(&_row_read_ctxs[i]._primary_key, |
498 | 0 | _tablet->tablet_schema()->num_key_columns(), true); |
499 | 0 | } |
500 | 0 | return Status::OK(); |
501 | 0 | } |
502 | | |
503 | 0 | Status PointQueryExecutor::_lookup_row_key() { |
504 | 0 | SCOPED_TIMER(&_profile_metrics.lookup_key_ns); |
505 | | // 2. lookup row location |
506 | 0 | Status st; |
507 | 0 | if (_version >= 0) { |
508 | 0 | CHECK(config::is_cloud_mode()) << "Only cloud mode support snapshot read at present"; |
509 | 0 | SyncOptions options; |
510 | 0 | options.query_version = _version; |
511 | 0 | RETURN_IF_ERROR(std::dynamic_pointer_cast<CloudTablet>(_tablet)->sync_rowsets(options)); |
512 | 0 | } |
513 | 0 | std::vector<RowsetSharedPtr> specified_rowsets; |
514 | 0 | { |
515 | 0 | std::shared_lock rlock(_tablet->get_header_lock()); |
516 | 0 | specified_rowsets = _tablet->get_rowset_by_ids(nullptr); |
517 | 0 | } |
518 | 0 | io::IOContext io_ctx; |
519 | 0 | io_ctx.reader_type = ReaderType::READER_QUERY; |
520 | 0 | io_ctx.file_cache_stats = &_profile_metrics.read_stats.file_cache_stats; |
521 | 0 | io_ctx.remote_scan_cache_write_limiter = _remote_scan_cache_write_limiter.get(); |
522 | 0 | std::vector<std::unique_ptr<SegmentCacheHandle>> segment_caches(specified_rowsets.size()); |
523 | 0 | for (size_t i = 0; i < _row_read_ctxs.size(); ++i) { |
524 | 0 | RowLocation location; |
525 | | // The row cache contains the physical JSONB row but not the owning rowset's version/TSO. |
526 | | // A query projecting read-time hidden columns must resolve the rowset before decoding it. |
527 | 0 | if (!config::disable_storage_row_cache && !_reusable->has_read_time_hidden_columns()) { |
528 | 0 | RowCache::CacheHandle cache_handle; |
529 | 0 | auto hit_cache = RowCache::instance()->lookup( |
530 | 0 | {_tablet->tablet_id(), _row_read_ctxs[i]._primary_key}, &cache_handle); |
531 | 0 | if (hit_cache) { |
532 | 0 | _row_read_ctxs[i]._cached_row_data = std::move(cache_handle); |
533 | 0 | ++_profile_metrics.row_cache_hits; |
534 | 0 | continue; |
535 | 0 | } |
536 | 0 | } |
537 | | // Get rowlocation and rowset, ctx._rowset_ptr will acquire wrap this ptr |
538 | 0 | auto rowset_ptr = std::make_unique<RowsetSharedPtr>(); |
539 | 0 | st = (_tablet->lookup_row_key(_row_read_ctxs[i]._primary_key, nullptr, false, |
540 | 0 | specified_rowsets, &location, INT32_MAX /*rethink?*/, |
541 | 0 | segment_caches, rowset_ptr.get(), false, nullptr, |
542 | 0 | &_profile_metrics.read_stats, nullptr, &io_ctx)); |
543 | 0 | if (st.is<ErrorCode::KEY_NOT_FOUND>()) { |
544 | 0 | continue; |
545 | 0 | } |
546 | 0 | RETURN_IF_ERROR(st); |
547 | 0 | _row_read_ctxs[i]._row_location = location; |
548 | | // acquire and wrap this rowset |
549 | 0 | (*rowset_ptr)->acquire(); |
550 | 0 | VLOG_DEBUG << "aquire rowset " << (*rowset_ptr)->rowset_id(); |
551 | 0 | _row_read_ctxs[i]._rowset_ptr = std::unique_ptr<RowsetSharedPtr, decltype(&release_rowset)>( |
552 | 0 | rowset_ptr.release(), &release_rowset); |
553 | 0 | _row_hits++; |
554 | 0 | } |
555 | 0 | return Status::OK(); |
556 | 0 | } |
557 | | |
558 | 0 | Status PointQueryExecutor::_lookup_row_data() { |
559 | | // 3. get values |
560 | 0 | SCOPED_TIMER(&_profile_metrics.lookup_data_ns); |
561 | 0 | { |
562 | 0 | auto result_columns_guard = _result_block->mutate_columns_scoped(); |
563 | 0 | MutableColumns& result_columns = result_columns_guard.mutable_columns(); |
564 | 0 | for (size_t i = 0; i < _row_read_ctxs.size(); ++i) { |
565 | 0 | if (_row_read_ctxs[i]._cached_row_data.valid()) { |
566 | 0 | DORIS_CHECK(!_reusable->has_read_time_hidden_columns()); |
567 | 0 | RETURN_IF_ERROR(JsonbSerializeUtil::jsonb_to_columns( |
568 | 0 | _reusable->get_data_type_serdes(), |
569 | 0 | _row_read_ctxs[i]._cached_row_data.data().data, |
570 | 0 | _row_read_ctxs[i]._cached_row_data.data().size, |
571 | 0 | _reusable->get_col_uid_to_idx(), result_columns, |
572 | 0 | _reusable->get_col_default_values(), _reusable->include_col_uids())); |
573 | 0 | continue; |
574 | 0 | } |
575 | 0 | if (!_row_read_ctxs[i]._row_location.has_value()) { |
576 | 0 | continue; |
577 | 0 | } |
578 | 0 | auto rowset = get_beta_rowset(_row_read_ctxs[i]._rowset_ptr.get()); |
579 | 0 | std::string value; |
580 | | // fill block by row store |
581 | 0 | if (_reusable->rs_column_uid() != -1) { |
582 | 0 | bool use_row_cache = !config::disable_storage_row_cache; |
583 | 0 | io::IOContext io_ctx; |
584 | 0 | io_ctx.reader_type = ReaderType::READER_QUERY; |
585 | 0 | io_ctx.file_cache_stats = &_profile_metrics.read_stats.file_cache_stats; |
586 | 0 | io_ctx.remote_scan_cache_write_limiter = _remote_scan_cache_write_limiter.get(); |
587 | 0 | RETURN_IF_ERROR(_tablet->lookup_row_data( |
588 | 0 | _row_read_ctxs[i]._primary_key, _row_read_ctxs[i]._row_location.value(), |
589 | 0 | *(_row_read_ctxs[i]._rowset_ptr), _profile_metrics.read_stats, value, |
590 | 0 | use_row_cache, &io_ctx)); |
591 | | // serialize value to block, currently only jsonb row format |
592 | 0 | RETURN_IF_ERROR(JsonbSerializeUtil::jsonb_to_columns( |
593 | 0 | _reusable->get_data_type_serdes(), value.data(), value.size(), |
594 | 0 | _reusable->get_col_uid_to_idx(), result_columns, |
595 | 0 | _reusable->get_col_default_values(), _reusable->include_col_uids())); |
596 | 0 | } |
597 | 0 | if (!_reusable->missing_col_uids().empty()) { |
598 | 0 | if (!_reusable->runtime_state()->enable_short_circuit_query_access_column_store()) { |
599 | 0 | std::string missing_columns; |
600 | 0 | for (int cid : _reusable->missing_col_uids()) { |
601 | 0 | missing_columns += |
602 | 0 | _tablet->tablet_schema()->column_by_uid(cid).name() + ","; |
603 | 0 | } |
604 | 0 | return Status::InternalError( |
605 | 0 | "Not support column store, set store_row_column=true or " |
606 | 0 | "row_store_columns in table properties, missing columns: " + |
607 | 0 | missing_columns + " should be added to row store"); |
608 | 0 | } |
609 | | // fill missing columns by column store |
610 | 0 | RowLocation row_loc = _row_read_ctxs[i]._row_location.value(); |
611 | 0 | SegmentCacheHandle segment_cache; |
612 | 0 | io::IOContext io_ctx; |
613 | 0 | io_ctx.reader_type = ReaderType::READER_QUERY; |
614 | 0 | io_ctx.file_cache_stats = &_read_stats.file_cache_stats; |
615 | 0 | io_ctx.remote_scan_cache_write_limiter = _remote_scan_cache_write_limiter.get(); |
616 | 0 | { |
617 | 0 | SCOPED_TIMER(&_profile_metrics.load_segment_data_stage_ns); |
618 | 0 | RETURN_IF_ERROR(SegmentLoader::instance()->load_segments( |
619 | 0 | rowset, &segment_cache, true, false, &_read_stats, &io_ctx)); |
620 | 0 | } |
621 | | // find segment |
622 | 0 | auto it = std::find_if(segment_cache.get_segments().cbegin(), |
623 | 0 | segment_cache.get_segments().cend(), |
624 | 0 | [&](const segment_v2::SegmentSharedPtr& seg) { |
625 | 0 | return seg->id() == row_loc.segment_id; |
626 | 0 | }); |
627 | 0 | const auto& segment = *it; |
628 | 0 | for (int cid : _reusable->missing_col_uids()) { |
629 | 0 | int pos = _reusable->get_col_uid_to_idx().at(cid); |
630 | 0 | std::vector<segment_v2::rowid_t> row_ids { |
631 | 0 | static_cast<segment_v2::rowid_t>(row_loc.row_id)}; |
632 | 0 | auto& column = result_columns[pos]; |
633 | 0 | std::unique_ptr<ColumnIterator> iter; |
634 | 0 | SlotDescriptor* slot = _reusable->tuple_desc()->slots()[pos]; |
635 | 0 | StorageReadOptions storage_read_options; |
636 | 0 | storage_read_options.stats = &_read_stats; |
637 | 0 | storage_read_options.io_ctx = io_ctx; |
638 | 0 | RETURN_IF_ERROR(segment->seek_and_read_by_rowid(*_tablet->tablet_schema(), slot, |
639 | 0 | row_ids, column, |
640 | 0 | storage_read_options, iter)); |
641 | 0 | } |
642 | 0 | } |
643 | 0 | replace_point_query_read_time_hidden_columns(_reusable->read_time_hidden_columns(), |
644 | 0 | *_tablet->tablet_schema(), *rowset, |
645 | 0 | result_columns); |
646 | 0 | } |
647 | 0 | if (result_columns.size() > _reusable->include_col_uids().size()) { |
648 | | // Padding rows for some columns that no need to output to mysql client |
649 | | // eg. SELECT k1,v1,v2 FROM TABLE WHERE k1 = 1, k1 is not in output slots, tuple as bellow |
650 | | // TupleDescriptor{id=1, tbl=table_with_column_group} |
651 | | // SlotDescriptor{id=8, col=v1, colUniqueId=1 ...} |
652 | | // SlotDescriptor{id=9, col=v2, colUniqueId=2 ...} |
653 | | // thus missing in include_col_uids and missing_col_uids |
654 | 0 | for (auto& column : result_columns) { |
655 | 0 | int padding_rows = _row_hits - cast_set<int>(column->size()); |
656 | 0 | if (padding_rows > 0) { |
657 | 0 | column->insert_many_defaults(padding_rows); |
658 | 0 | } |
659 | 0 | } |
660 | 0 | } |
661 | 0 | } |
662 | | // filter rows by delete sign |
663 | 0 | if (_row_hits > 0 && _reusable->delete_sign_idx() != -1) { |
664 | 0 | size_t filtered = 0; |
665 | 0 | size_t total = 0; |
666 | 0 | { |
667 | | // clear_column_data will check reference of ColumnPtr, so we need to release |
668 | | // reference before clear_column_data |
669 | 0 | ColumnPtr delete_filter_columns = |
670 | 0 | _result_block->get_columns()[_reusable->delete_sign_idx()]; |
671 | 0 | const auto& filter = |
672 | 0 | assert_cast<const ColumnInt8*>(delete_filter_columns.get())->get_data(); |
673 | 0 | filtered = filter.size() - simd::count_zero_num((int8_t*)filter.data(), filter.size()); |
674 | 0 | total = filter.size(); |
675 | 0 | } |
676 | |
|
677 | 0 | if (filtered == total) { |
678 | 0 | _result_block->clear_column_data(); |
679 | 0 | } else if (filtered > 0) { |
680 | 0 | return Status::NotSupported("Not implemented since only single row at present"); |
681 | 0 | } |
682 | 0 | } |
683 | 0 | return Status::OK(); |
684 | 0 | } |
685 | | |
686 | 0 | Status serialize_block(std::shared_ptr<TFetchDataResult> res, PTabletKeyLookupResponse* response) { |
687 | 0 | uint8_t* buf = nullptr; |
688 | 0 | uint32_t len = 0; |
689 | 0 | ThriftSerializer ser(false, 4096); |
690 | 0 | RETURN_IF_ERROR(ser.serialize(&(res->result_batch), &len, &buf)); |
691 | 0 | response->set_row_batch(std::string((const char*)buf, len)); |
692 | 0 | return Status::OK(); |
693 | 0 | } |
694 | | |
695 | 0 | Status PointQueryExecutor::_output_data() { |
696 | | // 4. exprs exec and serialize to mysql row batches |
697 | 0 | SCOPED_TIMER(&_profile_metrics.output_data_ns); |
698 | 0 | if (_result_block->rows()) { |
699 | 0 | RuntimeState state; |
700 | 0 | auto buffer = std::make_shared<PointQueryResultBlockBuffer>(&state); |
701 | | // TODO reuse mysql_writer |
702 | 0 | VMysqlResultWriter mysql_writer(buffer, _reusable->output_exprs(), nullptr, |
703 | 0 | _binary_row_format); |
704 | 0 | RETURN_IF_ERROR(mysql_writer.init(_reusable->runtime_state())); |
705 | 0 | _result_block->clear_names(); |
706 | 0 | RETURN_IF_ERROR(mysql_writer.write(_reusable->runtime_state(), *_result_block)); |
707 | 0 | RETURN_IF_ERROR(serialize_block(buffer->get_block(), _response)); |
708 | 0 | VLOG_DEBUG << "dump block " << _result_block->dump_data(); |
709 | 0 | } else { |
710 | 0 | _response->set_empty_batch(true); |
711 | 0 | } |
712 | 0 | _profile_metrics.result_data_bytes = _result_block->bytes(); |
713 | 0 | _reusable->return_block(_result_block); |
714 | 0 | return Status::OK(); |
715 | 0 | } |
716 | | |
717 | | } // namespace doris |