Coverage Report

Created: 2026-04-01 07:52

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/runtime/query_context.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 "runtime/query_context.h"
19
20
#include <fmt/core.h>
21
#include <gen_cpp/FrontendService_types.h>
22
#include <gen_cpp/RuntimeProfile_types.h>
23
#include <gen_cpp/Types_types.h>
24
#include <glog/logging.h>
25
26
#include <algorithm>
27
#include <exception>
28
#include <memory>
29
#include <mutex>
30
#include <utility>
31
#include <vector>
32
33
#include "common/logging.h"
34
#include "common/status.h"
35
#include "exec/operator/rec_cte_scan_operator.h"
36
#include "exec/pipeline/dependency.h"
37
#include "exec/pipeline/pipeline_fragment_context.h"
38
#include "exec/runtime_filter/runtime_filter_definitions.h"
39
#include "exec/spill/spill_file_manager.h"
40
#include "runtime/exec_env.h"
41
#include "runtime/fragment_mgr.h"
42
#include "runtime/memory/heap_profiler.h"
43
#include "runtime/runtime_query_statistics_mgr.h"
44
#include "runtime/runtime_state.h"
45
#include "runtime/thread_context.h"
46
#include "runtime/workload_group/workload_group_manager.h"
47
#include "runtime/workload_management/query_task_controller.h"
48
#include "storage/olap_common.h"
49
#include "util/mem_info.h"
50
#include "util/uid_util.h"
51
52
namespace doris {
53
54
class DelayReleaseToken : public Runnable {
55
    ENABLE_FACTORY_CREATOR(DelayReleaseToken);
56
57
public:
58
0
    DelayReleaseToken(std::unique_ptr<ThreadPoolToken>&& token) { token_ = std::move(token); }
59
    ~DelayReleaseToken() override = default;
60
0
    void run() override {}
61
    std::unique_ptr<ThreadPoolToken> token_;
62
};
63
64
0
const std::string toString(QuerySource queryType) {
65
0
    switch (queryType) {
66
0
    case QuerySource::INTERNAL_FRONTEND:
67
0
        return "INTERNAL_FRONTEND";
68
0
    case QuerySource::STREAM_LOAD:
69
0
        return "STREAM_LOAD";
70
0
    case QuerySource::GROUP_COMMIT_LOAD:
71
0
        return "EXTERNAL_QUERY";
72
0
    case QuerySource::ROUTINE_LOAD:
73
0
        return "ROUTINE_LOAD";
74
0
    case QuerySource::EXTERNAL_CONNECTOR:
75
0
        return "EXTERNAL_CONNECTOR";
76
0
    default:
77
0
        return "UNKNOWN";
78
0
    }
79
0
}
80
81
std::shared_ptr<QueryContext> QueryContext::create(TUniqueId query_id, ExecEnv* exec_env,
82
                                                   const TQueryOptions& query_options,
83
                                                   TNetworkAddress coord_addr, bool is_nereids,
84
                                                   TNetworkAddress current_connect_fe,
85
286k
                                                   QuerySource query_type) {
86
286k
    auto ctx = QueryContext::create_shared(query_id, exec_env, query_options, coord_addr,
87
286k
                                           is_nereids, current_connect_fe, query_type);
88
286k
    ctx->init_query_task_controller();
89
286k
    return ctx;
90
286k
}
91
92
QueryContext::QueryContext(TUniqueId query_id, ExecEnv* exec_env,
93
                           const TQueryOptions& query_options, TNetworkAddress coord_addr,
94
                           bool is_nereids, TNetworkAddress current_connect_fe,
95
                           QuerySource query_source)
96
408k
        : _timeout_second(-1),
97
408k
          _query_id(std::move(query_id)),
98
408k
          _exec_env(exec_env),
99
408k
          _is_nereids(is_nereids),
100
408k
          _query_options(query_options),
101
408k
          _query_source(query_source) {
102
408k
    _init_resource_context();
103
408k
    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(query_mem_tracker());
104
408k
    _query_watcher.start();
105
408k
    _execution_dependency = Dependency::create_unique(-1, -1, "ExecutionDependency", false);
106
408k
    _memory_sufficient_dependency =
107
408k
            Dependency::create_unique(-1, -1, "MemorySufficientDependency", true);
108
109
408k
    _runtime_filter_mgr = std::make_unique<RuntimeFilterMgr>(true);
110
111
408k
    _timeout_second = query_options.execution_timeout;
112
113
408k
    bool initialize_context_holder =
114
408k
            config::enable_file_cache && config::enable_file_cache_query_limit &&
115
408k
            query_options.__isset.enable_file_cache && query_options.enable_file_cache &&
116
408k
            query_options.__isset.file_cache_query_limit_percent &&
117
408k
            query_options.file_cache_query_limit_percent < 100;
118
119
    // Initialize file cache context holders
120
408k
    if (initialize_context_holder) {
121
10
        _query_context_holders = io::FileCacheFactory::instance()->get_query_context_holders(
122
10
                _query_id, query_options.file_cache_query_limit_percent);
123
10
    }
124
125
408k
    bool is_query_type_valid = query_options.query_type == TQueryType::SELECT ||
126
408k
                               query_options.query_type == TQueryType::LOAD ||
127
408k
                               query_options.query_type == TQueryType::EXTERNAL;
128
408k
    DCHECK_EQ(is_query_type_valid, true);
129
130
408k
    this->coord_addr = coord_addr;
131
    // current_connect_fe is used for report query statistics
132
408k
    this->current_connect_fe = current_connect_fe;
133
    // external query has no current_connect_fe
134
408k
    if (query_options.query_type != TQueryType::EXTERNAL) {
135
285k
        bool is_report_fe_addr_valid =
136
286k
                !this->current_connect_fe.hostname.empty() && this->current_connect_fe.port != 0;
137
285k
        DCHECK_EQ(is_report_fe_addr_valid, true);
138
285k
    }
139
408k
    clock_gettime(CLOCK_MONOTONIC, &this->_query_arrival_timestamp);
140
408k
    DorisMetrics::instance()->query_ctx_cnt->increment(1);
141
408k
    _mem_arb = MemShareArbitrator::create_shared(
142
            query_id, query_options.mem_limit,
143
406k
            query_options.__isset.max_scan_mem_ratio ? query_options.max_scan_mem_ratio : 1.0);
144
406k
}
145
18.4E
146
406k
void QueryContext::_init_query_mem_tracker() {
147
18.4E
    bool has_query_mem_limit = _query_options.__isset.mem_limit && (_query_options.mem_limit > 0);
148
18.4E
    int64_t bytes_limit = has_query_mem_limit ? _query_options.mem_limit : -1;
149
18.4E
    if (bytes_limit > MemInfo::mem_limit() || bytes_limit == -1) {
150
18.4E
        VLOG_NOTICE << "Query memory limit " << PrettyPrinter::print(bytes_limit, TUnit::BYTES)
151
165k
                    << " exceeds process memory limit of "
152
165k
                    << PrettyPrinter::print(MemInfo::mem_limit(), TUnit::BYTES)
153
                    << " OR is -1. Using process memory limit instead.";
154
        bytes_limit = MemInfo::mem_limit();
155
406k
    }
156
124k
    // If the query is a pure load task(streamload, routine load, group commit), then it should not use
157
124k
    // memlimit per query to limit their memory usage.
158
406k
    if (is_pure_load_task()) {
159
406k
        bytes_limit = MemInfo::mem_limit();
160
251k
    }
161
251k
    std::shared_ptr<MemTrackerLimiter> query_mem_tracker;
162
251k
    if (_query_options.query_type == TQueryType::SELECT) {
163
251k
        query_mem_tracker = MemTrackerLimiter::create_shared(
164
33.7k
                MemTrackerLimiter::Type::QUERY, fmt::format("Query#Id={}", print_id(_query_id)),
165
33.7k
                bytes_limit);
166
33.7k
    } else if (_query_options.query_type == TQueryType::LOAD) {
167
121k
        query_mem_tracker = MemTrackerLimiter::create_shared(
168
121k
                MemTrackerLimiter::Type::LOAD, fmt::format("Load#Id={}", print_id(_query_id)),
169
121k
                bytes_limit);
170
121k
    } else if (_query_options.query_type == TQueryType::EXTERNAL) { // spark/flink/etc..
171
18.4E
        query_mem_tracker = MemTrackerLimiter::create_shared(
172
18.4E
                MemTrackerLimiter::Type::QUERY, fmt::format("External#Id={}", print_id(_query_id)),
173
18.4E
                bytes_limit);
174
18.4E
    } else {
175
407k
        LOG(FATAL) << "__builtin_unreachable";
176
32.8k
        __builtin_unreachable();
177
32.8k
    }
178
    if (_query_options.__isset.is_report_success && _query_options.is_report_success) {
179
        query_mem_tracker->enable_print_log_usage();
180
    }
181
182
    // If enable reserve memory, not enable check limit, because reserve memory will check it.
183
    // If reserve enabled, even if the reserved memory size is smaller than the actual requested memory,
184
    // and the query memory consumption is larger than the limit, we do not expect the query to fail
185
407k
    // after `check_limit` returns an error, but to run as long as possible,
186
407k
    // and will enter the paused state and try to spill when the query reserves next time.
187
407k
    // If the workload group or process runs out of memory, it will be forced to cancel.
188
407k
    query_mem_tracker->set_enable_check_limit(!(_query_options.__isset.enable_reserve_memory &&
189
                                                _query_options.enable_reserve_memory));
190
408k
    _resource_ctx->memory_context()->set_mem_tracker(query_mem_tracker);
191
408k
}
192
408k
193
408k
void QueryContext::_init_resource_context() {
194
    _resource_ctx = ResourceContext::create_shared();
195
404k
    _init_query_mem_tracker();
196
404k
}
197
404k
198
404k
void QueryContext::init_query_task_controller() {
199
404k
    _resource_ctx->set_task_controller(QueryTaskController::create(shared_from_this()));
200
404k
    _resource_ctx->task_controller()->set_task_id(_query_id);
201
404k
    _resource_ctx->task_controller()->set_fe_addr(current_connect_fe);
202
404k
    _resource_ctx->task_controller()->set_query_type(_query_options.query_type);
203
404k
#ifndef BE_TEST
204
404k
    _exec_env->runtime_query_statistics_mgr()->register_resource_context(print_id(_query_id),
205
                                                                         _resource_ctx);
206
286k
#endif
207
286k
}
208
209
QueryContext::~QueryContext() {
210
    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(query_mem_tracker());
211
    // query mem tracker consumption is equal to 0, it means that after QueryContext is created,
212
286k
    // it is found that query already exists in _query_ctx_map, and query mem tracker is not used.
213
286k
    // query mem tracker consumption is not equal to 0 after use, because there is memory consumed
214
285k
    // on query mem tracker, released on other trackers.
215
285k
    std::string mem_tracker_msg;
216
285k
    if (query_mem_tracker()->peak_consumption() != 0) {
217
285k
        mem_tracker_msg = fmt::format(
218
285k
                "deregister query/load memory tracker, queryId={}, Limit={}, CurrUsed={}, "
219
285k
                "PeakUsed={}",
220
285k
                print_id(_query_id), PrettyPrinter::print_bytes(query_mem_tracker()->limit()),
221
286k
                PrettyPrinter::print_bytes(query_mem_tracker()->consumption()),
222
286k
                PrettyPrinter::print_bytes(query_mem_tracker()->peak_consumption()));
223
286k
    }
224
286k
    [[maybe_unused]] uint64_t group_id = 0;
225
    if (workload_group()) {
226
286k
        group_id = workload_group()->id(); // before remove
227
    }
228
286k
229
1.45k
    _resource_ctx->task_controller()->finish();
230
1.45k
231
    if (enable_profile()) {
232
286k
        _report_query_profile();
233
286k
    }
234
0
235
0
#ifndef BE_TEST
236
0
    if (ExecEnv::GetInstance()->pipeline_tracer_context()->enabled()) [[unlikely]] {
237
0
        try {
238
0
            ExecEnv::GetInstance()->pipeline_tracer_context()->end_query(_query_id, group_id);
239
0
        } catch (std::exception& e) {
240
286k
            LOG(WARNING) << "Dump trace log failed bacause " << e.what();
241
286k
        }
242
286k
    }
243
286k
#endif
244
286k
    _runtime_filter_mgr.reset();
245
286k
    _execution_dependency.reset();
246
286k
    _runtime_predicates.clear();
247
    file_scan_range_params_map.clear();
248
286k
    obj_pool.clear();
249
    _merge_controller_handler.reset();
250
286k
251
286k
    DorisMetrics::instance()->query_ctx_cnt->increment(-1);
252
286k
    // fragment_mgr is nullptr in unittest
253
    if (ExecEnv::GetInstance()->fragment_mgr()) {
254
286k
        ExecEnv::GetInstance()->fragment_mgr()->remove_query_context(this->_query_id);
255
286k
    }
256
    // the only one msg shows query's end. any other msg should append to it if need.
257
123k
    LOG_INFO("Query {} deconstructed, mem_tracker: {}", print_id(this->_query_id), mem_tracker_msg);
258
123k
}
259
123k
260
123k
void QueryContext::set_ready_to_execute(Status reason) {
261
    set_execution_dependency_ready();
262
167k
    _exec_status.update(reason);
263
167k
}
264
167k
265
void QueryContext::set_ready_to_execute_only() {
266
290k
    set_execution_dependency_ready();
267
290k
}
268
290k
269
void QueryContext::set_execution_dependency_ready() {
270
22
    _execution_dependency->set_ready();
271
22
}
272
10
273
10
void QueryContext::set_memory_sufficient(bool sufficient) {
274
10
    if (sufficient) {
275
10
        {
276
12
            _memory_sufficient_dependency->set_ready();
277
12
            _resource_ctx->task_controller()->reset_paused_reason();
278
12
        }
279
12
    } else {
280
22
        _memory_sufficient_dependency->block();
281
        _resource_ctx->task_controller()->add_paused_count();
282
116k
    }
283
116k
}
284
110k
285
110k
void QueryContext::cancel(Status new_status, int fragment_id) {
286
    if (!_exec_status.update(new_status)) {
287
5.98k
        return;
288
5.98k
    }
289
5.98k
    // Tasks should be always runnable.
290
5.98k
    _execution_dependency->set_always_ready();
291
5.98k
    _memory_sufficient_dependency->set_always_ready();
292
5.98k
    if ((new_status.is<ErrorCode::MEM_LIMIT_EXCEEDED>() ||
293
         new_status.is<ErrorCode::MEM_ALLOC_FAILED>()) &&
294
        _query_options.__isset.dump_heap_profile_when_mem_limit_exceeded &&
295
0
        _query_options.dump_heap_profile_when_mem_limit_exceeded) {
296
0
        // if query is cancelled because of query mem limit exceeded, dump heap profile
297
0
        // at the time of cancellation can get the most accurate memory usage for problem analysis
298
0
        auto wg = workload_group();
299
0
        auto log_str = fmt::format(
300
0
                "Query {} canceled because of memory limit exceeded, dumping memory "
301
0
                "detail profiles. wg: {}. {}",
302
0
                print_id(_query_id), wg ? wg->debug_string() : "null",
303
0
                doris::ProcessProfile::instance()->memory_profile()->process_memory_detail_str());
304
0
        LOG_LONG_STRING(INFO, log_str);
305
0
        std::string dot = HeapProfiler::instance()->dump_heap_profile_to_dot();
306
0
        if (!dot.empty()) {
307
0
            dot += "\n-------------------------------------------------------\n";
308
0
            dot += "Copy the text after `digraph` in the above output to "
309
0
                   "http://www.webgraphviz.com to generate a dot graph.\n"
310
0
                   "after start heap profiler, if there is no operation, will print `No nodes "
311
0
                   "to "
312
0
                   "print`."
313
0
                   "If there are many errors: `addr2line: Dwarf Error`,"
314
0
                   "or other FAQ, reference doc: "
315
0
                   "https://doris.apache.org/community/developer-guide/debug-tool/#4-qa\n";
316
0
            auto nest_log_str =
317
0
                    fmt::format("Query {}, dump heap profile to dot: {}", print_id(_query_id), dot);
318
            LOG_LONG_STRING(INFO, nest_log_str);
319
5.98k
        }
320
5.98k
    }
321
5.98k
322
    set_ready_to_execute(new_status);
323
27
    cancel_all_pipeline_context(new_status, fragment_id);
324
27
}
325
27
326
27
void QueryContext::set_load_error_url(std::string error_url) {
327
    std::lock_guard<std::mutex> lock(_error_url_lock);
328
47.3k
    _load_error_url = error_url;
329
47.3k
}
330
47.3k
331
47.3k
std::string QueryContext::get_load_error_url() {
332
    std::lock_guard<std::mutex> lock(_error_url_lock);
333
27
    return _load_error_url;
334
27
}
335
27
336
27
void QueryContext::set_first_error_msg(std::string error_msg) {
337
    std::lock_guard<std::mutex> lock(_error_url_lock);
338
47.3k
    _first_error_msg = error_msg;
339
47.3k
}
340
47.3k
341
47.3k
std::string QueryContext::get_first_error_msg() {
342
    std::lock_guard<std::mutex> lock(_error_url_lock);
343
5.97k
    return _first_error_msg;
344
5.97k
}
345
5.97k
346
5.97k
void QueryContext::cancel_all_pipeline_context(const Status& reason, int fragment_id) {
347
8.82k
    std::vector<std::weak_ptr<PipelineFragmentContext>> ctx_to_cancel;
348
8.82k
    {
349
1.36k
        std::lock_guard<std::mutex> lock(_pipeline_map_write_lock);
350
1.36k
        for (auto& [f_id, f_context] : _fragment_id_to_pipeline_ctx) {
351
7.45k
            if (fragment_id == f_id) {
352
7.45k
                continue;
353
5.97k
            }
354
7.45k
            ctx_to_cancel.push_back(f_context);
355
7.45k
        }
356
759
    }
357
759
    for (auto& f_context : ctx_to_cancel) {
358
7.45k
        if (auto pipeline_ctx = f_context.lock()) {
359
5.97k
            pipeline_ctx->cancel(reason);
360
        }
361
0
    }
362
0
}
363
0
364
0
std::string QueryContext::print_all_pipeline_context() {
365
0
    std::vector<std::weak_ptr<PipelineFragmentContext>> ctx_to_print;
366
0
    fmt::memory_buffer debug_string_buffer;
367
0
    size_t i = 0;
368
    {
369
0
        fmt::format_to(debug_string_buffer, "{} pipeline fragment contexts in query {}. \n",
370
0
                       _fragment_id_to_pipeline_ctx.size(), print_id(_query_id));
371
0
372
0
        {
373
0
            std::lock_guard<std::mutex> lock(_pipeline_map_write_lock);
374
0
            for (auto& [f_id, f_context] : _fragment_id_to_pipeline_ctx) {
375
0
                ctx_to_print.push_back(f_context);
376
0
            }
377
0
        }
378
0
        for (auto& f_context : ctx_to_print) {
379
0
            if (auto pipeline_ctx = f_context.lock()) {
380
0
                auto elapsed = pipeline_ctx->elapsed_time() / 1000000000.0;
381
0
                fmt::format_to(debug_string_buffer,
382
0
                               "No.{} (elapse_second={}s, fragment_id={}) : {}\n", i, elapsed,
383
0
                               pipeline_ctx->get_fragment_id(), pipeline_ctx->debug_string());
384
0
                i++;
385
0
            }
386
0
        }
387
    }
388
    return fmt::to_string(debug_string_buffer);
389
430k
}
390
430k
391
void QueryContext::set_pipeline_context(const int fragment_id,
392
                                        std::shared_ptr<PipelineFragmentContext> pip_ctx) {
393
430k
    std::lock_guard<std::mutex> lock(_pipeline_map_write_lock);
394
430k
    // Use insert_or_assign instead of insert to support overwriting old entries
395
    // when recursive CTE recreates PipelineFragmentContext between rounds.
396
5.43M
    _fragment_id_to_pipeline_ctx.insert_or_assign(fragment_id, pip_ctx);
397
5.43M
}
398
0
399
0
doris::TaskScheduler* QueryContext::get_pipe_exec_scheduler() {
400
5.43M
    if (!_task_scheduler) {
401
5.43M
        throw Exception(Status::InternalError("task_scheduler is null"));
402
    }
403
285k
    return _task_scheduler;
404
285k
}
405
406
Status QueryContext::set_workload_group(WorkloadGroupPtr& wg) {
407
    _resource_ctx->set_workload_group(wg);
408
285k
    // Should add query first, the workload group will not be deleted,
409
    // then visit workload group's resource
410
285k
    // see task_group_manager::delete_workload_group_by_ids
411
285k
    RETURN_IF_ERROR(workload_group()->add_resource_ctx(_query_id, _resource_ctx));
412
285k
413
285k
    workload_group()->get_query_scheduler(&_task_scheduler, &_scan_task_scheduler,
414
                                          &_remote_scan_task_scheduler);
415
    return Status::OK();
416
}
417
2.84k
418
2.84k
void QueryContext::add_fragment_profile(
419
0
        int fragment_id, const std::vector<std::shared_ptr<TRuntimeProfileTree>>& pipeline_profiles,
420
0
        std::shared_ptr<TRuntimeProfileTree> load_channel_profile) {
421
0
    if (pipeline_profiles.empty()) {
422
0
        std::string msg = fmt::format("Add pipeline profile failed, query {}, fragment {}",
423
0
                                      print_id(this->_query_id), fragment_id);
424
0
        LOG_ERROR(msg);
425
        DCHECK(false) << msg;
426
2.84k
        return;
427
8.17k
    }
428
8.17k
429
0
#ifndef NDEBUG
430
8.17k
    for (const auto& p : pipeline_profiles) {
431
2.84k
        DCHECK(p != nullptr) << fmt::format("Add pipeline profile failed, query {}, fragment {}",
432
                                            print_id(this->_query_id), fragment_id);
433
2.84k
    }
434
2.84k
#endif
435
0
436
0
    std::lock_guard<std::mutex> l(_profile_mutex);
437
    VLOG_ROW << fmt::format(
438
2.84k
            "Query add fragment profile, query {}, fragment {}, pipeline profile count {} ",
439
            print_id(this->_query_id), fragment_id, pipeline_profiles.size());
440
2.84k
441
2.84k
    _profile_map.insert(std::make_pair(fragment_id, pipeline_profiles));
442
2.84k
443
2.84k
    if (load_channel_profile != nullptr) {
444
        _load_channel_profile_map.insert(std::make_pair(fragment_id, load_channel_profile));
445
1.45k
    }
446
1.45k
}
447
448
2.78k
void QueryContext::_report_query_profile() {
449
2.78k
    std::lock_guard<std::mutex> lg(_profile_mutex);
450
451
2.78k
    for (auto& [fragment_id, fragment_profile] : _profile_map) {
452
2.78k
        std::shared_ptr<TRuntimeProfileTree> load_channel_profile = nullptr;
453
2.78k
454
        if (_load_channel_profile_map.contains(fragment_id)) {
455
2.78k
            load_channel_profile = _load_channel_profile_map[fragment_id];
456
2.78k
        }
457
2.78k
458
        ExecEnv::GetInstance()->runtime_query_statistics_mgr()->register_fragment_profile(
459
1.45k
                _query_id, this->coord_addr, fragment_id, fragment_profile, load_channel_profile);
460
1.45k
    }
461
462
    ExecEnv::GetInstance()->runtime_query_statistics_mgr()->trigger_profile_reporting();
463
0
}
464
0
465
0
std::unordered_map<int, std::vector<std::shared_ptr<TRuntimeProfileTree>>>
466
0
QueryContext::_collect_realtime_query_profile() {
467
0
    std::unordered_map<int, std::vector<std::shared_ptr<TRuntimeProfileTree>>> res;
468
0
    std::lock_guard<std::mutex> lock(_pipeline_map_write_lock);
469
0
    for (const auto& [fragment_id, fragment_ctx_wptr] : _fragment_id_to_pipeline_ctx) {
470
0
        if (auto fragment_ctx = fragment_ctx_wptr.lock()) {
471
0
            if (fragment_ctx == nullptr) {
472
0
                std::string msg =
473
0
                        fmt::format("PipelineFragmentContext is nullptr, query {} fragment_id: {}",
474
0
                                    print_id(_query_id), fragment_id);
475
0
                LOG_ERROR(msg);
476
                DCHECK(false) << msg;
477
0
                continue;
478
            }
479
0
480
0
            auto profile = fragment_ctx->collect_realtime_profile();
481
0
482
0
            if (profile.empty()) {
483
0
                std::string err_msg = fmt::format(
484
0
                        "Get nothing when collecting profile, query {}, fragment_id: {}",
485
0
                        print_id(_query_id), fragment_id);
486
0
                LOG_ERROR(err_msg);
487
                DCHECK(false) << err_msg;
488
0
                continue;
489
0
            }
490
0
491
            res.insert(std::make_pair(fragment_id, profile));
492
0
        }
493
0
    }
494
495
0
    return res;
496
0
}
497
498
0
TReportExecStatusParams QueryContext::get_realtime_exec_status() {
499
0
    TReportExecStatusParams exec_status;
500
501
0
    auto realtime_query_profile = _collect_realtime_query_profile();
502
0
    std::vector<std::shared_ptr<TRuntimeProfileTree>> load_channel_profiles;
503
0
504
0
    for (auto load_channel_profile : _load_channel_profile_map) {
505
0
        if (load_channel_profile.second != nullptr) {
506
            load_channel_profiles.push_back(load_channel_profile.second);
507
0
        }
508
0
    }
509
0
510
    exec_status = RuntimeQueryStatisticsMgr::create_report_exec_status_params(
511
0
            this->_query_id, std::move(realtime_query_profile), std::move(load_channel_profiles),
512
0
            /*is_done=*/false);
513
514
    return exec_status;
515
}
516
3.59k
517
3.59k
Status QueryContext::send_block_to_cte_scan(
518
3.59k
        const TUniqueId& instance_id, int node_id,
519
3.59k
        const google::protobuf::RepeatedPtrField<doris::PBlock>& pblocks, bool eos) {
520
0
    std::unique_lock<std::mutex> l(_cte_scan_lock);
521
0
    auto it = _cte_scan.find(std::make_pair(instance_id, node_id));
522
0
    if (it == _cte_scan.end()) {
523
3.59k
        return Status::InternalError("RecCTEScan not found for instance {}, node {}",
524
1.96k
                                     print_id(instance_id), node_id);
525
1.96k
    }
526
3.59k
    for (const auto& pblock : pblocks) {
527
1.85k
        RETURN_IF_ERROR(it->second->add_block(pblock));
528
1.85k
    }
529
3.59k
    if (eos) {
530
3.59k
        it->second->set_ready();
531
    }
532
    return Status::OK();
533
1.85k
}
534
1.85k
535
1.85k
void QueryContext::registe_cte_scan(const TUniqueId& instance_id, int node_id,
536
1.85k
                                    RecCTEScanLocalState* scan) {
537
0
    std::unique_lock<std::mutex> l(_cte_scan_lock);
538
1.85k
    auto key = std::make_pair(instance_id, node_id);
539
1.85k
    DCHECK(!_cte_scan.contains(key)) << "Duplicate registe cte scan for instance "
540
                                     << print_id(instance_id) << ", node " << node_id;
541
1.85k
    _cte_scan.emplace(key, scan);
542
1.85k
}
543
1.85k
544
1.85k
void QueryContext::deregiste_cte_scan(const TUniqueId& instance_id, int node_id) {
545
0
    std::lock_guard<std::mutex> l(_cte_scan_lock);
546
1.85k
    auto key = std::make_pair(instance_id, node_id);
547
1.85k
    DCHECK(_cte_scan.contains(key)) << "Duplicate deregiste cte scan for instance "
548
                                    << print_id(instance_id) << ", node " << node_id;
549
1.74k
    _cte_scan.erase(key);
550
1.74k
}
551
1.74k
552
1.74k
Status QueryContext::reset_global_rf(const google::protobuf::RepeatedField<int32_t>& filter_ids) {
553
0
    if (_merge_controller_handler) {
554
1.74k
        return _merge_controller_handler->reset_global_rf(this, filter_ids);
555
    }
556
    return Status::OK();
557
}
558
559
} // namespace doris