Coverage Report

Created: 2026-09-20 01:59

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
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