Coverage Report

Created: 2026-08-06 18:25

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/scan/meta_scanner.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 "exec/scan/meta_scanner.h"
19
20
#include <fmt/format.h>
21
#include <gen_cpp/FrontendService.h>
22
#include <gen_cpp/FrontendService_types.h>
23
#include <gen_cpp/HeartbeatService_types.h>
24
#include <gen_cpp/PaloInternalService_types.h>
25
#include <gen_cpp/PlanNodes_types.h>
26
27
#include <ostream>
28
#include <string>
29
#include <unordered_map>
30
31
#include "common/cast_set.h"
32
#include "common/logging.h"
33
#include "core/block/block.h"
34
#include "core/column/column.h"
35
#include "core/column/column_nullable.h"
36
#include "core/column/column_string.h"
37
#include "core/column/column_vector.h"
38
#include "core/data_type/define_primitive_type.h"
39
#include "core/types.h"
40
#include "format/table/parquet_metadata_reader.h"
41
#include "runtime/cluster_info.h"
42
#include "runtime/descriptors.h"
43
#include "runtime/exec_env.h"
44
#include "runtime/runtime_state.h"
45
#include "util/client_cache.h"
46
#include "util/thrift_rpc_helper.h"
47
48
namespace doris {
49
class RuntimeProfile;
50
class VExprContext;
51
} // namespace doris
52
53
namespace doris {
54
55
MetaScanner::MetaScanner(RuntimeState* state, ScanLocalStateBase* local_state, TupleId tuple_id,
56
                         const TScanRangeParams& scan_range, int64_t limit, RuntimeProfile* profile,
57
                         TUserIdentity user_identity)
58
0
        : Scanner(state, local_state, limit, profile),
59
0
          _meta_eos(false),
60
0
          _tuple_id(tuple_id),
61
0
          _user_identity(user_identity),
62
0
          _scan_range(scan_range.scan_range) {}
63
64
0
Status MetaScanner::_open_impl(RuntimeState* state) {
65
0
    VLOG_CRITICAL << "MetaScanner::open";
66
0
    RETURN_IF_ERROR(Scanner::_open_impl(state));
67
0
    if (_scan_range.meta_scan_range.metadata_type == TMetadataType::PARQUET) {
68
0
        auto reader = ParquetMetadataReader::create_unique(_tuple_desc->slots(), state, _profile,
69
0
                                                           _scan_range.meta_scan_range);
70
0
        RETURN_IF_ERROR(reader->init_reader());
71
0
        _reader = std::move(reader);
72
0
    } else {
73
0
        RETURN_IF_ERROR(_fetch_metadata(_scan_range.meta_scan_range));
74
0
    }
75
0
    return Status::OK();
76
0
}
77
78
0
Status MetaScanner::init(RuntimeState* state, const VExprContextSPtrs& conjuncts) {
79
0
    VLOG_CRITICAL << "MetaScanner::init";
80
0
    RETURN_IF_ERROR(Scanner::init(_state, conjuncts));
81
0
    _tuple_desc = state->desc_tbl().get_tuple_descriptor(_tuple_id);
82
0
    return Status::OK();
83
0
}
84
85
0
Status MetaScanner::_get_block_impl(RuntimeState* state, Block* block, bool* eof) {
86
0
    VLOG_CRITICAL << "MetaScanner::_get_block_impl";
87
0
    if (nullptr == state || nullptr == block || nullptr == eof) {
88
0
        return Status::InternalError("input is NULL pointer");
89
0
    }
90
91
    // Build name to index map only once on first call
92
0
    if (_src_block_name_to_idx.empty()) {
93
0
        _src_block_name_to_idx = block->get_name_to_pos_map();
94
0
    }
95
96
0
    if (_reader) {
97
        // TODO: This is a temporary workaround; the code is planned to be refactored later.
98
0
        size_t read_rows = 0;
99
0
        return _reader->get_next_block(block, &read_rows, eof);
100
0
    }
101
102
0
    if (_meta_eos == true) {
103
0
        *eof = true;
104
0
        return Status::OK();
105
0
    }
106
107
0
    auto column_size = _tuple_desc->slots().size();
108
0
    std::vector<MutableColumnPtr> columns(column_size);
109
0
    bool mem_reuse = block->mem_reuse();
110
0
    do {
111
0
        RETURN_IF_CANCELLED(state);
112
113
0
        columns.resize(column_size);
114
0
        for (auto i = 0; i < column_size; i++) {
115
0
            if (mem_reuse) {
116
0
                columns[i] = IColumn::mutate(std::move(block->get_by_position(i).column));
117
0
            } else {
118
0
                columns[i] = _tuple_desc->slots()[i]->get_empty_mutable_column();
119
0
            }
120
0
        }
121
        // fill block
122
0
        RETURN_IF_ERROR(_fill_block_with_remote_data(columns));
123
0
        const bool empty_result = columns.empty() || columns.front()->empty();
124
0
        if (!mem_reuse) {
125
0
            int column_index = 0;
126
0
            for (const auto slot_desc : _tuple_desc->slots()) {
127
0
                block->insert(ColumnWithTypeAndName(std::move(columns[column_index++]),
128
0
                                                    slot_desc->get_data_type_ptr(),
129
0
                                                    slot_desc->col_name()));
130
0
            }
131
0
        } else {
132
0
            block->set_columns(std::move(columns));
133
0
        }
134
0
        if (_meta_eos == true) {
135
0
            if (empty_result) {
136
0
                *eof = true;
137
0
            }
138
0
            break;
139
0
        }
140
0
        VLOG_ROW << "VMetaScanNode output rows: " << block->rows();
141
0
    } while (block->rows() == 0 && !(*eof));
142
0
    return Status::OK();
143
0
}
144
145
0
Status MetaScanner::_fill_block_with_remote_data(const std::vector<MutableColumnPtr>& columns) {
146
0
    VLOG_CRITICAL << "MetaScanner::_fill_block_with_remote_data";
147
0
    for (int col_idx = 0; col_idx < columns.size(); col_idx++) {
148
0
        auto slot_desc = _tuple_desc->slots()[col_idx];
149
150
0
        for (int _row_idx = 0; _row_idx < _batch_data.size(); _row_idx++) {
151
0
            IColumn* col_ptr = columns[col_idx].get();
152
0
            TCell& cell = _batch_data[_row_idx].column_value[col_idx];
153
0
            if (cell.__isset.isNull && cell.isNull) {
154
0
                DCHECK(slot_desc->is_nullable())
155
0
                        << "cell is null but column is not nullable: " << slot_desc->col_name();
156
0
                auto& null_col = reinterpret_cast<ColumnNullable&>(*col_ptr);
157
0
                null_col.get_nested_column().insert_default();
158
0
                null_col.get_null_map_data().push_back(1);
159
0
            } else {
160
0
                if (slot_desc->is_nullable()) {
161
0
                    auto& null_col = reinterpret_cast<ColumnNullable&>(*col_ptr);
162
0
                    null_col.get_null_map_data().push_back(0);
163
0
                    col_ptr = null_col.get_nested_column_ptr().get();
164
0
                }
165
0
                switch (slot_desc->type()->get_primitive_type()) {
166
0
                case TYPE_BOOLEAN: {
167
0
                    bool data = cell.boolVal;
168
0
                    assert_cast<ColumnBool*>(col_ptr)->insert_value((uint8_t)data);
169
0
                    break;
170
0
                }
171
0
                case TYPE_TINYINT: {
172
0
                    int8_t data = (int8_t)cell.intVal;
173
0
                    assert_cast<ColumnInt8*>(col_ptr)->insert_value(data);
174
0
                    break;
175
0
                }
176
0
                case TYPE_SMALLINT: {
177
0
                    int16_t data = (int16_t)cell.intVal;
178
0
                    assert_cast<ColumnInt16*>(col_ptr)->insert_value(data);
179
0
                    break;
180
0
                }
181
0
                case TYPE_INT: {
182
0
                    int32_t data = cell.intVal;
183
0
                    assert_cast<ColumnInt32*>(col_ptr)->insert_value(data);
184
0
                    break;
185
0
                }
186
0
                case TYPE_BIGINT: {
187
0
                    int64_t data = cell.longVal;
188
0
                    assert_cast<ColumnInt64*>(col_ptr)->insert_value(data);
189
0
                    break;
190
0
                }
191
0
                case TYPE_FLOAT: {
192
0
                    auto data = static_cast<float>(cell.doubleVal);
193
0
                    assert_cast<ColumnFloat32*>(col_ptr)->insert_value(data);
194
0
                    break;
195
0
                }
196
0
                case TYPE_DOUBLE: {
197
0
                    double data = cell.doubleVal;
198
0
                    assert_cast<ColumnFloat64*>(col_ptr)->insert_value(data);
199
0
                    break;
200
0
                }
201
0
                case TYPE_DATEV2: {
202
0
                    uint32_t data = (uint32_t)cell.longVal;
203
0
                    assert_cast<ColumnDateV2*>(col_ptr)->insert_value(data);
204
0
                    break;
205
0
                }
206
0
                case TYPE_DATETIMEV2: {
207
0
                    uint64_t data = cell.longVal;
208
0
                    assert_cast<ColumnDateTimeV2*>(col_ptr)->insert_value(data);
209
0
                    break;
210
0
                }
211
0
                case TYPE_STRING:
212
0
                case TYPE_CHAR:
213
0
                case TYPE_VARCHAR: {
214
0
                    std::string data = cell.stringVal;
215
0
                    assert_cast<ColumnString*>(col_ptr)->insert_data(data.c_str(), data.length());
216
0
                    break;
217
0
                }
218
0
                default: {
219
0
                    std::string error_msg =
220
0
                            fmt::format("Invalid column type {} on column: {}.",
221
0
                                        slot_desc->type()->get_name(), slot_desc->col_name());
222
0
                    return Status::InternalError(std::string(error_msg));
223
0
                }
224
0
                }
225
0
            }
226
0
        }
227
0
    }
228
0
    _meta_eos = true;
229
0
    return Status::OK();
230
0
}
231
232
0
Status MetaScanner::_fetch_metadata(const TMetaScanRange& meta_scan_range) {
233
0
    VLOG_CRITICAL << "MetaScanner::_fetch_metadata";
234
0
    TFetchSchemaTableDataRequest request;
235
0
    switch (meta_scan_range.metadata_type) {
236
0
    case TMetadataType::BACKENDS:
237
0
        RETURN_IF_ERROR(_build_backends_metadata_request(meta_scan_range, &request));
238
0
        break;
239
0
    case TMetadataType::FRONTENDS:
240
0
        RETURN_IF_ERROR(_build_frontends_metadata_request(meta_scan_range, &request));
241
0
        break;
242
0
    case TMetadataType::FRONTENDS_DISKS:
243
0
        RETURN_IF_ERROR(_build_frontends_disks_metadata_request(meta_scan_range, &request));
244
0
        break;
245
0
    case TMetadataType::WORKLOAD_SCHED_POLICY:
246
0
        RETURN_IF_ERROR(_build_workload_sched_policy_metadata_request(meta_scan_range, &request));
247
0
        break;
248
0
    case TMetadataType::CATALOGS:
249
0
        RETURN_IF_ERROR(_build_catalogs_metadata_request(meta_scan_range, &request));
250
0
        break;
251
0
    case TMetadataType::MATERIALIZED_VIEWS:
252
0
        RETURN_IF_ERROR(_build_materialized_views_metadata_request(meta_scan_range, &request));
253
0
        break;
254
0
    case TMetadataType::PARTITIONS:
255
0
        RETURN_IF_ERROR(_build_partitions_metadata_request(meta_scan_range, &request));
256
0
        break;
257
0
    case TMetadataType::JOBS:
258
0
        RETURN_IF_ERROR(_build_jobs_metadata_request(meta_scan_range, &request));
259
0
        break;
260
0
    case TMetadataType::TASKS:
261
0
        RETURN_IF_ERROR(_build_tasks_metadata_request(meta_scan_range, &request));
262
0
        break;
263
0
    case TMetadataType::PARTITION_VALUES:
264
0
        RETURN_IF_ERROR(_build_partition_values_metadata_request(meta_scan_range, &request));
265
0
        break;
266
0
    default:
267
0
        _meta_eos = true;
268
0
        return Status::OK();
269
0
    }
270
271
    // set filter columns
272
0
    std::vector<std::string> filter_columns;
273
0
    for (const auto& slot : _tuple_desc->slots()) {
274
0
        filter_columns.emplace_back(slot->col_name_lower_case());
275
0
    }
276
0
    request.metada_table_params.__set_columns_name(filter_columns);
277
278
    // _state->execution_timeout() is seconds, change to milliseconds
279
0
    int time_out = _state->execution_timeout() * 1000;
280
0
    TNetworkAddress master_addr = ExecEnv::GetInstance()->cluster_info()->master_fe_addr;
281
0
    TFetchSchemaTableDataResult result;
282
0
    RETURN_IF_ERROR(ThriftRpcHelper::rpc<FrontendServiceClient>(
283
0
            master_addr.hostname, master_addr.port,
284
0
            [&request, &result](FrontendServiceConnection& client) {
285
0
                client->fetchSchemaTableData(result, request);
286
0
            },
287
0
            time_out));
288
289
0
    Status status(Status::create(result.status));
290
0
    if (!status.ok()) {
291
0
        LOG(WARNING) << "fetch schema table data from master failed, errmsg=" << status;
292
0
        return status;
293
0
    }
294
0
    _batch_data = std::move(result.data_batch);
295
0
    return Status::OK();
296
0
}
297
298
Status MetaScanner::_build_backends_metadata_request(const TMetaScanRange& meta_scan_range,
299
0
                                                     TFetchSchemaTableDataRequest* request) {
300
0
    VLOG_CRITICAL << "MetaScanner::_build_backends_metadata_request";
301
0
    if (!meta_scan_range.__isset.backends_params) {
302
0
        return Status::InternalError("Can not find TBackendsMetadataParams from meta_scan_range.");
303
0
    }
304
    // create request
305
0
    request->__set_cluster_name("");
306
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
307
308
    // create TMetadataTableRequestParams
309
0
    TMetadataTableRequestParams metadata_table_params;
310
0
    metadata_table_params.__set_metadata_type(TMetadataType::BACKENDS);
311
0
    metadata_table_params.__set_backends_metadata_params(meta_scan_range.backends_params);
312
313
0
    request->__set_metada_table_params(metadata_table_params);
314
0
    return Status::OK();
315
0
}
316
317
Status MetaScanner::_build_frontends_metadata_request(const TMetaScanRange& meta_scan_range,
318
0
                                                      TFetchSchemaTableDataRequest* request) {
319
0
    VLOG_CRITICAL << "MetaScanner::_build_frontends_metadata_request";
320
0
    if (!meta_scan_range.__isset.frontends_params) {
321
0
        return Status::InternalError("Can not find TFrontendsMetadataParams from meta_scan_range.");
322
0
    }
323
    // create request
324
0
    request->__set_cluster_name("");
325
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
326
327
    // create TMetadataTableRequestParams
328
0
    TMetadataTableRequestParams metadata_table_params;
329
0
    metadata_table_params.__set_metadata_type(TMetadataType::FRONTENDS);
330
0
    metadata_table_params.__set_frontends_metadata_params(meta_scan_range.frontends_params);
331
332
0
    request->__set_metada_table_params(metadata_table_params);
333
0
    return Status::OK();
334
0
}
335
336
Status MetaScanner::_build_frontends_disks_metadata_request(const TMetaScanRange& meta_scan_range,
337
0
                                                            TFetchSchemaTableDataRequest* request) {
338
0
    VLOG_CRITICAL << "MetaScanner::_build_frontends_metadata_request";
339
0
    if (!meta_scan_range.__isset.frontends_params) {
340
0
        return Status::InternalError("Can not find TFrontendsMetadataParams from meta_scan_range.");
341
0
    }
342
    // create request
343
0
    request->__set_cluster_name("");
344
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
345
346
    // create TMetadataTableRequestParams
347
0
    TMetadataTableRequestParams metadata_table_params;
348
0
    metadata_table_params.__set_metadata_type(TMetadataType::FRONTENDS_DISKS);
349
0
    metadata_table_params.__set_frontends_metadata_params(meta_scan_range.frontends_params);
350
351
0
    request->__set_metada_table_params(metadata_table_params);
352
0
    return Status::OK();
353
0
}
354
355
Status MetaScanner::_build_workload_sched_policy_metadata_request(
356
0
        const TMetaScanRange& meta_scan_range, TFetchSchemaTableDataRequest* request) {
357
0
    VLOG_CRITICAL << "MetaScanner::_build_workload_sched_policy_metadata_request";
358
359
    // create request
360
0
    request->__set_cluster_name("");
361
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
362
363
    // create TMetadataTableRequestParams
364
0
    TMetadataTableRequestParams metadata_table_params;
365
0
    metadata_table_params.__set_metadata_type(TMetadataType::WORKLOAD_SCHED_POLICY);
366
0
    metadata_table_params.__set_current_user_ident(_user_identity);
367
368
0
    request->__set_metada_table_params(metadata_table_params);
369
0
    return Status::OK();
370
0
}
371
372
Status MetaScanner::_build_catalogs_metadata_request(const TMetaScanRange& meta_scan_range,
373
0
                                                     TFetchSchemaTableDataRequest* request) {
374
0
    VLOG_CRITICAL << "MetaScanner::_build_catalogs_metadata_request";
375
376
    // create request
377
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
378
379
    // create TMetadataTableRequestParams
380
0
    TMetadataTableRequestParams metadata_table_params;
381
0
    metadata_table_params.__set_metadata_type(TMetadataType::CATALOGS);
382
0
    metadata_table_params.__set_current_user_ident(_user_identity);
383
384
0
    request->__set_metada_table_params(metadata_table_params);
385
0
    return Status::OK();
386
0
}
387
388
Status MetaScanner::_build_materialized_views_metadata_request(
389
0
        const TMetaScanRange& meta_scan_range, TFetchSchemaTableDataRequest* request) {
390
0
    VLOG_CRITICAL << "MetaScanner::_build_materialized_views_metadata_request";
391
0
    if (!meta_scan_range.__isset.materialized_views_params) {
392
0
        return Status::InternalError(
393
0
                "Can not find TMaterializedViewsMetadataParams from meta_scan_range.");
394
0
    }
395
396
    // create request
397
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
398
399
    // create TMetadataTableRequestParams
400
0
    TMetadataTableRequestParams metadata_table_params;
401
0
    metadata_table_params.__set_metadata_type(TMetadataType::MATERIALIZED_VIEWS);
402
0
    metadata_table_params.__set_materialized_views_metadata_params(
403
0
            meta_scan_range.materialized_views_params);
404
405
0
    request->__set_metada_table_params(metadata_table_params);
406
0
    return Status::OK();
407
0
}
408
409
Status MetaScanner::_build_partitions_metadata_request(const TMetaScanRange& meta_scan_range,
410
0
                                                       TFetchSchemaTableDataRequest* request) {
411
0
    VLOG_CRITICAL << "MetaScanner::_build_partitions_metadata_request";
412
0
    if (!meta_scan_range.__isset.partitions_params) {
413
0
        return Status::InternalError(
414
0
                "Can not find TPartitionsMetadataParams from meta_scan_range.");
415
0
    }
416
417
    // create request
418
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
419
420
    // create TMetadataTableRequestParams
421
0
    TMetadataTableRequestParams metadata_table_params;
422
0
    metadata_table_params.__set_metadata_type(TMetadataType::PARTITIONS);
423
0
    metadata_table_params.__set_partitions_metadata_params(meta_scan_range.partitions_params);
424
425
0
    request->__set_metada_table_params(metadata_table_params);
426
0
    return Status::OK();
427
0
}
428
429
Status MetaScanner::_build_jobs_metadata_request(const TMetaScanRange& meta_scan_range,
430
0
                                                 TFetchSchemaTableDataRequest* request) {
431
0
    VLOG_CRITICAL << "MetaScanner::_build_jobs_metadata_request";
432
0
    if (!meta_scan_range.__isset.jobs_params) {
433
0
        return Status::InternalError("Can not find TJobsMetadataParams from meta_scan_range.");
434
0
    }
435
436
    // create request
437
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
438
439
    // create TMetadataTableRequestParams
440
0
    TMetadataTableRequestParams metadata_table_params;
441
0
    metadata_table_params.__set_metadata_type(TMetadataType::JOBS);
442
0
    metadata_table_params.__set_jobs_metadata_params(meta_scan_range.jobs_params);
443
444
0
    request->__set_metada_table_params(metadata_table_params);
445
0
    return Status::OK();
446
0
}
447
448
Status MetaScanner::_build_tasks_metadata_request(const TMetaScanRange& meta_scan_range,
449
0
                                                  TFetchSchemaTableDataRequest* request) {
450
0
    VLOG_CRITICAL << "MetaScanner::_build_tasks_metadata_request";
451
0
    if (!meta_scan_range.__isset.tasks_params) {
452
0
        return Status::InternalError("Can not find TTasksMetadataParams from meta_scan_range.");
453
0
    }
454
455
    // create request
456
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
457
458
    // create TMetadataTableRequestParams
459
0
    TMetadataTableRequestParams metadata_table_params;
460
0
    metadata_table_params.__set_metadata_type(TMetadataType::TASKS);
461
0
    metadata_table_params.__set_tasks_metadata_params(meta_scan_range.tasks_params);
462
463
0
    request->__set_metada_table_params(metadata_table_params);
464
0
    return Status::OK();
465
0
}
466
467
Status MetaScanner::_build_partition_values_metadata_request(
468
0
        const TMetaScanRange& meta_scan_range, TFetchSchemaTableDataRequest* request) {
469
0
    VLOG_CRITICAL << "MetaScanner::_build_partition_values_metadata_request";
470
0
    if (!meta_scan_range.__isset.partition_values_params) {
471
0
        return Status::InternalError(
472
0
                "Can not find TPartitionValuesMetadataParams from meta_scan_range.");
473
0
    }
474
475
    // create request
476
0
    request->__set_schema_table_name(TSchemaTableName::METADATA_TABLE);
477
478
    // create TMetadataTableRequestParams
479
0
    TMetadataTableRequestParams metadata_table_params;
480
0
    metadata_table_params.__set_metadata_type(TMetadataType::PARTITION_VALUES);
481
0
    metadata_table_params.__set_partition_values_metadata_params(
482
0
            meta_scan_range.partition_values_params);
483
484
0
    request->__set_metada_table_params(metadata_table_params);
485
0
    return Status::OK();
486
0
}
487
488
0
Status MetaScanner::close(RuntimeState* state) {
489
0
    VLOG_CRITICAL << "MetaScanner::close";
490
0
    if (!_try_close()) {
491
0
        return Status::OK();
492
0
    }
493
0
    if (_reader) {
494
0
        RETURN_IF_ERROR(_reader->close());
495
0
    }
496
0
    RETURN_IF_ERROR(Scanner::close(state));
497
0
    return Status::OK();
498
0
}
499
500
} // namespace doris