Coverage Report

Created: 2026-08-06 13:04

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/pipeline/pipeline_fragment_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 "exec/pipeline/pipeline_fragment_context.h"
19
20
#include <gen_cpp/DataSinks_types.h>
21
#include <gen_cpp/FrontendService.h>
22
#include <gen_cpp/FrontendService_types.h>
23
#include <gen_cpp/PaloInternalService_types.h>
24
#include <gen_cpp/PlanNodes_types.h>
25
#include <pthread.h>
26
27
#include <algorithm>
28
#include <cstdlib>
29
// IWYU pragma: no_include <bits/chrono.h>
30
#include <fmt/format.h>
31
#include <thrift/Thrift.h>
32
#include <thrift/protocol/TDebugProtocol.h>
33
#include <thrift/transport/TTransportException.h>
34
35
#include <chrono> // IWYU pragma: keep
36
#include <map>
37
#include <memory>
38
#include <ostream>
39
#include <utility>
40
41
#include "cloud/config.h"
42
#include "common/cast_set.h"
43
#include "common/config.h"
44
#include "common/exception.h"
45
#include "common/logging.h"
46
#include "common/status.h"
47
#include "exec/exchange/local_exchange_sink_operator.h"
48
#include "exec/exchange/local_exchange_source_operator.h"
49
#include "exec/exchange/local_exchanger.h"
50
#include "exec/exchange/vdata_stream_mgr.h"
51
#include "exec/operator/aggregation_sink_operator.h"
52
#include "exec/operator/aggregation_source_operator.h"
53
#include "exec/operator/analytic_sink_operator.h"
54
#include "exec/operator/analytic_source_operator.h"
55
#include "exec/operator/assert_num_rows_operator.h"
56
#include "exec/operator/blackhole_sink_operator.h"
57
#include "exec/operator/bucketed_aggregation_sink_operator.h"
58
#include "exec/operator/bucketed_aggregation_source_operator.h"
59
#include "exec/operator/cache_sink_operator.h"
60
#include "exec/operator/cache_source_operator.h"
61
#include "exec/operator/datagen_operator.h"
62
#include "exec/operator/dict_sink_operator.h"
63
#include "exec/operator/distinct_streaming_aggregation_operator.h"
64
#include "exec/operator/empty_set_operator.h"
65
#include "exec/operator/exchange_sink_operator.h"
66
#include "exec/operator/exchange_source_operator.h"
67
#include "exec/operator/file_scan_operator.h"
68
#include "exec/operator/group_commit_block_sink_operator.h"
69
#include "exec/operator/group_commit_scan_operator.h"
70
#include "exec/operator/hashjoin_build_sink.h"
71
#include "exec/operator/hashjoin_probe_operator.h"
72
#include "exec/operator/hive_table_sink_operator.h"
73
#include "exec/operator/iceberg_delete_sink_operator.h"
74
#include "exec/operator/iceberg_merge_sink_operator.h"
75
#include "exec/operator/iceberg_table_sink_operator.h"
76
#include "exec/operator/jdbc_scan_operator.h"
77
#include "exec/operator/jdbc_table_sink_operator.h"
78
#include "exec/operator/local_merge_sort_source_operator.h"
79
#include "exec/operator/materialization_opertor.h"
80
#include "exec/operator/maxcompute_table_sink_operator.h"
81
#include "exec/operator/memory_scratch_sink_operator.h"
82
#include "exec/operator/meta_scan_operator.h"
83
#include "exec/operator/multi_cast_data_stream_sink.h"
84
#include "exec/operator/multi_cast_data_stream_source.h"
85
#include "exec/operator/nested_loop_join_build_operator.h"
86
#include "exec/operator/nested_loop_join_probe_operator.h"
87
#include "exec/operator/olap_scan_operator.h"
88
#include "exec/operator/olap_table_sink_operator.h"
89
#include "exec/operator/olap_table_sink_v2_operator.h"
90
#include "exec/operator/partition_sort_sink_operator.h"
91
#include "exec/operator/partition_sort_source_operator.h"
92
#include "exec/operator/partitioned_aggregation_sink_operator.h"
93
#include "exec/operator/partitioned_aggregation_source_operator.h"
94
#include "exec/operator/partitioned_hash_join_probe_operator.h"
95
#include "exec/operator/partitioned_hash_join_sink_operator.h"
96
#include "exec/operator/rec_cte_anchor_sink_operator.h"
97
#include "exec/operator/rec_cte_scan_operator.h"
98
#include "exec/operator/rec_cte_sink_operator.h"
99
#include "exec/operator/rec_cte_source_operator.h"
100
#include "exec/operator/repeat_operator.h"
101
#include "exec/operator/result_file_sink_operator.h"
102
#include "exec/operator/result_sink_operator.h"
103
#include "exec/operator/schema_scan_operator.h"
104
#include "exec/operator/select_operator.h"
105
#include "exec/operator/set_probe_sink_operator.h"
106
#include "exec/operator/set_sink_operator.h"
107
#include "exec/operator/set_source_operator.h"
108
#include "exec/operator/sort_sink_operator.h"
109
#include "exec/operator/sort_source_operator.h"
110
#include "exec/operator/spill_iceberg_table_sink_operator.h"
111
#include "exec/operator/spill_sort_sink_operator.h"
112
#include "exec/operator/spill_sort_source_operator.h"
113
#include "exec/operator/streaming_aggregation_operator.h"
114
#include "exec/operator/table_function_operator.h"
115
#include "exec/operator/tvf_table_sink_operator.h"
116
#include "exec/operator/union_sink_operator.h"
117
#include "exec/operator/union_source_operator.h"
118
#include "exec/pipeline/dependency.h"
119
#include "exec/pipeline/pipeline_task.h"
120
#include "exec/pipeline/task_scheduler.h"
121
#include "exec/runtime_filter/runtime_filter_mgr.h"
122
#include "exec/sort/topn_sorter.h"
123
#include "exec/spill/spill_file.h"
124
#include "io/fs/stream_load_pipe.h"
125
#include "load/stream_load/new_load_stream_mgr.h"
126
#include "runtime/cluster_info.h"
127
#include "runtime/exec_env.h"
128
#include "runtime/fragment_mgr.h"
129
#include "runtime/result_buffer_mgr.h"
130
#include "runtime/runtime_state.h"
131
#include "runtime/thread_context.h"
132
#include "service/backend_options.h"
133
#include "util/client_cache.h"
134
#include "util/countdown_latch.h"
135
#include "util/debug_util.h"
136
#include "util/network_util.h"
137
#include "util/uid_util.h"
138
139
namespace doris {
140
PipelineFragmentContext::PipelineFragmentContext(
141
        TUniqueId query_id, const TPipelineFragmentParams& request,
142
        std::shared_ptr<QueryContext> query_ctx, ExecEnv* exec_env,
143
        const std::function<void(RuntimeState*, Status*)>& call_back)
144
34
        : _query_id(std::move(query_id)),
145
34
          _fragment_id(request.fragment_id),
146
34
          _exec_env(exec_env),
147
34
          _query_ctx(std::move(query_ctx)),
148
34
          _call_back(call_back),
149
34
          _params(request),
150
34
          _parallel_instances(_params.__isset.parallel_instances ? _params.parallel_instances : 0),
151
34
          _need_notify_close(request.__isset.need_notify_close ? request.need_notify_close
152
34
                                                               : false) {
153
34
    _fragment_watcher.start();
154
34
}
155
156
34
PipelineFragmentContext::~PipelineFragmentContext() {
157
34
    LOG_INFO("PipelineFragmentContext::~PipelineFragmentContext")
158
34
            .tag("query_id", print_id(_query_id))
159
34
            .tag("fragment_id", _fragment_id);
160
34
    _release_resource();
161
34
    {
162
        // The memory released by the query end is recorded in the query mem tracker.
163
34
        SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(_query_ctx->query_mem_tracker());
164
34
        _runtime_state.reset();
165
34
        _query_ctx.reset();
166
34
    }
167
34
}
168
169
0
bool PipelineFragmentContext::is_timeout(timespec now) const {
170
0
    if (_timeout <= 0) {
171
0
        return false;
172
0
    }
173
0
    return _fragment_watcher.elapsed_time_seconds(now) > _timeout;
174
0
}
175
176
// notify_close() transitions the PFC from "waiting for external close notification" to
177
// "self-managed close". A recursive CTE PFC normally remains registered until rerun_fragment()
178
// calls this for WAIT_FOR_DESTROY or FINAL_CLOSE; cancellation can also call it.
179
// Returns true if all tasks have already closed (i.e., the PFC can be safely destroyed).
180
3
bool PipelineFragmentContext::notify_close() {
181
3
    bool all_closed = false;
182
3
    bool need_remove = false;
183
3
    {
184
3
        std::lock_guard<std::mutex> l(_task_mutex);
185
3
        if (_closed_tasks >= _total_tasks) {
186
2
            if (_need_notify_close) {
187
                // The fragment finished while waiting for the external close notification.
188
                // Record that we need to remove from fragment mgr, but do it
189
                // after releasing _task_mutex to avoid ABBA deadlock with
190
                // dump_pipeline_tasks() (which acquires _pipeline_map lock
191
                // first, then _task_mutex via debug_string()).
192
2
                need_remove = true;
193
2
            }
194
2
            all_closed = true;
195
2
        }
196
        // Allow the fragment to be removed now or after its remaining tasks close.
197
3
        _need_notify_close = false;
198
3
    }
199
3
    if (need_remove) {
200
2
        _exec_env->fragment_mgr()->remove_pipeline_context({_query_id, _fragment_id});
201
2
    }
202
3
    return all_closed;
203
3
}
204
205
// Must not add lock in this method. Because it will call query ctx cancel. And
206
// QueryCtx cancel will call fragment ctx cancel. And Also Fragment ctx's running
207
// Method like exchange sink buffer will call query ctx cancel. If we add lock here
208
// There maybe dead lock.
209
1
void PipelineFragmentContext::cancel(const Status reason) {
210
1
    LOG_INFO("PipelineFragmentContext::cancel")
211
1
            .tag("query_id", print_id(_query_id))
212
1
            .tag("fragment_id", _fragment_id)
213
1
            .tag("reason", reason.to_string());
214
1
    if (notify_close()) {
215
0
        return;
216
0
    }
217
    // Timeout is a special error code, we need print current stack to debug timeout issue.
218
1
    if (reason.is<ErrorCode::TIMEOUT>()) {
219
0
        auto dbg_str = fmt::format("PipelineFragmentContext is cancelled due to timeout:\n{}",
220
0
                                   debug_string());
221
0
        LOG_LONG_STRING(WARNING, dbg_str);
222
0
    }
223
224
    // `ILLEGAL_STATE` means queries this fragment belongs to was not found in FE (maybe finished)
225
1
    if (reason.is<ErrorCode::ILLEGAL_STATE>()) {
226
0
        LOG_WARNING("PipelineFragmentContext is cancelled due to illegal state : {}",
227
0
                    debug_string());
228
0
    }
229
230
1
    if (reason.is<ErrorCode::MEM_LIMIT_EXCEEDED>() || reason.is<ErrorCode::MEM_ALLOC_FAILED>()) {
231
0
        print_profile("cancel pipeline, reason: " + reason.to_string());
232
0
    }
233
234
1
    if (auto error_url = get_load_error_url(); !error_url.empty()) {
235
0
        _query_ctx->set_load_error_url(error_url);
236
0
    }
237
238
1
    if (auto first_error_msg = get_first_error_msg(); !first_error_msg.empty()) {
239
0
        _query_ctx->set_first_error_msg(first_error_msg);
240
0
    }
241
242
1
    _query_ctx->cancel(reason, _fragment_id);
243
1
    if (!reason.is<ErrorCode::LIMIT_REACH>() && !reason.is<ErrorCode::FINISHED>()) {
244
1
        for (auto& id : _fragment_instance_ids) {
245
0
            LOG(WARNING) << "PipelineFragmentContext cancel instance: " << print_id(id);
246
0
        }
247
1
    }
248
    // Get pipe from new load stream manager and send cancel to it or the fragment may hang to wait read from pipe
249
    // For stream load the fragment's query_id == load id, it is set in FE.
250
1
    auto stream_load_ctx = _exec_env->new_load_stream_mgr()->get(_query_id);
251
1
    if (stream_load_ctx != nullptr) {
252
0
        stream_load_ctx->pipe->cancel(reason.to_string());
253
        // Set error URL here because after pipe is cancelled, stream load execution may return early.
254
        // We need to set the error URL at this point to ensure error information is properly
255
        // propagated to the client.
256
0
        stream_load_ctx->error_url = get_load_error_url();
257
0
        stream_load_ctx->first_error_msg = get_first_error_msg();
258
0
    }
259
260
1
    for (auto& tasks : _tasks) {
261
0
        for (auto& task : tasks) {
262
0
            task.first->unblock_all_dependencies();
263
0
        }
264
0
    }
265
1
}
266
267
0
PipelinePtr PipelineFragmentContext::add_pipeline(PipelinePtr parent, int idx) {
268
0
    PipelineId id = _next_pipeline_id++;
269
0
    auto pipeline = std::make_shared<Pipeline>(
270
0
            id, parent ? std::min(parent->num_tasks(), _num_instances) : _num_instances,
271
0
            parent ? parent->num_tasks() : _num_instances);
272
0
    if (idx >= 0) {
273
0
        _pipelines.insert(_pipelines.begin() + idx, pipeline);
274
0
    } else {
275
0
        _pipelines.emplace_back(pipeline);
276
0
    }
277
0
    if (parent) {
278
0
        parent->set_children(pipeline);
279
0
    }
280
0
    return pipeline;
281
0
}
282
283
0
Status PipelineFragmentContext::_build_and_prepare_full_pipeline(ThreadPool* thread_pool) {
284
0
    {
285
0
        SCOPED_TIMER(_build_pipelines_timer);
286
        // 2. Build pipelines with operators in this fragment.
287
0
        auto root_pipeline = add_pipeline();
288
0
        RETURN_IF_ERROR(_build_pipelines(_runtime_state->obj_pool(), *_query_ctx->desc_tbl,
289
0
                                         &_root_op, root_pipeline));
290
291
        // Propagate _num_instances from LOCAL_EXCHANGE pipelines to ancestor pipelines
292
        // that inherited reduced num_tasks from a serial operator.
293
0
        _propagate_local_exchange_num_tasks();
294
295
        // Create deferred local exchangers now that all pipelines have final num_tasks.
296
0
        RETURN_IF_ERROR(_create_deferred_local_exchangers());
297
298
        // Raise num_tasks for pipelines whose serial non-scan operators (e.g.,
299
        // UNPARTITIONED Exchange) reduced num_tasks below _num_instances.
300
        // Without this, fragment instances 1+ have no task for these pipelines
301
        // and downstream operators fail with "must set shared state".
302
        //
303
        // This applies to ALL pipelines (not just deferred exchanger upstreams):
304
        // fragments with UNION/INTERSECT/EXCEPT + serial Exchange in child
305
        // pipelines also need the raise, even without FE-planned local exchange.
306
        //
307
        // Exception: serial scan sources (pooling scan) keep num_tasks=1 — the
308
        // PassthroughExchanger(1, N) handles the fan-out correctly.
309
        // NOTE: Do NOT raise pipelines whose source is a serial operator
310
        // (Exchange or scan) — they legitimately have 1 task, and raising
311
        // them causes crashes (e.g., 4 Exchange tasks but only 1 receives
312
        // data).  The correct fix for shared state injection across
313
        // instances is handled by the FE: it inserts local exchange nodes
314
        // between serial operators and their downstream consumers, creating
315
        // proper pipeline boundaries with _num_instances tasks.
316
317
        // 3. Create sink operator
318
0
        if (!_params.fragment.__isset.output_sink) {
319
0
            return Status::InternalError("No output sink in this fragment!");
320
0
        }
321
0
        RETURN_IF_ERROR(_create_data_sink(_runtime_state->obj_pool(), _params.fragment.output_sink,
322
0
                                          _params.fragment.output_exprs, _params,
323
0
                                          root_pipeline->output_row_desc(), _runtime_state.get(),
324
0
                                          *_desc_tbl, root_pipeline->id()));
325
0
        RETURN_IF_ERROR(_sink->init(_params.fragment.output_sink));
326
0
        RETURN_IF_ERROR(root_pipeline->set_sink(_sink));
327
328
0
        for (PipelinePtr& pipeline : _pipelines) {
329
0
            DCHECK(pipeline->sink() != nullptr) << pipeline->operators().size();
330
0
            RETURN_IF_ERROR(pipeline->sink()->set_child(pipeline->operators().back()));
331
0
        }
332
0
    }
333
    // 4. Build local exchanger
334
0
    if (_runtime_state->plan_local_shuffle()) {
335
0
        SCOPED_TIMER(_plan_local_exchanger_timer);
336
0
        RETURN_IF_ERROR(_plan_local_exchange(_params.num_buckets,
337
0
                                             _params.bucket_seq_to_instance_idx,
338
0
                                             _params.shuffle_idx_to_instance_idx));
339
0
    }
340
341
    // 5. Initialize global states in pipelines.
342
0
    for (PipelinePtr& pipeline : _pipelines) {
343
0
        SCOPED_TIMER(_prepare_all_pipelines_timer);
344
0
        pipeline->children().clear();
345
0
        RETURN_IF_ERROR(pipeline->prepare(_runtime_state.get()));
346
0
    }
347
348
0
    {
349
0
        SCOPED_TIMER(_build_tasks_timer);
350
        // 6. Build pipeline tasks and initialize local state.
351
0
        RETURN_IF_ERROR(_build_pipeline_tasks(thread_pool));
352
0
    }
353
354
0
    return Status::OK();
355
0
}
356
357
0
Status PipelineFragmentContext::prepare(ThreadPool* thread_pool) {
358
0
    if (_prepared) {
359
0
        return Status::InternalError("Already prepared");
360
0
    }
361
0
    if (_params.__isset.query_options && _params.query_options.__isset.execution_timeout) {
362
0
        _timeout = _params.query_options.execution_timeout;
363
0
    }
364
365
0
    _fragment_level_profile = std::make_unique<RuntimeProfile>("PipelineContext");
366
0
    _prepare_timer = ADD_TIMER(_fragment_level_profile, "PrepareTime");
367
0
    SCOPED_TIMER(_prepare_timer);
368
0
    _build_pipelines_timer = ADD_TIMER(_fragment_level_profile, "BuildPipelinesTime");
369
0
    _init_context_timer = ADD_TIMER(_fragment_level_profile, "InitContextTime");
370
0
    _plan_local_exchanger_timer = ADD_TIMER(_fragment_level_profile, "PlanLocalLocalExchangerTime");
371
0
    _build_tasks_timer = ADD_TIMER(_fragment_level_profile, "BuildTasksTime");
372
0
    _prepare_all_pipelines_timer = ADD_TIMER(_fragment_level_profile, "PrepareAllPipelinesTime");
373
0
    {
374
0
        SCOPED_TIMER(_init_context_timer);
375
0
        cast_set(_num_instances, _params.local_params.size());
376
0
        _total_instances =
377
0
                _params.__isset.total_instances ? _params.total_instances : _num_instances;
378
379
0
        auto* fragment_context = this;
380
381
0
        if (_params.query_options.__isset.is_report_success) {
382
0
            fragment_context->set_is_report_success(_params.query_options.is_report_success);
383
0
        }
384
385
        // 1. Set up the global runtime state.
386
0
        _runtime_state = RuntimeState::create_unique(
387
0
                _params.query_id, _params.fragment_id, _params.query_options,
388
0
                _query_ctx->query_globals, _exec_env, _query_ctx.get());
389
0
        _runtime_state->set_task_execution_context(shared_from_this());
390
0
        SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(_runtime_state->query_mem_tracker());
391
0
        if (_params.__isset.backend_id) {
392
0
            _runtime_state->set_backend_id(_params.backend_id);
393
0
        }
394
0
        if (_params.__isset.import_label) {
395
0
            _runtime_state->set_import_label(_params.import_label);
396
0
        }
397
0
        if (_params.__isset.db_name) {
398
0
            _runtime_state->set_db_name(_params.db_name);
399
0
        }
400
0
        if (_params.__isset.load_job_id) {
401
0
            _runtime_state->set_load_job_id(_params.load_job_id);
402
0
        }
403
404
0
        if (_params.is_simplified_param) {
405
0
            _desc_tbl = _query_ctx->desc_tbl;
406
0
        } else {
407
0
            DCHECK(_params.__isset.desc_tbl);
408
0
            RETURN_IF_ERROR(DescriptorTbl::create(_runtime_state->obj_pool(), _params.desc_tbl,
409
0
                                                  &_desc_tbl));
410
0
        }
411
0
        _runtime_state->set_desc_tbl(_desc_tbl);
412
0
        _runtime_state->set_num_per_fragment_instances(_params.num_senders);
413
0
        _runtime_state->set_load_stream_per_node(_params.load_stream_per_node);
414
0
        _runtime_state->set_total_load_streams(_params.total_load_streams);
415
0
        _runtime_state->set_num_local_sink(_params.num_local_sink);
416
417
        // init fragment_instance_ids
418
0
        const auto target_size = _params.local_params.size();
419
0
        _fragment_instance_ids.resize(target_size);
420
0
        for (size_t i = 0; i < _params.local_params.size(); i++) {
421
0
            auto fragment_instance_id = _params.local_params[i].fragment_instance_id;
422
0
            _fragment_instance_ids[i] = fragment_instance_id;
423
0
        }
424
0
    }
425
426
0
    RETURN_IF_ERROR(_build_and_prepare_full_pipeline(thread_pool));
427
428
0
    _init_next_report_time();
429
430
0
    _prepared = true;
431
0
    return Status::OK();
432
0
}
433
434
Status PipelineFragmentContext::_build_pipeline_tasks_for_instance(
435
        int instance_idx,
436
0
        const std::vector<std::shared_ptr<RuntimeProfile>>& pipeline_id_to_profile) {
437
0
    const auto& local_params = _params.local_params[instance_idx];
438
0
    auto fragment_instance_id = local_params.fragment_instance_id;
439
0
    auto runtime_filter_mgr = std::make_unique<RuntimeFilterMgr>(false);
440
0
    std::map<PipelineId, PipelineTask*> pipeline_id_to_task;
441
0
    auto get_shared_state = [&](PipelinePtr pipeline)
442
0
            -> std::map<int, std::pair<std::shared_ptr<BasicSharedState>,
443
0
                                       std::vector<std::shared_ptr<Dependency>>>> {
444
0
        std::map<int, std::pair<std::shared_ptr<BasicSharedState>,
445
0
                                std::vector<std::shared_ptr<Dependency>>>>
446
0
                shared_state_map;
447
0
        for (auto& op : pipeline->operators()) {
448
0
            auto source_id = op->operator_id();
449
0
            if (auto iter = _op_id_to_shared_state.find(source_id);
450
0
                iter != _op_id_to_shared_state.end()) {
451
0
                shared_state_map.insert({source_id, iter->second});
452
0
            }
453
0
        }
454
0
        for (auto sink_to_source_id : pipeline->sink()->dests_id()) {
455
0
            if (auto iter = _op_id_to_shared_state.find(sink_to_source_id);
456
0
                iter != _op_id_to_shared_state.end()) {
457
0
                shared_state_map.insert({sink_to_source_id, iter->second});
458
0
            }
459
0
        }
460
0
        return shared_state_map;
461
0
    };
462
463
0
    for (size_t pip_idx = 0; pip_idx < _pipelines.size(); pip_idx++) {
464
0
        auto& pipeline = _pipelines[pip_idx];
465
0
        if (pipeline->num_tasks() > 1 || instance_idx == 0) {
466
0
            auto task_runtime_state = RuntimeState::create_unique(
467
0
                    local_params.fragment_instance_id, _params.query_id, _params.fragment_id,
468
0
                    _params.query_options, _query_ctx->query_globals, _exec_env, _query_ctx.get());
469
0
            {
470
                // Initialize runtime state for this task
471
0
                task_runtime_state->set_query_mem_tracker(_query_ctx->query_mem_tracker());
472
473
0
                task_runtime_state->set_task_execution_context(shared_from_this());
474
0
                task_runtime_state->set_be_number(local_params.backend_num);
475
476
0
                if (_params.__isset.backend_id) {
477
0
                    task_runtime_state->set_backend_id(_params.backend_id);
478
0
                }
479
0
                if (_params.__isset.import_label) {
480
0
                    task_runtime_state->set_import_label(_params.import_label);
481
0
                }
482
0
                if (_params.__isset.db_name) {
483
0
                    task_runtime_state->set_db_name(_params.db_name);
484
0
                }
485
0
                if (_params.__isset.load_job_id) {
486
0
                    task_runtime_state->set_load_job_id(_params.load_job_id);
487
0
                }
488
0
                if (_params.__isset.wal_id) {
489
0
                    task_runtime_state->set_wal_id(_params.wal_id);
490
0
                }
491
0
                if (_params.__isset.content_length) {
492
0
                    task_runtime_state->set_content_length(_params.content_length);
493
0
                }
494
495
0
                task_runtime_state->set_desc_tbl(_desc_tbl);
496
0
                task_runtime_state->set_per_fragment_instance_idx(local_params.sender_id);
497
0
                task_runtime_state->set_num_per_fragment_instances(_params.num_senders);
498
0
                task_runtime_state->resize_op_id_to_local_state(max_operator_id());
499
0
                task_runtime_state->set_max_operator_id(max_operator_id());
500
0
                task_runtime_state->set_load_stream_per_node(_params.load_stream_per_node);
501
0
                task_runtime_state->set_total_load_streams(_params.total_load_streams);
502
0
                task_runtime_state->set_num_local_sink(_params.num_local_sink);
503
504
0
                task_runtime_state->set_runtime_filter_mgr(runtime_filter_mgr.get());
505
0
            }
506
0
            auto cur_task_id = _total_tasks++;
507
0
            task_runtime_state->set_task_id(cur_task_id);
508
0
            task_runtime_state->set_task_num(pipeline->num_tasks());
509
0
            auto task = std::make_shared<PipelineTask>(
510
0
                    pipeline, cur_task_id, task_runtime_state.get(),
511
0
                    std::dynamic_pointer_cast<PipelineFragmentContext>(shared_from_this()),
512
0
                    pipeline_id_to_profile[pip_idx].get(), get_shared_state(pipeline),
513
0
                    instance_idx);
514
0
            pipeline->incr_created_tasks(instance_idx, task.get());
515
0
            pipeline_id_to_task.insert({pipeline->id(), task.get()});
516
0
            _tasks[instance_idx].emplace_back(
517
0
                    std::pair<std::shared_ptr<PipelineTask>, std::unique_ptr<RuntimeState>> {
518
0
                            std::move(task), std::move(task_runtime_state)});
519
0
        }
520
0
    }
521
522
    /**
523
         * Build DAG for pipeline tasks.
524
         * For example, we have
525
         *
526
         *   ExchangeSink (Pipeline1)     JoinBuildSink (Pipeline2)
527
         *            \                      /
528
         *          JoinProbeOperator1 (Pipeline1)    JoinBuildSink (Pipeline3)
529
         *                 \                          /
530
         *               JoinProbeOperator2 (Pipeline1)
531
         *
532
         * In this fragment, we have three pipelines and pipeline 1 depends on pipeline 2 and pipeline 3.
533
         * To build this DAG, `_dag` manage dependencies between pipelines by pipeline ID and
534
         * `pipeline_id_to_task` is used to find the task by a unique pipeline ID.
535
         *
536
         * Finally, we have two upstream dependencies in Pipeline1 corresponding to JoinProbeOperator1
537
         * and JoinProbeOperator2.
538
         */
539
0
    for (auto& _pipeline : _pipelines) {
540
0
        if (pipeline_id_to_task.contains(_pipeline->id())) {
541
0
            auto* task = pipeline_id_to_task[_pipeline->id()];
542
0
            DCHECK(task != nullptr);
543
544
            // If this task has upstream dependency, then inject it into this task.
545
0
            if (_dag.contains(_pipeline->id())) {
546
0
                auto& deps = _dag[_pipeline->id()];
547
0
                for (auto& dep : deps) {
548
0
                    if (pipeline_id_to_task.contains(dep)) {
549
0
                        auto ss = pipeline_id_to_task[dep]->get_sink_shared_state();
550
0
                        if (ss) {
551
0
                            task->inject_shared_state(ss);
552
0
                        } else {
553
0
                            pipeline_id_to_task[dep]->inject_shared_state(
554
0
                                    task->get_source_shared_state());
555
0
                        }
556
0
                    }
557
0
                }
558
0
            }
559
0
        }
560
0
    }
561
0
    for (size_t pip_idx = 0; pip_idx < _pipelines.size(); pip_idx++) {
562
0
        if (pipeline_id_to_task.contains(_pipelines[pip_idx]->id())) {
563
0
            auto* task = pipeline_id_to_task[_pipelines[pip_idx]->id()];
564
0
            DCHECK(pipeline_id_to_profile[pip_idx]);
565
0
            std::vector<TScanRangeParams> scan_ranges;
566
0
            auto node_id = _pipelines[pip_idx]->operators().front()->node_id();
567
0
            if (local_params.per_node_scan_ranges.contains(node_id)) {
568
0
                scan_ranges = local_params.per_node_scan_ranges.find(node_id)->second;
569
0
            }
570
0
            RETURN_IF_ERROR_OR_CATCH_EXCEPTION(task->prepare(scan_ranges, local_params.sender_id,
571
0
                                                             _params.fragment.output_sink));
572
0
        }
573
0
    }
574
0
    {
575
0
        std::lock_guard<std::mutex> l(_state_map_lock);
576
0
        _runtime_filter_mgr_map[instance_idx] = std::move(runtime_filter_mgr);
577
0
    }
578
0
    return Status::OK();
579
0
}
580
581
0
Status PipelineFragmentContext::_build_pipeline_tasks(ThreadPool* thread_pool) {
582
0
    _total_tasks = 0;
583
0
    _closed_tasks = 0;
584
0
    const auto target_size = _params.local_params.size();
585
0
    _tasks.resize(target_size);
586
0
    _runtime_filter_mgr_map.resize(target_size);
587
0
    for (size_t pip_idx = 0; pip_idx < _pipelines.size(); pip_idx++) {
588
0
        _pip_id_to_pipeline[_pipelines[pip_idx]->id()] = _pipelines[pip_idx].get();
589
0
    }
590
0
    auto pipeline_id_to_profile = _runtime_state->build_pipeline_profile(_pipelines.size());
591
592
0
    if (target_size > 1 &&
593
0
        (_runtime_state->query_options().__isset.parallel_prepare_threshold &&
594
0
         target_size > _runtime_state->query_options().parallel_prepare_threshold)) {
595
        // If instances parallelism is big enough ( > parallel_prepare_threshold), we will prepare all tasks by multi-threads
596
0
        std::vector<Status> prepare_status(target_size);
597
0
        int submitted_tasks = 0;
598
0
        Status submit_status;
599
0
        CountDownLatch latch((int)target_size);
600
0
        for (int i = 0; i < target_size; i++) {
601
0
            submit_status = thread_pool->submit_func([&, i]() {
602
0
                SCOPED_ATTACH_TASK(_query_ctx.get());
603
0
                prepare_status[i] = _build_pipeline_tasks_for_instance(i, pipeline_id_to_profile);
604
0
                latch.count_down();
605
0
            });
606
0
            if (LIKELY(submit_status.ok())) {
607
0
                submitted_tasks++;
608
0
            } else {
609
0
                break;
610
0
            }
611
0
        }
612
0
        latch.arrive_and_wait(target_size - submitted_tasks);
613
0
        if (UNLIKELY(!submit_status.ok())) {
614
0
            return submit_status;
615
0
        }
616
0
        for (int i = 0; i < submitted_tasks; i++) {
617
0
            if (!prepare_status[i].ok()) {
618
0
                return prepare_status[i];
619
0
            }
620
0
        }
621
0
    } else {
622
0
        for (int i = 0; i < target_size; i++) {
623
0
            RETURN_IF_ERROR(_build_pipeline_tasks_for_instance(i, pipeline_id_to_profile));
624
0
        }
625
0
    }
626
0
    _pipeline_parent_map.clear();
627
0
    _op_id_to_shared_state.clear();
628
    // Record task cardinality once when this fragment context finishes task initialization.
629
0
    _query_ctx->add_total_task_num(_total_tasks.load(std::memory_order_relaxed));
630
631
0
    return Status::OK();
632
0
}
633
634
0
void PipelineFragmentContext::_init_next_report_time() {
635
0
    auto interval_s = config::pipeline_status_report_interval;
636
0
    if (_is_report_success && interval_s > 0 && _timeout > interval_s) {
637
0
        VLOG_FILE << "enable period report: fragment id=" << _fragment_id;
638
0
        uint64_t report_fragment_offset = (uint64_t)(rand() % interval_s) * NANOS_PER_SEC;
639
        // We don't want to wait longer than it takes to run the entire fragment.
640
0
        _previous_report_time =
641
0
                MonotonicNanos() + report_fragment_offset - (uint64_t)(interval_s)*NANOS_PER_SEC;
642
0
        _disable_period_report = false;
643
0
    }
644
0
}
645
646
0
void PipelineFragmentContext::refresh_next_report_time() {
647
0
    auto disable = _disable_period_report.load(std::memory_order_acquire);
648
0
    DCHECK(disable == true);
649
0
    _previous_report_time.store(MonotonicNanos(), std::memory_order_release);
650
0
    _disable_period_report.compare_exchange_strong(disable, false);
651
0
}
652
653
1
void PipelineFragmentContext::trigger_report_if_necessary() {
654
1
    if (!_is_report_success) {
655
1
        return;
656
1
    }
657
0
    auto disable = _disable_period_report.load(std::memory_order_acquire);
658
0
    if (disable) {
659
0
        return;
660
0
    }
661
0
    int32_t interval_s = config::pipeline_status_report_interval;
662
0
    if (interval_s <= 0) {
663
0
        LOG(WARNING) << "config::status_report_interval is equal to or less than zero, do not "
664
0
                        "trigger "
665
0
                        "report.";
666
0
    }
667
0
    uint64_t next_report_time = _previous_report_time.load(std::memory_order_acquire) +
668
0
                                (uint64_t)(interval_s)*NANOS_PER_SEC;
669
0
    if (MonotonicNanos() > next_report_time) {
670
0
        if (!_disable_period_report.compare_exchange_strong(disable, true,
671
0
                                                            std::memory_order_acq_rel)) {
672
0
            return;
673
0
        }
674
0
        if (VLOG_FILE_IS_ON) {
675
0
            VLOG_FILE << "Reporting "
676
0
                      << "profile for query_id " << print_id(_query_id)
677
0
                      << ", fragment id: " << _fragment_id;
678
679
0
            std::stringstream ss;
680
0
            _runtime_state->runtime_profile()->compute_time_in_profile();
681
0
            _runtime_state->runtime_profile()->pretty_print(&ss);
682
0
            if (_runtime_state->load_channel_profile()) {
683
0
                _runtime_state->load_channel_profile()->pretty_print(&ss);
684
0
            }
685
686
0
            VLOG_FILE << "Query " << print_id(get_query_id()) << " fragment " << get_fragment_id()
687
0
                      << " profile:\n"
688
0
                      << ss.str();
689
0
        }
690
0
        auto st = send_report(false);
691
0
        if (!st.ok()) {
692
0
            disable = true;
693
0
            _disable_period_report.compare_exchange_strong(disable, false,
694
0
                                                           std::memory_order_acq_rel);
695
0
        }
696
0
    }
697
0
}
698
699
Status PipelineFragmentContext::_build_pipelines(ObjectPool* pool, const DescriptorTbl& descs,
700
0
                                                 OperatorPtr* root, PipelinePtr cur_pipe) {
701
0
    if (_params.fragment.plan.nodes.empty()) {
702
0
        throw Exception(ErrorCode::INTERNAL_ERROR, "Invalid plan which has no plan node!");
703
0
    }
704
705
0
    int node_idx = 0;
706
707
0
    RETURN_IF_ERROR(_create_tree_helper(pool, _params.fragment.plan.nodes, descs, nullptr,
708
0
                                        &node_idx, root, cur_pipe, 0, false, false));
709
710
0
    if (node_idx + 1 != _params.fragment.plan.nodes.size()) {
711
0
        return Status::InternalError(
712
0
                "Plan tree only partially reconstructed. Not all thrift nodes were used.");
713
0
    }
714
0
    return Status::OK();
715
0
}
716
717
0
Status PipelineFragmentContext::_create_deferred_local_exchangers() {
718
0
    for (auto& info : _deferred_exchangers) {
719
        // DANGER ZONE — do not "fix" this line without reading the history.
720
        //
721
        // sender_count seeds Exchanger::_running_sink_operators, which the source side
722
        // waits to reach 0 via sub_running_sink_operators on each sink LocalState close.
723
        // The correct value is THIS pipeline-instance's sink task count, which is exactly
724
        // info.upstream_pipe->num_tasks() — one PipelineTask per task, one close per task.
725
        //
726
        // Tempting wrong fix #1: `std::max(num_tasks, _num_instances)` to mirror the
727
        //   BE-planned path in _add_local_exchange_impl (~line 1023).  THIS BREAKS the
728
        //   common FE-planned shape of `serial scan → LE(PT) → ...`: upstream_pipe
729
        //   genuinely has num_tasks=1, only 1 close arrives, but seed becomes
730
        //   _num_instances so _running_sink_operators never reaches 0 — downstream
731
        //   sources hang on SHUFFLE_DATA_DEPENDENCY (e.g. MTMV refresh from
732
        //   mtmv_up_down_job_p0/load.groovy stays at Status=RUNNING and regressed
733
        //   exactly this way).  BE-planned mode uses max() because its
734
        //   `cur_pipe` is the source-side pipeline (always raised to _num_instances by
735
        //   add_pipeline) — not analogous to our `upstream_pipe` here, which is the
736
        //   sink-side pipeline that may legitimately stay at 1 for serial sources.
737
        //
738
        // Tempting wrong fix #2: multiply by _num_instances on the theory shared_state
739
        //   is shared across all instances.  Same hang — each fragment-instance
740
        //   PipelineFragmentContext has its OWN _op_id_to_shared_state map, so the
741
        //   exchanger is per-instance, not per-BE.  num_tasks() is already the right
742
        //   close-count for one instance.
743
        //
744
        // If a hang shows up with `_running_sink_operators < 0`, the bug is upstream:
745
        // _propagate_local_exchange_num_tasks left num_tasks too low (or too high) for
746
        // this fragment shape.  Fix THAT pass, not this seed value.
747
0
        const int sender_count = info.upstream_pipe->num_tasks();
748
0
        switch (info.partition_type) {
749
0
        case TLocalPartitionType::LOCAL_EXECUTION_HASH_SHUFFLE:
750
0
        case TLocalPartitionType::GLOBAL_EXECUTION_HASH_SHUFFLE:
751
0
            info.shared_state->exchanger = ShuffleExchanger::create_unique(
752
0
                    sender_count, _num_instances, info.num_partitions, info.free_blocks_limit,
753
0
                    info.partition_type);
754
0
            break;
755
0
        case TLocalPartitionType::BUCKET_HASH_SHUFFLE:
756
0
            info.shared_state->exchanger = BucketShuffleExchanger::create_unique(
757
0
                    sender_count, _num_instances, info.num_partitions, info.free_blocks_limit);
758
0
            break;
759
0
        case TLocalPartitionType::PASSTHROUGH:
760
0
            info.shared_state->exchanger = PassthroughExchanger::create_unique(
761
0
                    sender_count, _num_instances, info.free_blocks_limit);
762
0
            break;
763
0
        case TLocalPartitionType::BROADCAST:
764
0
            info.shared_state->exchanger = BroadcastExchanger::create_unique(
765
0
                    sender_count, _num_instances, info.free_blocks_limit);
766
0
            break;
767
0
        case TLocalPartitionType::PASS_TO_ONE:
768
0
            if (_runtime_state->enable_share_hash_table_for_broadcast_join()) {
769
0
                info.shared_state->exchanger = PassToOneExchanger::create_unique(
770
0
                        sender_count, _num_instances, info.free_blocks_limit);
771
0
            } else {
772
0
                info.shared_state->exchanger = BroadcastExchanger::create_unique(
773
0
                        sender_count, _num_instances, info.free_blocks_limit);
774
0
            }
775
0
            break;
776
0
        case TLocalPartitionType::ADAPTIVE_PASSTHROUGH:
777
0
            info.shared_state->exchanger = AdaptivePassthroughExchanger::create_unique(
778
0
                    sender_count, _num_instances, info.free_blocks_limit);
779
0
            break;
780
0
        case TLocalPartitionType::NOOP:
781
0
        case TLocalPartitionType::LOCAL_MERGE_SORT:
782
            // FE-planned LocalExchangeNode currently never emits NOOP or LOCAL_MERGE_SORT
783
            // through the deferred-exchanger path.  NOOP means "no exchange needed" and
784
            // is filtered out before reaching here; LOCAL_MERGE_SORT is planned by the
785
            // legacy BE path only.  Crash in debug to surface the protocol violation if
786
            // that ever changes; return an error in release to avoid silently corrupting
787
            // execution.
788
0
            DCHECK(false) << "FE-planned local exchange should not emit partition_type="
789
0
                          << static_cast<int>(info.partition_type);
790
0
            return Status::InternalError("FE-planned local exchange emitted unsupported type: " +
791
0
                                         std::to_string(static_cast<int>(info.partition_type)));
792
0
        default:
793
            // New TLocalPartitionType added on FE side without a BE handler here.
794
0
            DCHECK(false) << "Unhandled TLocalPartitionType in deferred exchangers: "
795
0
                          << static_cast<int>(info.partition_type);
796
0
            return Status::InternalError("Unsupported FE-planned local exchange type: " +
797
0
                                         std::to_string(static_cast<int>(info.partition_type)));
798
0
        }
799
0
    }
800
0
    _deferred_exchangers.clear();
801
0
    return Status::OK();
802
0
}
803
804
0
void PipelineFragmentContext::_propagate_local_exchange_num_tasks() {
805
    // Only runs when FE has planned local exchanges and BE deferred their construction.
806
    // In legacy mode (enable_local_shuffle_planner=false) BE plans LE itself via
807
    // _plan_local_exchange and _deferred_exchangers stays empty — the legacy path
808
    // already gets its num_tasks right at construction time, so the propagate passes
809
    // would be no-ops and are skipped.  This is a transitional design: once the FE
810
    // planner is the only planner, the propagation logic itself should degrade into
811
    // a pure assertion that the FE plan already wired the right num_tasks everywhere.
812
0
    if (_deferred_exchangers.empty()) {
813
0
        return;
814
0
    }
815
    // Reconcile num_tasks across paired pipelines created by pipeline-splitting operators
816
    // (AGG, SORT, JOIN): they share state via inject_shared_state and must agree, or
817
    // instance 1+ tasks access null shared_state.  A pipeline's num_tasks is fully
818
    // determined by its source operator plus its upstreams:
819
    //   - LocalExchangeSource  -> _num_instances (the LE re-parallelizes)
820
    //   - serial source        -> its reduced count (kept as-is, typically 1)
821
    //   - otherwise (splitter) -> inherit from upstreams: raise to _num_instances if any
822
    //                             upstream was raised by an LE, then lower to a serial
823
    //                             upstream's count (lower wins).
824
    // Visiting each pipeline only after all its upstreams (topological order over _dag) lets
825
    // a single sweep reach the same fixpoint the previous two while-loops iterated to — those
826
    // only existed to reconcile the top-down build's parent-inherited num_tasks guesses.
827
0
    std::map<PipelineId, PipelinePtr> id_to_pipe;
828
0
    std::map<PipelineId, std::vector<PipelineId>> downstreams_of;
829
0
    std::map<PipelineId, int> in_degree;
830
0
    for (auto& p : _pipelines) {
831
0
        id_to_pipe[p->id()] = p;
832
0
        in_degree.try_emplace(p->id(), 0);
833
0
    }
834
0
    for (const auto& [downstream_id, upstream_ids] : _dag) {
835
0
        for (auto upstream_id : upstream_ids) {
836
0
            downstreams_of[upstream_id].push_back(downstream_id);
837
0
            in_degree[downstream_id]++;
838
0
        }
839
0
    }
840
0
    std::vector<PipelineId> ready;
841
0
    for (const auto& [id, deg] : in_degree) {
842
0
        if (deg == 0) {
843
0
            ready.push_back(id);
844
0
        }
845
0
    }
846
0
    size_t visited = 0;
847
0
    while (!ready.empty()) {
848
0
        const auto id = ready.back();
849
0
        ready.pop_back();
850
0
        visited++;
851
0
        auto pit = id_to_pipe.find(id);
852
0
        if (pit != id_to_pipe.end()) {
853
0
            auto& pipe = pit->second;
854
0
            const auto& ops = pipe->operators();
855
0
            const bool le_source =
856
0
                    !ops.empty() && dynamic_cast<LocalExchangeSourceOperatorX*>(ops.front().get());
857
0
            const bool serial_source = !ops.empty() && ops.front()->is_serial_operator();
858
0
            if (le_source) {
859
0
                pipe->set_num_tasks(_num_instances);
860
0
            } else if (!serial_source) {
861
0
                int target = pipe->num_tasks();
862
0
                const auto up_it = _dag.find(id);
863
0
                if (up_it != _dag.end()) {
864
                    // raise: any upstream already at _num_instances (e.g. an LE source)
865
0
                    for (auto upstream_id : up_it->second) {
866
0
                        auto uit = id_to_pipe.find(upstream_id);
867
0
                        if (uit != id_to_pipe.end() && uit->second->num_tasks() >= _num_instances) {
868
0
                            target = _num_instances;
869
0
                            break;
870
0
                        }
871
0
                    }
872
                    // lower: a serial upstream with fewer tasks (wins over the raise above)
873
0
                    for (auto upstream_id : up_it->second) {
874
0
                        auto uit = id_to_pipe.find(upstream_id);
875
0
                        if (uit != id_to_pipe.end() && uit->second->num_tasks() < target &&
876
0
                            !uit->second->operators().empty() &&
877
0
                            uit->second->operators().front()->is_serial_operator()) {
878
0
                            target = uit->second->num_tasks();
879
0
                        }
880
0
                    }
881
0
                }
882
0
                pipe->set_num_tasks(target);
883
0
            }
884
0
        }
885
0
        for (auto down : downstreams_of[id]) {
886
0
            if (--in_degree[down] == 0) {
887
0
                ready.push_back(down);
888
0
            }
889
0
        }
890
0
    }
891
    // The pipeline DAG is acyclic; if a future change introduces a back-edge, some pipelines
892
    // stay unvisited (in_degree never reaches 0) — fail loudly rather than silently leaving
893
    // their num_tasks unreconciled.
894
0
    DCHECK_EQ(visited, in_degree.size())
895
0
            << "pipeline num_tasks topological sweep visited " << visited << " of "
896
0
            << in_degree.size() << " pipelines (cycle in _dag?)";
897
0
}
898
899
Status PipelineFragmentContext::_create_tree_helper(
900
        ObjectPool* pool, const std::vector<TPlanNode>& tnodes, const DescriptorTbl& descs,
901
        OperatorPtr parent, int* node_idx, OperatorPtr* root, PipelinePtr& cur_pipe, int child_idx,
902
0
        const bool followed_by_shuffled_operator, const bool require_bucket_distribution) {
903
    // propagate error case
904
0
    if (*node_idx >= tnodes.size()) {
905
0
        return Status::InternalError(
906
0
                "Failed to reconstruct plan tree from thrift. Node id: {}, number of nodes: {}",
907
0
                *node_idx, tnodes.size());
908
0
    }
909
0
    const TPlanNode& tnode = tnodes[*node_idx];
910
911
0
    int num_children = tnodes[*node_idx].num_children;
912
0
    bool current_followed_by_shuffled_operator = followed_by_shuffled_operator;
913
0
    bool current_require_bucket_distribution = require_bucket_distribution;
914
    // TODO: Create CacheOperator is confused now
915
0
    OperatorPtr op = nullptr;
916
0
    OperatorPtr cache_op = nullptr;
917
0
    RETURN_IF_ERROR(_create_operator(pool, tnodes[*node_idx], descs, op, cur_pipe,
918
0
                                     parent == nullptr ? -1 : parent->node_id(), child_idx,
919
0
                                     followed_by_shuffled_operator,
920
0
                                     current_require_bucket_distribution, cache_op));
921
    // Initialization must be done here. For example, group by expressions in agg will be used to
922
    // decide if a local shuffle should be planed, so it must be initialized here.
923
0
    RETURN_IF_ERROR(op->init(tnode, _runtime_state.get()));
924
    // assert(parent != nullptr || (node_idx == 0 && root_expr != nullptr));
925
0
    if (parent != nullptr) {
926
        // add to parent's child(s)
927
0
        RETURN_IF_ERROR(parent->set_child(cache_op ? cache_op : op));
928
0
    } else {
929
0
        *root = op;
930
0
    }
931
    /**
932
     * `TLocalPartitionType::GLOBAL_EXECUTION_HASH_SHUFFLE` should be used if an operator is followed by a shuffled operator (shuffled hash join, union operator followed by co-located operators).
933
     *
934
     * For plan:
935
     * LocalExchange(id=0) -> Aggregation(id=1) -> ShuffledHashJoin(id=2)
936
     *                           Exchange(id=3) -> ShuffledHashJoinBuild(id=2)
937
     * We must ensure data distribution of `LocalExchange(id=0)` is same as Exchange(id=3).
938
     *
939
     * If an operator's is followed by a local exchange without shuffle (e.g. passthrough), a
940
     * shuffled local exchanger will be used before join so it is not followed by shuffle join.
941
     */
942
0
    auto required_data_distribution =
943
0
            cur_pipe->operators().empty()
944
0
                    ? cur_pipe->sink()->required_data_distribution(_runtime_state.get())
945
0
                    : op->required_data_distribution(_runtime_state.get());
946
0
    current_followed_by_shuffled_operator =
947
0
            ((followed_by_shuffled_operator ||
948
0
              (cur_pipe->operators().empty() ? cur_pipe->sink()->is_shuffled_operator()
949
0
                                             : op->is_shuffled_operator())) &&
950
0
             Pipeline::is_hash_exchange(required_data_distribution.distribution_type)) ||
951
0
            (followed_by_shuffled_operator &&
952
0
             required_data_distribution.distribution_type == TLocalPartitionType::NOOP);
953
954
0
    current_require_bucket_distribution =
955
0
            ((require_bucket_distribution ||
956
0
              (cur_pipe->operators().empty() ? cur_pipe->sink()->is_colocated_operator()
957
0
                                             : op->is_colocated_operator())) &&
958
0
             Pipeline::is_hash_exchange(required_data_distribution.distribution_type)) ||
959
0
            (require_bucket_distribution &&
960
0
             required_data_distribution.distribution_type == TLocalPartitionType::NOOP);
961
962
0
    if (num_children == 0) {
963
0
        _use_serial_source = op->is_serial_operator();
964
0
    }
965
    // rely on that tnodes is preorder of the plan
966
0
    for (int i = 0; i < num_children; i++) {
967
0
        ++*node_idx;
968
0
        RETURN_IF_ERROR(_create_tree_helper(pool, tnodes, descs, op, node_idx, nullptr, cur_pipe, i,
969
0
                                            current_followed_by_shuffled_operator,
970
0
                                            current_require_bucket_distribution));
971
972
        // we are expecting a child, but have used all nodes
973
        // this means we have been given a bad tree and must fail
974
0
        if (*node_idx >= tnodes.size()) {
975
0
            return Status::InternalError(
976
0
                    "Failed to reconstruct plan tree from thrift. Node id: {}, number of "
977
0
                    "nodes: {}",
978
0
                    *node_idx, tnodes.size());
979
0
        }
980
0
    }
981
982
0
    return Status::OK();
983
0
}
984
985
void PipelineFragmentContext::_inherit_pipeline_properties(
986
        const DataDistribution& data_distribution, PipelinePtr pipe_with_source,
987
0
        PipelinePtr pipe_with_sink) {
988
0
    pipe_with_sink->set_num_tasks(pipe_with_source->num_tasks());
989
0
    pipe_with_source->set_num_tasks(_num_instances);
990
0
    pipe_with_source->set_data_distribution(data_distribution);
991
0
}
992
993
Status PipelineFragmentContext::_add_local_exchange_impl(
994
        int idx, ObjectPool* pool, PipelinePtr cur_pipe, PipelinePtr new_pip,
995
        DataDistribution data_distribution, bool* do_local_exchange, int num_buckets,
996
        const std::map<int, int>& bucket_seq_to_instance_idx,
997
0
        const std::map<int, int>& shuffle_idx_to_instance_idx) {
998
0
    auto& operators = cur_pipe->operators();
999
0
    const auto downstream_pipeline_id = cur_pipe->id();
1000
0
    auto local_exchange_id = next_operator_id();
1001
    // 1. Create a new pipeline with local exchange sink.
1002
0
    DataSinkOperatorPtr sink;
1003
0
    auto sink_id = next_sink_operator_id();
1004
1005
    /**
1006
     * `bucket_seq_to_instance_idx` is empty if no scan operator is contained in this fragment.
1007
     * So co-located operators(e.g. Agg, Analytic) should use `HASH_SHUFFLE` instead of `BUCKET_HASH_SHUFFLE`.
1008
     */
1009
0
    const bool followed_by_shuffled_operator =
1010
0
            operators.size() > idx ? operators[idx]->followed_by_shuffled_operator()
1011
0
                                   : cur_pipe->sink()->followed_by_shuffled_operator();
1012
0
    const bool use_global_hash_shuffle = bucket_seq_to_instance_idx.empty() &&
1013
0
                                         !shuffle_idx_to_instance_idx.contains(-1) &&
1014
0
                                         followed_by_shuffled_operator && !_use_serial_source;
1015
0
    sink = std::make_shared<LocalExchangeSinkOperatorX>(
1016
0
            sink_id, local_exchange_id, use_global_hash_shuffle ? _total_instances : _num_instances,
1017
0
            data_distribution.partition_exprs, bucket_seq_to_instance_idx);
1018
0
    if (bucket_seq_to_instance_idx.empty() &&
1019
0
        data_distribution.distribution_type == TLocalPartitionType::BUCKET_HASH_SHUFFLE) {
1020
0
        data_distribution.distribution_type =
1021
0
                use_global_hash_shuffle ? TLocalPartitionType::GLOBAL_EXECUTION_HASH_SHUFFLE
1022
0
                                        : TLocalPartitionType::LOCAL_EXECUTION_HASH_SHUFFLE;
1023
0
    }
1024
0
    if (!use_global_hash_shuffle &&
1025
0
        data_distribution.distribution_type == TLocalPartitionType::GLOBAL_EXECUTION_HASH_SHUFFLE) {
1026
0
        data_distribution.distribution_type = TLocalPartitionType::LOCAL_EXECUTION_HASH_SHUFFLE;
1027
0
    }
1028
0
    RETURN_IF_ERROR(new_pip->set_sink(sink));
1029
0
    RETURN_IF_ERROR(new_pip->sink()->init(_runtime_state.get(), data_distribution.distribution_type,
1030
0
                                          num_buckets, shuffle_idx_to_instance_idx));
1031
1032
    // 2. Create and initialize LocalExchangeSharedState.
1033
0
    std::shared_ptr<LocalExchangeSharedState> shared_state =
1034
0
            LocalExchangeSharedState::create_shared(_num_instances);
1035
0
    switch (data_distribution.distribution_type) {
1036
0
    case TLocalPartitionType::LOCAL_EXECUTION_HASH_SHUFFLE:
1037
0
    case TLocalPartitionType::GLOBAL_EXECUTION_HASH_SHUFFLE:
1038
0
        shared_state->exchanger = ShuffleExchanger::create_unique(
1039
0
                std::max(cur_pipe->num_tasks(), _num_instances), _num_instances,
1040
0
                use_global_hash_shuffle ? _total_instances : _num_instances,
1041
0
                _runtime_state->query_options().__isset.local_exchange_free_blocks_limit
1042
0
                        ? cast_set<int>(
1043
0
                                  _runtime_state->query_options().local_exchange_free_blocks_limit)
1044
0
                        : 0,
1045
0
                data_distribution.distribution_type);
1046
0
        break;
1047
0
    case TLocalPartitionType::BUCKET_HASH_SHUFFLE:
1048
0
        shared_state->exchanger = BucketShuffleExchanger::create_unique(
1049
0
                std::max(cur_pipe->num_tasks(), _num_instances), _num_instances, num_buckets,
1050
0
                _runtime_state->query_options().__isset.local_exchange_free_blocks_limit
1051
0
                        ? cast_set<int>(
1052
0
                                  _runtime_state->query_options().local_exchange_free_blocks_limit)
1053
0
                        : 0);
1054
0
        break;
1055
0
    case TLocalPartitionType::PASSTHROUGH:
1056
0
        shared_state->exchanger = PassthroughExchanger::create_unique(
1057
0
                cur_pipe->num_tasks(), _num_instances,
1058
0
                _runtime_state->query_options().__isset.local_exchange_free_blocks_limit
1059
0
                        ? cast_set<int>(
1060
0
                                  _runtime_state->query_options().local_exchange_free_blocks_limit)
1061
0
                        : 0);
1062
0
        break;
1063
0
    case TLocalPartitionType::BROADCAST:
1064
0
        shared_state->exchanger = BroadcastExchanger::create_unique(
1065
0
                cur_pipe->num_tasks(), _num_instances,
1066
0
                _runtime_state->query_options().__isset.local_exchange_free_blocks_limit
1067
0
                        ? cast_set<int>(
1068
0
                                  _runtime_state->query_options().local_exchange_free_blocks_limit)
1069
0
                        : 0);
1070
0
        break;
1071
0
    case TLocalPartitionType::PASS_TO_ONE:
1072
0
        if (_runtime_state->enable_share_hash_table_for_broadcast_join()) {
1073
            // If shared hash table is enabled for BJ, hash table will be built by only one task
1074
0
            shared_state->exchanger = PassToOneExchanger::create_unique(
1075
0
                    cur_pipe->num_tasks(), _num_instances,
1076
0
                    _runtime_state->query_options().__isset.local_exchange_free_blocks_limit
1077
0
                            ? cast_set<int>(_runtime_state->query_options()
1078
0
                                                    .local_exchange_free_blocks_limit)
1079
0
                            : 0);
1080
0
        } else {
1081
0
            shared_state->exchanger = BroadcastExchanger::create_unique(
1082
0
                    cur_pipe->num_tasks(), _num_instances,
1083
0
                    _runtime_state->query_options().__isset.local_exchange_free_blocks_limit
1084
0
                            ? cast_set<int>(_runtime_state->query_options()
1085
0
                                                    .local_exchange_free_blocks_limit)
1086
0
                            : 0);
1087
0
        }
1088
0
        break;
1089
0
    case TLocalPartitionType::ADAPTIVE_PASSTHROUGH:
1090
0
        shared_state->exchanger = AdaptivePassthroughExchanger::create_unique(
1091
0
                std::max(cur_pipe->num_tasks(), _num_instances), _num_instances,
1092
0
                _runtime_state->query_options().__isset.local_exchange_free_blocks_limit
1093
0
                        ? cast_set<int>(
1094
0
                                  _runtime_state->query_options().local_exchange_free_blocks_limit)
1095
0
                        : 0);
1096
0
        break;
1097
0
    default:
1098
0
        return Status::InternalError("Unsupported local exchange type : " +
1099
0
                                     std::to_string((int)data_distribution.distribution_type));
1100
0
    }
1101
0
    shared_state->create_source_dependencies(_num_instances, local_exchange_id, local_exchange_id,
1102
0
                                             "LOCAL_EXCHANGE_OPERATOR");
1103
0
    shared_state->create_sink_dependency(sink_id, local_exchange_id, "LOCAL_EXCHANGE_SINK");
1104
0
    _op_id_to_shared_state.insert({local_exchange_id, {shared_state, shared_state->sink_deps}});
1105
1106
    // 3. Set two pipelines' operator list. For example, split pipeline [Scan - AggSink] to
1107
    // pipeline1 [Scan - LocalExchangeSink] and pipeline2 [LocalExchangeSource - AggSink].
1108
1109
    // 3.1 Initialize new pipeline's operator list.
1110
0
    std::copy(operators.begin(), operators.begin() + idx,
1111
0
              std::inserter(new_pip->operators(), new_pip->operators().end()));
1112
1113
    // 3.2 Erase unused operators in previous pipeline.
1114
0
    operators.erase(operators.begin(), operators.begin() + idx);
1115
1116
    // 4. Initialize LocalExchangeSource and insert it into this pipeline.
1117
0
    OperatorPtr source_op;
1118
0
    source_op = std::make_shared<LocalExchangeSourceOperatorX>(pool, local_exchange_id);
1119
0
    RETURN_IF_ERROR(source_op->set_child(new_pip->operators().back()));
1120
0
    RETURN_IF_ERROR(source_op->init(data_distribution.distribution_type));
1121
0
    if (!operators.empty()) {
1122
0
        RETURN_IF_ERROR(operators.front()->set_child(nullptr));
1123
0
        RETURN_IF_ERROR(operators.front()->set_child(source_op));
1124
0
    }
1125
0
    operators.insert(operators.begin(), source_op);
1126
1127
    // 5. Set children for two pipelines separately.
1128
0
    std::vector<std::shared_ptr<Pipeline>> new_children;
1129
0
    std::vector<PipelineId> edges_with_source;
1130
0
    for (auto child : cur_pipe->children()) {
1131
0
        bool found = false;
1132
0
        for (auto op : new_pip->operators()) {
1133
0
            if (child->sink()->node_id() == op->node_id()) {
1134
0
                new_pip->set_children(child);
1135
0
                found = true;
1136
0
            };
1137
0
        }
1138
0
        if (!found) {
1139
0
            new_children.push_back(child);
1140
0
            edges_with_source.push_back(child->id());
1141
0
        }
1142
0
    }
1143
0
    new_children.push_back(new_pip);
1144
0
    edges_with_source.push_back(new_pip->id());
1145
1146
    // 6. Set DAG for new pipelines.
1147
0
    if (!new_pip->children().empty()) {
1148
0
        std::vector<PipelineId> edges_with_sink;
1149
0
        for (auto child : new_pip->children()) {
1150
0
            edges_with_sink.push_back(child->id());
1151
0
        }
1152
0
        _dag.insert({new_pip->id(), edges_with_sink});
1153
0
    }
1154
0
    cur_pipe->set_children(new_children);
1155
0
    _dag[downstream_pipeline_id] = edges_with_source;
1156
0
    RETURN_IF_ERROR(new_pip->sink()->set_child(new_pip->operators().back()));
1157
0
    RETURN_IF_ERROR(cur_pipe->sink()->set_child(nullptr));
1158
0
    RETURN_IF_ERROR(cur_pipe->sink()->set_child(cur_pipe->operators().back()));
1159
1160
    // 7. Inherit properties from current pipeline.
1161
0
    _inherit_pipeline_properties(data_distribution, cur_pipe, new_pip);
1162
0
    return Status::OK();
1163
0
}
1164
1165
Status PipelineFragmentContext::_add_local_exchange(
1166
        int pip_idx, int idx, int node_id, ObjectPool* pool, PipelinePtr cur_pipe,
1167
        DataDistribution data_distribution, bool* do_local_exchange, int num_buckets,
1168
        const std::map<int, int>& bucket_seq_to_instance_idx,
1169
0
        const std::map<int, int>& shuffle_idx_to_instance_idx) {
1170
0
    if (_num_instances <= 1 || cur_pipe->num_tasks_of_parent() <= 1) {
1171
0
        return Status::OK();
1172
0
    }
1173
1174
0
    if (!cur_pipe->need_to_local_exchange(data_distribution, idx)) {
1175
0
        return Status::OK();
1176
0
    }
1177
0
    *do_local_exchange = true;
1178
1179
0
    auto& operators = cur_pipe->operators();
1180
0
    auto total_op_num = operators.size();
1181
0
    auto new_pip = add_pipeline(cur_pipe, pip_idx + 1);
1182
0
    RETURN_IF_ERROR(_add_local_exchange_impl(
1183
0
            idx, pool, cur_pipe, new_pip, data_distribution, do_local_exchange, num_buckets,
1184
0
            bucket_seq_to_instance_idx, shuffle_idx_to_instance_idx));
1185
1186
0
    CHECK(total_op_num + 1 == cur_pipe->operators().size() + new_pip->operators().size())
1187
0
            << "total_op_num: " << total_op_num
1188
0
            << " cur_pipe->operators().size(): " << cur_pipe->operators().size()
1189
0
            << " new_pip->operators().size(): " << new_pip->operators().size();
1190
1191
    // There are some local shuffles with relatively heavy operations on the sink.
1192
    // If the local sink concurrency is 1 and the local source concurrency is n, the sink becomes a bottleneck.
1193
    // Therefore, local passthrough is used to increase the concurrency of the sink.
1194
    // op -> local sink(1) -> local source (n)
1195
    // op -> local passthrough(1) -> local passthrough(n) ->  local sink(n) -> local source (n)
1196
0
    if (cur_pipe->num_tasks() > 1 && new_pip->num_tasks() == 1 &&
1197
0
        Pipeline::heavy_operations_on_the_sink(data_distribution.distribution_type)) {
1198
0
        RETURN_IF_ERROR(_add_local_exchange_impl(
1199
0
                cast_set<int>(new_pip->operators().size()), pool, new_pip,
1200
0
                add_pipeline(new_pip, pip_idx + 2),
1201
0
                DataDistribution(TLocalPartitionType::PASSTHROUGH), do_local_exchange, num_buckets,
1202
0
                bucket_seq_to_instance_idx, shuffle_idx_to_instance_idx));
1203
0
    }
1204
0
    return Status::OK();
1205
0
}
1206
1207
Status PipelineFragmentContext::_plan_local_exchange(
1208
        int num_buckets, const std::map<int, int>& bucket_seq_to_instance_idx,
1209
0
        const std::map<int, int>& shuffle_idx_to_instance_idx) {
1210
0
    for (int pip_idx = cast_set<int>(_pipelines.size()) - 1; pip_idx >= 0; pip_idx--) {
1211
0
        _pipelines[pip_idx]->init_data_distribution(_runtime_state.get());
1212
        // Set property if child pipeline is not join operator's child.
1213
0
        if (!_pipelines[pip_idx]->children().empty()) {
1214
0
            for (auto& child : _pipelines[pip_idx]->children()) {
1215
0
                if (child->sink()->node_id() ==
1216
0
                    _pipelines[pip_idx]->operators().front()->node_id()) {
1217
0
                    _pipelines[pip_idx]->set_data_distribution(child->data_distribution());
1218
0
                }
1219
0
            }
1220
0
        }
1221
1222
        // if 'num_buckets == 0' means the fragment is colocated by exchange node not the
1223
        // scan node. so here use `_num_instance` to replace the `num_buckets` to prevent dividing 0
1224
        // still keep colocate plan after local shuffle
1225
0
        RETURN_IF_ERROR(_plan_local_exchange(num_buckets, pip_idx, _pipelines[pip_idx],
1226
0
                                             bucket_seq_to_instance_idx,
1227
0
                                             shuffle_idx_to_instance_idx));
1228
0
    }
1229
0
    return Status::OK();
1230
0
}
1231
1232
Status PipelineFragmentContext::_plan_local_exchange(
1233
        int num_buckets, int pip_idx, PipelinePtr pip,
1234
        const std::map<int, int>& bucket_seq_to_instance_idx,
1235
0
        const std::map<int, int>& shuffle_idx_to_instance_idx) {
1236
0
    int idx = 1;
1237
0
    bool do_local_exchange = false;
1238
0
    do {
1239
0
        auto& ops = pip->operators();
1240
0
        do_local_exchange = false;
1241
        // Plan local exchange for each operator.
1242
0
        for (; idx < ops.size();) {
1243
0
            auto _le_req = ops[idx]->required_data_distribution(_runtime_state.get());
1244
0
            if (_le_req.need_local_exchange()) {
1245
0
                RETURN_IF_ERROR(_add_local_exchange(
1246
0
                        pip_idx, idx, ops[idx]->node_id(), _runtime_state->obj_pool(), pip, _le_req,
1247
0
                        &do_local_exchange, num_buckets, bucket_seq_to_instance_idx,
1248
0
                        shuffle_idx_to_instance_idx));
1249
0
            }
1250
0
            if (do_local_exchange) {
1251
                // If local exchange is needed for current operator, we will split this pipeline to
1252
                // two pipelines by local exchange sink/source. And then we need to process remaining
1253
                // operators in this pipeline so we set idx to 2 (0 is local exchange source and 1
1254
                // is current operator was already processed) and continue to plan local exchange.
1255
0
                idx = 2;
1256
0
                break;
1257
0
            }
1258
0
            idx++;
1259
0
        }
1260
0
    } while (do_local_exchange);
1261
0
    if (pip->sink()->required_data_distribution(_runtime_state.get()).need_local_exchange()) {
1262
0
        RETURN_IF_ERROR(_add_local_exchange(
1263
0
                pip_idx, idx, pip->sink()->node_id(), _runtime_state->obj_pool(), pip,
1264
0
                pip->sink()->required_data_distribution(_runtime_state.get()), &do_local_exchange,
1265
0
                num_buckets, bucket_seq_to_instance_idx, shuffle_idx_to_instance_idx));
1266
0
    }
1267
0
    return Status::OK();
1268
0
}
1269
1270
Status PipelineFragmentContext::_create_data_sink(ObjectPool* pool, const TDataSink& thrift_sink,
1271
                                                  const std::vector<TExpr>& output_exprs,
1272
                                                  const TPipelineFragmentParams& params,
1273
                                                  const RowDescriptor& row_desc,
1274
                                                  RuntimeState* state, DescriptorTbl& desc_tbl,
1275
0
                                                  PipelineId cur_pipeline_id) {
1276
0
    switch (thrift_sink.type) {
1277
0
    case TDataSinkType::DATA_STREAM_SINK: {
1278
0
        if (!thrift_sink.__isset.stream_sink) {
1279
0
            return Status::InternalError("Missing data stream sink.");
1280
0
        }
1281
0
        _sink = std::make_shared<ExchangeSinkOperatorX>(
1282
0
                state, row_desc, next_sink_operator_id(), thrift_sink.stream_sink,
1283
0
                params.destinations, _fragment_instance_ids);
1284
0
        break;
1285
0
    }
1286
0
    case TDataSinkType::RESULT_SINK: {
1287
0
        if (!thrift_sink.__isset.result_sink) {
1288
0
            return Status::InternalError("Missing data buffer sink.");
1289
0
        }
1290
1291
0
        auto& pipeline = _pipelines[cur_pipeline_id];
1292
0
        int child_node_id = pipeline->operators().back()->node_id();
1293
0
        _sink = std::make_shared<ResultSinkOperatorX>(next_sink_operator_id(), child_node_id + 1,
1294
0
                                                      row_desc, output_exprs,
1295
0
                                                      thrift_sink.result_sink);
1296
0
        break;
1297
0
    }
1298
0
    case TDataSinkType::DICTIONARY_SINK: {
1299
0
        if (!thrift_sink.__isset.dictionary_sink) {
1300
0
            return Status::InternalError("Missing dict sink.");
1301
0
        }
1302
1303
0
        _sink = std::make_shared<DictSinkOperatorX>(next_sink_operator_id(), row_desc, output_exprs,
1304
0
                                                    thrift_sink.dictionary_sink);
1305
0
        break;
1306
0
    }
1307
0
    case TDataSinkType::GROUP_COMMIT_OLAP_TABLE_SINK:
1308
0
    case TDataSinkType::OLAP_TABLE_SINK: {
1309
0
        auto& pipeline = _pipelines[cur_pipeline_id];
1310
0
        int child_node_id = pipeline->operators().back()->node_id();
1311
0
        if (state->query_options().enable_memtable_on_sink_node &&
1312
0
            !_has_inverted_index_v1_or_partial_update(thrift_sink.olap_table_sink) &&
1313
0
            !_has_row_binlog(thrift_sink.olap_table_sink) && !config::is_cloud_mode()) {
1314
0
            _sink = std::make_shared<OlapTableSinkV2OperatorX>(
1315
0
                    pool, next_sink_operator_id(), child_node_id + 1, row_desc, output_exprs);
1316
0
        } else {
1317
0
            _sink = std::make_shared<OlapTableSinkOperatorX>(
1318
0
                    pool, next_sink_operator_id(), child_node_id + 1, row_desc, output_exprs);
1319
0
        }
1320
0
        break;
1321
0
    }
1322
0
    case TDataSinkType::GROUP_COMMIT_BLOCK_SINK: {
1323
0
        DCHECK(thrift_sink.__isset.olap_table_sink);
1324
0
        DCHECK(state->get_query_ctx() != nullptr);
1325
0
        state->get_query_ctx()->query_mem_tracker()->is_group_commit_load = true;
1326
0
        _sink = std::make_shared<GroupCommitBlockSinkOperatorX>(next_sink_operator_id(), row_desc,
1327
0
                                                                output_exprs);
1328
0
        break;
1329
0
    }
1330
0
    case TDataSinkType::HIVE_TABLE_SINK: {
1331
0
        if (!thrift_sink.__isset.hive_table_sink) {
1332
0
            return Status::InternalError("Missing hive table sink.");
1333
0
        }
1334
0
        _sink = std::make_shared<HiveTableSinkOperatorX>(pool, next_sink_operator_id(), row_desc,
1335
0
                                                         output_exprs);
1336
0
        break;
1337
0
    }
1338
0
    case TDataSinkType::ICEBERG_TABLE_SINK: {
1339
0
        if (!thrift_sink.__isset.iceberg_table_sink) {
1340
0
            return Status::InternalError("Missing iceberg table sink.");
1341
0
        }
1342
0
        if (thrift_sink.iceberg_table_sink.__isset.sort_info) {
1343
0
            _sink = std::make_shared<SpillIcebergTableSinkOperatorX>(pool, next_sink_operator_id(),
1344
0
                                                                     row_desc, output_exprs);
1345
0
        } else {
1346
0
            _sink = std::make_shared<IcebergTableSinkOperatorX>(pool, next_sink_operator_id(),
1347
0
                                                                row_desc, output_exprs);
1348
0
        }
1349
0
        break;
1350
0
    }
1351
0
    case TDataSinkType::ICEBERG_DELETE_SINK: {
1352
0
        if (!thrift_sink.__isset.iceberg_delete_sink) {
1353
0
            return Status::InternalError("Missing iceberg delete sink.");
1354
0
        }
1355
0
        _sink = std::make_shared<IcebergDeleteSinkOperatorX>(pool, next_sink_operator_id(),
1356
0
                                                             row_desc, output_exprs);
1357
0
        break;
1358
0
    }
1359
0
    case TDataSinkType::ICEBERG_MERGE_SINK: {
1360
0
        if (!thrift_sink.__isset.iceberg_merge_sink) {
1361
0
            return Status::InternalError("Missing iceberg merge sink.");
1362
0
        }
1363
0
        _sink = std::make_shared<IcebergMergeSinkOperatorX>(pool, next_sink_operator_id(), row_desc,
1364
0
                                                            output_exprs);
1365
0
        break;
1366
0
    }
1367
0
    case TDataSinkType::MAXCOMPUTE_TABLE_SINK: {
1368
0
        if (!thrift_sink.__isset.max_compute_table_sink) {
1369
0
            return Status::InternalError("Missing max compute table sink.");
1370
0
        }
1371
0
        _sink = std::make_shared<MCTableSinkOperatorX>(pool, next_sink_operator_id(), row_desc,
1372
0
                                                       output_exprs);
1373
0
        break;
1374
0
    }
1375
0
    case TDataSinkType::JDBC_TABLE_SINK: {
1376
0
        if (!thrift_sink.__isset.jdbc_table_sink) {
1377
0
            return Status::InternalError("Missing data jdbc sink.");
1378
0
        }
1379
0
        if (config::enable_java_support) {
1380
0
            _sink = std::make_shared<JdbcTableSinkOperatorX>(row_desc, next_sink_operator_id(),
1381
0
                                                             output_exprs);
1382
0
        } else {
1383
0
            return Status::InternalError(
1384
0
                    "Jdbc table sink is not enabled, you can change be config "
1385
0
                    "enable_java_support to true and restart be.");
1386
0
        }
1387
0
        break;
1388
0
    }
1389
0
    case TDataSinkType::MEMORY_SCRATCH_SINK: {
1390
0
        if (!thrift_sink.__isset.memory_scratch_sink) {
1391
0
            return Status::InternalError("Missing data buffer sink.");
1392
0
        }
1393
1394
0
        _sink = std::make_shared<MemoryScratchSinkOperatorX>(row_desc, next_sink_operator_id(),
1395
0
                                                             output_exprs);
1396
0
        break;
1397
0
    }
1398
0
    case TDataSinkType::RESULT_FILE_SINK: {
1399
0
        if (!thrift_sink.__isset.result_file_sink) {
1400
0
            return Status::InternalError("Missing result file sink.");
1401
0
        }
1402
1403
        // Result file sink is not the top sink
1404
0
        if (params.__isset.destinations && !params.destinations.empty()) {
1405
0
            _sink = std::make_shared<ResultFileSinkOperatorX>(
1406
0
                    next_sink_operator_id(), row_desc, thrift_sink.result_file_sink,
1407
0
                    params.destinations, output_exprs, desc_tbl);
1408
0
        } else {
1409
0
            _sink = std::make_shared<ResultFileSinkOperatorX>(next_sink_operator_id(), row_desc,
1410
0
                                                              output_exprs);
1411
0
        }
1412
0
        break;
1413
0
    }
1414
0
    case TDataSinkType::MULTI_CAST_DATA_STREAM_SINK: {
1415
0
        DCHECK(thrift_sink.__isset.multi_cast_stream_sink);
1416
0
        DCHECK_GT(thrift_sink.multi_cast_stream_sink.sinks.size(), 0);
1417
0
        auto sink_id = next_sink_operator_id();
1418
0
        const int multi_cast_node_id = sink_id;
1419
0
        auto sender_size = thrift_sink.multi_cast_stream_sink.sinks.size();
1420
        // one sink has multiple sources.
1421
0
        std::vector<int> sources;
1422
0
        for (int i = 0; i < sender_size; ++i) {
1423
0
            auto source_id = next_operator_id();
1424
0
            sources.push_back(source_id);
1425
0
        }
1426
1427
0
        _sink = std::make_shared<MultiCastDataStreamSinkOperatorX>(
1428
0
                sink_id, multi_cast_node_id, sources, pool, thrift_sink.multi_cast_stream_sink);
1429
0
        for (int i = 0; i < sender_size; ++i) {
1430
0
            auto new_pipeline = add_pipeline();
1431
            // use to exchange sink
1432
0
            RowDescriptor* exchange_row_desc = nullptr;
1433
0
            {
1434
0
                const auto& tmp_row_desc =
1435
0
                        !thrift_sink.multi_cast_stream_sink.sinks[i].output_exprs.empty()
1436
0
                                ? RowDescriptor(state->desc_tbl(),
1437
0
                                                {thrift_sink.multi_cast_stream_sink.sinks[i]
1438
0
                                                         .output_tuple_id})
1439
0
                                : row_desc;
1440
0
                exchange_row_desc = pool->add(new RowDescriptor(tmp_row_desc));
1441
0
            }
1442
0
            auto source_id = sources[i];
1443
0
            OperatorPtr source_op;
1444
            // 1. create and set the source operator of multi_cast_data_stream_source for new pipeline
1445
0
            source_op = std::make_shared<MultiCastDataStreamerSourceOperatorX>(
1446
0
                    /*node_id*/ source_id, /*consumer_id*/ i, pool,
1447
0
                    thrift_sink.multi_cast_stream_sink.sinks[i], row_desc,
1448
0
                    /*operator_id=*/source_id);
1449
0
            RETURN_IF_ERROR(new_pipeline->add_operator(
1450
0
                    source_op, params.__isset.parallel_instances ? params.parallel_instances : 0));
1451
            // 2. create and set sink operator of data stream sender for new pipeline
1452
1453
0
            DataSinkOperatorPtr sink_op;
1454
0
            sink_op = std::make_shared<ExchangeSinkOperatorX>(
1455
0
                    state, *exchange_row_desc, next_sink_operator_id(),
1456
0
                    thrift_sink.multi_cast_stream_sink.sinks[i],
1457
0
                    thrift_sink.multi_cast_stream_sink.destinations[i], _fragment_instance_ids);
1458
1459
0
            RETURN_IF_ERROR(new_pipeline->set_sink(sink_op));
1460
0
            {
1461
0
                TDataSink* t = pool->add(new TDataSink());
1462
0
                t->stream_sink = thrift_sink.multi_cast_stream_sink.sinks[i];
1463
0
                RETURN_IF_ERROR(sink_op->init(*t));
1464
0
            }
1465
1466
            // 3. set dependency dag
1467
0
            _dag[new_pipeline->id()].push_back(cur_pipeline_id);
1468
0
        }
1469
0
        if (sources.empty()) {
1470
0
            return Status::InternalError("size of sources must be greater than 0");
1471
0
        }
1472
0
        break;
1473
0
    }
1474
0
    case TDataSinkType::BLACKHOLE_SINK: {
1475
0
        if (!thrift_sink.__isset.blackhole_sink) {
1476
0
            return Status::InternalError("Missing blackhole sink.");
1477
0
        }
1478
1479
0
        _sink.reset(new BlackholeSinkOperatorX(next_sink_operator_id()));
1480
0
        break;
1481
0
    }
1482
0
    case TDataSinkType::TVF_TABLE_SINK: {
1483
0
        if (!thrift_sink.__isset.tvf_table_sink) {
1484
0
            return Status::InternalError("Missing TVF table sink.");
1485
0
        }
1486
0
        _sink = std::make_shared<TVFTableSinkOperatorX>(pool, next_sink_operator_id(), row_desc,
1487
0
                                                        output_exprs);
1488
0
        break;
1489
0
    }
1490
0
    default:
1491
0
        return Status::InternalError("Unsuported sink type in pipeline: {}", thrift_sink.type);
1492
0
    }
1493
0
    return Status::OK();
1494
0
}
1495
1496
// NOLINTBEGIN(readability-function-size)
1497
// NOLINTBEGIN(readability-function-cognitive-complexity)
1498
Status PipelineFragmentContext::_create_operator(ObjectPool* pool, const TPlanNode& tnode,
1499
                                                 const DescriptorTbl& descs, OperatorPtr& op,
1500
                                                 PipelinePtr& cur_pipe, int parent_idx,
1501
                                                 int child_idx,
1502
                                                 const bool followed_by_shuffled_operator,
1503
                                                 const bool require_bucket_distribution,
1504
2
                                                 OperatorPtr& cache_op) {
1505
2
    std::vector<DataSinkOperatorPtr> sink_ops;
1506
2
    Defer defer = Defer([&]() {
1507
2
        if (op) {
1508
0
            op->update_operator(tnode, followed_by_shuffled_operator, require_bucket_distribution);
1509
0
        }
1510
2
        for (auto& s : sink_ops) {
1511
0
            s->update_operator(tnode, followed_by_shuffled_operator, require_bucket_distribution);
1512
0
        }
1513
2
    });
1514
    // We directly construct the operator from Thrift because the given array is in the order of preorder traversal.
1515
    // Therefore, here we need to use a stack-like structure.
1516
2
    _pipeline_parent_map.pop(cur_pipe, parent_idx, child_idx);
1517
2
    std::stringstream error_msg;
1518
2
    bool enable_query_cache = _params.fragment.__isset.query_cache_param;
1519
1520
2
    bool fe_with_old_version = false;
1521
2
    switch (tnode.node_type) {
1522
2
    case TPlanNodeType::OLAP_SCAN_NODE: {
1523
2
        if (enable_query_cache) {
1524
2
            if (_query_cache_runtime == nullptr) {
1525
                // The plan tree is built in pre-order and the cache source
1526
                // sits above the scan, so the runtime it created must already
1527
                // exist here. Running the scan with its own runtime instead
1528
                // would silently drop data on a HIT (the scan skips scanning
1529
                // while no cache source emits the entry), so a malformed plan
1530
                // shape must fail loudly.
1531
2
                return Status::InternalError(
1532
2
                        "query cache runtime is absent at the scan node, node_id={}, "
1533
2
                        "cache node_id={}",
1534
2
                        tnode.node_id, _params.fragment.query_cache_param.node_id);
1535
2
            }
1536
0
            if (tnode.olap_scan_node.__isset.read_row_binlog &&
1537
0
                tnode.olap_scan_node.read_row_binlog) {
1538
                // Row-binlog scans read a different data stream: they must
1539
                // neither serve nor fill the query cache.
1540
0
                _query_cache_runtime->disable_for_binlog_scan();
1541
0
            }
1542
0
        }
1543
0
        op = std::make_shared<OlapScanOperatorX>(
1544
0
                pool, tnode, next_operator_id(), descs, _num_instances,
1545
0
                enable_query_cache ? _params.fragment.query_cache_param : TQueryCacheParam {},
1546
0
                enable_query_cache ? _query_cache_runtime : nullptr);
1547
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1548
0
        fe_with_old_version = !tnode.__isset.is_serial_operator;
1549
0
        break;
1550
0
    }
1551
0
    case TPlanNodeType::GROUP_COMMIT_SCAN_NODE: {
1552
0
        DCHECK(_query_ctx != nullptr);
1553
0
        _query_ctx->query_mem_tracker()->is_group_commit_load = true;
1554
0
        op = std::make_shared<GroupCommitOperatorX>(pool, tnode, next_operator_id(), descs,
1555
0
                                                    _num_instances);
1556
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1557
0
        fe_with_old_version = !tnode.__isset.is_serial_operator;
1558
0
        break;
1559
0
    }
1560
0
    case TPlanNodeType::JDBC_SCAN_NODE: {
1561
0
        if (config::enable_java_support) {
1562
0
            op = std::make_shared<JDBCScanOperatorX>(pool, tnode, next_operator_id(), descs,
1563
0
                                                     _num_instances);
1564
0
            RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1565
0
        } else {
1566
0
            return Status::InternalError(
1567
0
                    "Jdbc scan node is disabled, you can change be config enable_java_support "
1568
0
                    "to true and restart be.");
1569
0
        }
1570
0
        fe_with_old_version = !tnode.__isset.is_serial_operator;
1571
0
        break;
1572
0
    }
1573
0
    case TPlanNodeType::FILE_SCAN_NODE: {
1574
0
        op = std::make_shared<FileScanOperatorX>(pool, tnode, next_operator_id(), descs,
1575
0
                                                 _num_instances);
1576
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1577
0
        fe_with_old_version = !tnode.__isset.is_serial_operator;
1578
0
        break;
1579
0
    }
1580
0
    case TPlanNodeType::EXCHANGE_NODE: {
1581
0
        int num_senders = _params.per_exch_num_senders.contains(tnode.node_id)
1582
0
                                  ? _params.per_exch_num_senders.find(tnode.node_id)->second
1583
0
                                  : 0;
1584
0
        DCHECK_GT(num_senders, 0);
1585
0
        auto exchange_op = std::make_shared<ExchangeSourceOperatorX>(
1586
0
                pool, tnode, next_operator_id(), descs, num_senders);
1587
0
        if (!_params.bucket_seq_to_instance_idx.empty()) {
1588
            // Lets bucket-routed exchanges detect orphan instances (owning no bucket) that
1589
            // no sender channel will ever address — their receivers must start at EOS.
1590
0
            exchange_op->set_bucket_dest_instances(_params.bucket_seq_to_instance_idx);
1591
0
        }
1592
0
        op = exchange_op;
1593
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1594
0
        fe_with_old_version = !tnode.__isset.is_serial_operator;
1595
0
        break;
1596
0
    }
1597
0
    case TPlanNodeType::AGGREGATION_NODE: {
1598
0
        if (tnode.agg_node.grouping_exprs.empty() &&
1599
0
            descs.get_tuple_descriptor(tnode.agg_node.output_tuple_id)->slots().empty()) {
1600
0
            return Status::InternalError("Illegal aggregate node " + std::to_string(tnode.node_id) +
1601
0
                                         ": group by and output is empty");
1602
0
        }
1603
0
        bool need_create_cache_op =
1604
0
                enable_query_cache && tnode.node_id == _params.fragment.query_cache_param.node_id;
1605
0
        auto create_query_cache_operator = [&](PipelinePtr& new_pipe) {
1606
0
            auto cache_node_id = _params.local_params[0].per_node_scan_ranges.begin()->first;
1607
0
            auto cache_source_id = next_operator_id();
1608
0
            if (_query_cache_runtime == nullptr) {
1609
0
                _query_cache_runtime =
1610
0
                        std::make_shared<QueryCacheRuntime>(_params.fragment.query_cache_param);
1611
0
            }
1612
0
            op = std::make_shared<CacheSourceOperatorX>(pool, cache_node_id, cache_source_id,
1613
0
                                                        _params.fragment.query_cache_param,
1614
0
                                                        _query_cache_runtime);
1615
0
            RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1616
1617
0
            const auto downstream_pipeline_id = cur_pipe->id();
1618
0
            if (!_dag.contains(downstream_pipeline_id)) {
1619
0
                _dag.insert({downstream_pipeline_id, {}});
1620
0
            }
1621
0
            new_pipe = add_pipeline(cur_pipe);
1622
0
            _dag[downstream_pipeline_id].push_back(new_pipe->id());
1623
1624
0
            DataSinkOperatorPtr cache_sink(new CacheSinkOperatorX(
1625
0
                    next_sink_operator_id(), op->node_id(), op->operator_id()));
1626
0
            RETURN_IF_ERROR(new_pipe->set_sink(cache_sink));
1627
0
            return Status::OK();
1628
0
        };
1629
0
        const bool group_by_limit_opt =
1630
0
                tnode.agg_node.__isset.agg_sort_info_by_group_key && tnode.limit > 0;
1631
1632
        /// PartitionedAggSourceOperatorX does not support "group by limit opt(#29641)" yet.
1633
        /// If `group_by_limit_opt` is true, then it might not need to spill at all.
1634
0
        const bool enable_spill = _runtime_state->enable_spill() &&
1635
0
                                  !tnode.agg_node.grouping_exprs.empty() && !group_by_limit_opt;
1636
0
        const bool is_streaming_agg = tnode.agg_node.__isset.use_streaming_preaggregation &&
1637
0
                                      tnode.agg_node.use_streaming_preaggregation &&
1638
0
                                      !tnode.agg_node.grouping_exprs.empty();
1639
        // TODO: distinct streaming agg does not support spill.
1640
0
        const bool can_use_distinct_streaming_agg =
1641
0
                (!enable_spill || is_streaming_agg) && tnode.agg_node.aggregate_functions.empty() &&
1642
0
                !tnode.agg_node.__isset.agg_sort_info_by_group_key &&
1643
0
                _params.query_options.__isset.enable_distinct_streaming_aggregation &&
1644
0
                _params.query_options.enable_distinct_streaming_aggregation;
1645
1646
0
        if (can_use_distinct_streaming_agg) {
1647
0
            if (need_create_cache_op) {
1648
0
                PipelinePtr new_pipe;
1649
0
                RETURN_IF_ERROR(create_query_cache_operator(new_pipe));
1650
1651
0
                cache_op = op;
1652
0
                op = std::make_shared<DistinctStreamingAggOperatorX>(pool, next_operator_id(),
1653
0
                                                                     tnode, descs);
1654
0
                RETURN_IF_ERROR(new_pipe->add_operator(op, _parallel_instances));
1655
0
                RETURN_IF_ERROR(cur_pipe->operators().front()->set_child(op));
1656
0
                cur_pipe = new_pipe;
1657
0
            } else {
1658
0
                op = std::make_shared<DistinctStreamingAggOperatorX>(pool, next_operator_id(),
1659
0
                                                                     tnode, descs);
1660
0
                RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1661
0
            }
1662
0
        } else if (is_streaming_agg) {
1663
0
            if (need_create_cache_op) {
1664
0
                PipelinePtr new_pipe;
1665
0
                RETURN_IF_ERROR(create_query_cache_operator(new_pipe));
1666
0
                cache_op = op;
1667
0
                op = std::make_shared<StreamingAggOperatorX>(pool, next_operator_id(), tnode,
1668
0
                                                             descs);
1669
0
                RETURN_IF_ERROR(cur_pipe->operators().front()->set_child(op));
1670
0
                RETURN_IF_ERROR(new_pipe->add_operator(op, _parallel_instances));
1671
0
                cur_pipe = new_pipe;
1672
0
            } else {
1673
0
                op = std::make_shared<StreamingAggOperatorX>(pool, next_operator_id(), tnode,
1674
0
                                                             descs);
1675
0
                RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1676
0
            }
1677
0
        } else {
1678
            // create new pipeline to add query cache operator
1679
0
            PipelinePtr new_pipe;
1680
0
            if (need_create_cache_op) {
1681
0
                RETURN_IF_ERROR(create_query_cache_operator(new_pipe));
1682
0
                cache_op = op;
1683
0
            }
1684
1685
0
            if (enable_spill) {
1686
0
                op = std::make_shared<PartitionedAggSourceOperatorX>(pool, tnode,
1687
0
                                                                     next_operator_id(), descs);
1688
0
            } else {
1689
0
                op = std::make_shared<AggSourceOperatorX>(pool, tnode, next_operator_id(), descs);
1690
0
            }
1691
0
            if (need_create_cache_op) {
1692
0
                RETURN_IF_ERROR(cur_pipe->operators().front()->set_child(op));
1693
0
                RETURN_IF_ERROR(new_pipe->add_operator(op, _parallel_instances));
1694
0
                cur_pipe = new_pipe;
1695
0
            } else {
1696
0
                RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1697
0
            }
1698
1699
0
            const auto downstream_pipeline_id = cur_pipe->id();
1700
0
            if (!_dag.contains(downstream_pipeline_id)) {
1701
0
                _dag.insert({downstream_pipeline_id, {}});
1702
0
            }
1703
0
            cur_pipe = add_pipeline(cur_pipe);
1704
0
            _dag[downstream_pipeline_id].push_back(cur_pipe->id());
1705
1706
0
            if (enable_spill) {
1707
0
                sink_ops.push_back(std::make_shared<PartitionedAggSinkOperatorX>(
1708
0
                        pool, next_sink_operator_id(), op->operator_id(), tnode, descs));
1709
0
            } else {
1710
0
                sink_ops.push_back(std::make_shared<AggSinkOperatorX>(
1711
0
                        pool, next_sink_operator_id(), op->operator_id(), tnode, descs));
1712
0
            }
1713
0
            RETURN_IF_ERROR(cur_pipe->set_sink(sink_ops.back()));
1714
0
            RETURN_IF_ERROR(cur_pipe->sink()->init(tnode, _runtime_state.get()));
1715
0
        }
1716
0
        break;
1717
0
    }
1718
0
    case TPlanNodeType::BUCKETED_AGGREGATION_NODE: {
1719
0
        if (tnode.bucketed_agg_node.grouping_exprs.empty()) {
1720
0
            return Status::InternalError(
1721
0
                    "Bucketed aggregation node {} should not be used without group by keys",
1722
0
                    tnode.node_id);
1723
0
        }
1724
1725
        // Create source operator (goes on the current / downstream pipeline).
1726
0
        op = std::make_shared<BucketedAggSourceOperatorX>(pool, tnode, next_operator_id(), descs);
1727
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1728
1729
        // Create a new pipeline for the sink side.
1730
0
        const auto downstream_pipeline_id = cur_pipe->id();
1731
0
        if (!_dag.contains(downstream_pipeline_id)) {
1732
0
            _dag.insert({downstream_pipeline_id, {}});
1733
0
        }
1734
0
        cur_pipe = add_pipeline(cur_pipe);
1735
0
        _dag[downstream_pipeline_id].push_back(cur_pipe->id());
1736
1737
        // Create sink operator.
1738
0
        sink_ops.push_back(std::make_shared<BucketedAggSinkOperatorX>(
1739
0
                pool, next_sink_operator_id(), op->operator_id(), tnode, descs));
1740
0
        RETURN_IF_ERROR(cur_pipe->set_sink(sink_ops.back()));
1741
0
        RETURN_IF_ERROR(cur_pipe->sink()->init(tnode, _runtime_state.get()));
1742
1743
        // Pre-register a single shared state for ALL instances so that every
1744
        // sink instance writes its per-instance hash table into the same
1745
        // BucketedAggSharedState and every source instance can merge across
1746
        // all of them.
1747
0
        {
1748
0
            auto shared_state = BucketedAggSharedState::create_shared();
1749
0
            shared_state->id = op->operator_id();
1750
0
            shared_state->related_op_ids.insert(op->operator_id());
1751
1752
0
            for (int i = 0; i < _num_instances; i++) {
1753
0
                auto sink_dep = std::make_shared<Dependency>(op->operator_id(), op->node_id(),
1754
0
                                                             "BUCKETED_AGG_SINK_DEPENDENCY");
1755
0
                sink_dep->set_shared_state(shared_state.get());
1756
0
                shared_state->sink_deps.push_back(sink_dep);
1757
0
            }
1758
0
            shared_state->create_source_dependencies(_num_instances, op->operator_id(),
1759
0
                                                     op->node_id(), "BUCKETED_AGG_SOURCE");
1760
0
            _op_id_to_shared_state.insert(
1761
0
                    {op->operator_id(), {shared_state, shared_state->sink_deps}});
1762
0
        }
1763
0
        break;
1764
0
    }
1765
0
    case TPlanNodeType::HASH_JOIN_NODE: {
1766
0
        const auto is_broadcast_join = tnode.hash_join_node.__isset.is_broadcast_join &&
1767
0
                                       tnode.hash_join_node.is_broadcast_join;
1768
0
        const auto enable_spill = _runtime_state->enable_spill();
1769
0
        if (enable_spill && !is_broadcast_join) {
1770
0
            auto tnode_ = tnode;
1771
0
            tnode_.runtime_filters.clear();
1772
0
            auto inner_probe_operator =
1773
0
                    std::make_shared<HashJoinProbeOperatorX>(pool, tnode_, 0, descs);
1774
1775
            // probe side inner sink operator is used to build hash table on probe side when data is spilled.
1776
            // So here use `tnode_` which has no runtime filters.
1777
0
            auto probe_side_inner_sink_operator =
1778
0
                    std::make_shared<HashJoinBuildSinkOperatorX>(pool, 0, 0, tnode_, descs);
1779
1780
0
            RETURN_IF_ERROR(inner_probe_operator->init(tnode_, _runtime_state.get()));
1781
0
            RETURN_IF_ERROR(probe_side_inner_sink_operator->init(tnode_, _runtime_state.get()));
1782
1783
0
            auto probe_operator = std::make_shared<PartitionedHashJoinProbeOperatorX>(
1784
0
                    pool, tnode_, next_operator_id(), descs);
1785
0
            probe_operator->set_inner_operators(probe_side_inner_sink_operator,
1786
0
                                                inner_probe_operator);
1787
0
            op = std::move(probe_operator);
1788
0
            RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1789
1790
0
            const auto downstream_pipeline_id = cur_pipe->id();
1791
0
            if (!_dag.contains(downstream_pipeline_id)) {
1792
0
                _dag.insert({downstream_pipeline_id, {}});
1793
0
            }
1794
0
            PipelinePtr build_side_pipe = add_pipeline(cur_pipe);
1795
0
            _dag[downstream_pipeline_id].push_back(build_side_pipe->id());
1796
1797
0
            auto inner_sink_operator =
1798
0
                    std::make_shared<HashJoinBuildSinkOperatorX>(pool, 0, 0, tnode, descs);
1799
0
            auto sink_operator = std::make_shared<PartitionedHashJoinSinkOperatorX>(
1800
0
                    pool, next_sink_operator_id(), op->operator_id(), tnode_, descs);
1801
0
            RETURN_IF_ERROR(inner_sink_operator->init(tnode, _runtime_state.get()));
1802
1803
0
            sink_operator->set_inner_operators(inner_sink_operator, inner_probe_operator);
1804
0
            sink_ops.push_back(std::move(sink_operator));
1805
0
            RETURN_IF_ERROR(build_side_pipe->set_sink(sink_ops.back()));
1806
0
            RETURN_IF_ERROR(build_side_pipe->sink()->init(tnode_, _runtime_state.get()));
1807
1808
0
            _pipeline_parent_map.push(op->node_id(), cur_pipe);
1809
0
            _pipeline_parent_map.push(op->node_id(), build_side_pipe);
1810
0
        } else {
1811
0
            op = std::make_shared<HashJoinProbeOperatorX>(pool, tnode, next_operator_id(), descs);
1812
0
            RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1813
1814
0
            const auto downstream_pipeline_id = cur_pipe->id();
1815
0
            if (!_dag.contains(downstream_pipeline_id)) {
1816
0
                _dag.insert({downstream_pipeline_id, {}});
1817
0
            }
1818
0
            PipelinePtr build_side_pipe = add_pipeline(cur_pipe);
1819
0
            _dag[downstream_pipeline_id].push_back(build_side_pipe->id());
1820
1821
0
            sink_ops.push_back(std::make_shared<HashJoinBuildSinkOperatorX>(
1822
0
                    pool, next_sink_operator_id(), op->operator_id(), tnode, descs));
1823
0
            RETURN_IF_ERROR(build_side_pipe->set_sink(sink_ops.back()));
1824
0
            RETURN_IF_ERROR(build_side_pipe->sink()->init(tnode, _runtime_state.get()));
1825
1826
0
            _pipeline_parent_map.push(op->node_id(), cur_pipe);
1827
0
            _pipeline_parent_map.push(op->node_id(), build_side_pipe);
1828
0
        }
1829
0
        if (is_broadcast_join && _runtime_state->enable_share_hash_table_for_broadcast_join()) {
1830
0
            std::shared_ptr<HashJoinSharedState> shared_state =
1831
0
                    HashJoinSharedState::create_shared(_num_instances);
1832
0
            for (int i = 0; i < _num_instances; i++) {
1833
0
                auto sink_dep = std::make_shared<Dependency>(op->operator_id(), op->node_id(),
1834
0
                                                             "HASH_JOIN_BUILD_DEPENDENCY");
1835
0
                sink_dep->set_shared_state(shared_state.get());
1836
0
                shared_state->sink_deps.push_back(sink_dep);
1837
0
            }
1838
0
            shared_state->create_source_dependencies(_num_instances, op->operator_id(),
1839
0
                                                     op->node_id(), "HASH_JOIN_PROBE");
1840
0
            _op_id_to_shared_state.insert(
1841
0
                    {op->operator_id(), {shared_state, shared_state->sink_deps}});
1842
0
        }
1843
0
        break;
1844
0
    }
1845
0
    case TPlanNodeType::CROSS_JOIN_NODE: {
1846
0
        op = std::make_shared<NestedLoopJoinProbeOperatorX>(pool, tnode, next_operator_id(), descs);
1847
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1848
1849
0
        const auto downstream_pipeline_id = cur_pipe->id();
1850
0
        if (!_dag.contains(downstream_pipeline_id)) {
1851
0
            _dag.insert({downstream_pipeline_id, {}});
1852
0
        }
1853
0
        PipelinePtr build_side_pipe = add_pipeline(cur_pipe);
1854
0
        _dag[downstream_pipeline_id].push_back(build_side_pipe->id());
1855
1856
0
        sink_ops.push_back(std::make_shared<NestedLoopJoinBuildSinkOperatorX>(
1857
0
                pool, next_sink_operator_id(), op->operator_id(), tnode, descs));
1858
0
        RETURN_IF_ERROR(build_side_pipe->set_sink(sink_ops.back()));
1859
0
        RETURN_IF_ERROR(build_side_pipe->sink()->init(tnode, _runtime_state.get()));
1860
0
        _pipeline_parent_map.push(op->node_id(), cur_pipe);
1861
0
        _pipeline_parent_map.push(op->node_id(), build_side_pipe);
1862
0
        break;
1863
0
    }
1864
0
    case TPlanNodeType::UNION_NODE: {
1865
0
        int child_count = tnode.num_children;
1866
0
        op = std::make_shared<UnionSourceOperatorX>(pool, tnode, next_operator_id(), descs);
1867
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1868
1869
0
        const auto downstream_pipeline_id = cur_pipe->id();
1870
0
        if (!_dag.contains(downstream_pipeline_id)) {
1871
0
            _dag.insert({downstream_pipeline_id, {}});
1872
0
        }
1873
0
        for (int i = 0; i < child_count; i++) {
1874
0
            PipelinePtr build_side_pipe = add_pipeline(cur_pipe);
1875
0
            _dag[downstream_pipeline_id].push_back(build_side_pipe->id());
1876
0
            sink_ops.push_back(std::make_shared<UnionSinkOperatorX>(
1877
0
                    i, next_sink_operator_id(), op->operator_id(), pool, tnode, descs));
1878
0
            RETURN_IF_ERROR(build_side_pipe->set_sink(sink_ops.back()));
1879
0
            RETURN_IF_ERROR(build_side_pipe->sink()->init(tnode, _runtime_state.get()));
1880
            // preset children pipelines. if any pipeline found this as its father, will use the prepared pipeline to build.
1881
0
            _pipeline_parent_map.push(op->node_id(), build_side_pipe);
1882
0
        }
1883
0
        break;
1884
0
    }
1885
0
    case TPlanNodeType::SORT_NODE: {
1886
0
        const auto should_spill = _runtime_state->enable_spill() &&
1887
0
                                  tnode.sort_node.algorithm == TSortAlgorithm::FULL_SORT;
1888
0
        const bool use_local_merge =
1889
0
                tnode.sort_node.__isset.use_local_merge && tnode.sort_node.use_local_merge;
1890
0
        if (should_spill) {
1891
0
            op = std::make_shared<SpillSortSourceOperatorX>(pool, tnode, next_operator_id(), descs);
1892
0
        } else if (use_local_merge) {
1893
0
            op = std::make_shared<LocalMergeSortSourceOperatorX>(pool, tnode, next_operator_id(),
1894
0
                                                                 descs);
1895
0
        } else {
1896
0
            op = std::make_shared<SortSourceOperatorX>(pool, tnode, next_operator_id(), descs);
1897
0
        }
1898
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1899
1900
0
        const auto downstream_pipeline_id = cur_pipe->id();
1901
0
        if (!_dag.contains(downstream_pipeline_id)) {
1902
0
            _dag.insert({downstream_pipeline_id, {}});
1903
0
        }
1904
0
        cur_pipe = add_pipeline(cur_pipe);
1905
0
        _dag[downstream_pipeline_id].push_back(cur_pipe->id());
1906
1907
0
        if (should_spill) {
1908
0
            sink_ops.push_back(std::make_shared<SpillSortSinkOperatorX>(
1909
0
                    pool, next_sink_operator_id(), op->operator_id(), tnode, descs));
1910
0
        } else {
1911
0
            sink_ops.push_back(std::make_shared<SortSinkOperatorX>(
1912
0
                    pool, next_sink_operator_id(), op->operator_id(), tnode, descs));
1913
0
        }
1914
0
        RETURN_IF_ERROR(cur_pipe->set_sink(sink_ops.back()));
1915
0
        RETURN_IF_ERROR(cur_pipe->sink()->init(tnode, _runtime_state.get()));
1916
0
        break;
1917
0
    }
1918
0
    case TPlanNodeType::PARTITION_SORT_NODE: {
1919
0
        op = std::make_shared<PartitionSortSourceOperatorX>(pool, tnode, next_operator_id(), descs);
1920
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1921
1922
0
        const auto downstream_pipeline_id = cur_pipe->id();
1923
0
        if (!_dag.contains(downstream_pipeline_id)) {
1924
0
            _dag.insert({downstream_pipeline_id, {}});
1925
0
        }
1926
0
        cur_pipe = add_pipeline(cur_pipe);
1927
0
        _dag[downstream_pipeline_id].push_back(cur_pipe->id());
1928
1929
0
        sink_ops.push_back(std::make_shared<PartitionSortSinkOperatorX>(
1930
0
                pool, next_sink_operator_id(), op->operator_id(), tnode, descs));
1931
0
        RETURN_IF_ERROR(cur_pipe->set_sink(sink_ops.back()));
1932
0
        RETURN_IF_ERROR(cur_pipe->sink()->init(tnode, _runtime_state.get()));
1933
0
        break;
1934
0
    }
1935
0
    case TPlanNodeType::ANALYTIC_EVAL_NODE: {
1936
0
        op = std::make_shared<AnalyticSourceOperatorX>(pool, tnode, next_operator_id(), descs);
1937
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1938
1939
0
        const auto downstream_pipeline_id = cur_pipe->id();
1940
0
        if (!_dag.contains(downstream_pipeline_id)) {
1941
0
            _dag.insert({downstream_pipeline_id, {}});
1942
0
        }
1943
0
        cur_pipe = add_pipeline(cur_pipe);
1944
0
        _dag[downstream_pipeline_id].push_back(cur_pipe->id());
1945
1946
0
        sink_ops.push_back(std::make_shared<AnalyticSinkOperatorX>(
1947
0
                pool, next_sink_operator_id(), op->operator_id(), tnode, descs));
1948
0
        RETURN_IF_ERROR(cur_pipe->set_sink(sink_ops.back()));
1949
0
        RETURN_IF_ERROR(cur_pipe->sink()->init(tnode, _runtime_state.get()));
1950
0
        break;
1951
0
    }
1952
0
    case TPlanNodeType::MATERIALIZATION_NODE: {
1953
0
        op = std::make_shared<MaterializationOperator>(pool, tnode, next_operator_id(), descs);
1954
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1955
0
        break;
1956
0
    }
1957
0
    case TPlanNodeType::INTERSECT_NODE: {
1958
0
        RETURN_IF_ERROR(_build_operators_for_set_operation_node<true>(pool, tnode, descs, op,
1959
0
                                                                      cur_pipe, sink_ops));
1960
0
        break;
1961
0
    }
1962
0
    case TPlanNodeType::EXCEPT_NODE: {
1963
0
        RETURN_IF_ERROR(_build_operators_for_set_operation_node<false>(pool, tnode, descs, op,
1964
0
                                                                       cur_pipe, sink_ops));
1965
0
        break;
1966
0
    }
1967
0
    case TPlanNodeType::REPEAT_NODE: {
1968
0
        op = std::make_shared<RepeatOperatorX>(pool, tnode, next_operator_id(), descs);
1969
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1970
0
        break;
1971
0
    }
1972
0
    case TPlanNodeType::TABLE_FUNCTION_NODE: {
1973
0
        op = std::make_shared<TableFunctionOperatorX>(pool, tnode, next_operator_id(), descs);
1974
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1975
0
        break;
1976
0
    }
1977
0
    case TPlanNodeType::ASSERT_NUM_ROWS_NODE: {
1978
0
        op = std::make_shared<AssertNumRowsOperatorX>(pool, tnode, next_operator_id(), descs);
1979
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1980
0
        break;
1981
0
    }
1982
0
    case TPlanNodeType::EMPTY_SET_NODE: {
1983
0
        op = std::make_shared<EmptySetSourceOperatorX>(pool, tnode, next_operator_id(), descs);
1984
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1985
0
        break;
1986
0
    }
1987
0
    case TPlanNodeType::DATA_GEN_SCAN_NODE: {
1988
0
        op = std::make_shared<DataGenSourceOperatorX>(pool, tnode, next_operator_id(), descs);
1989
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1990
0
        fe_with_old_version = !tnode.__isset.is_serial_operator;
1991
0
        break;
1992
0
    }
1993
0
    case TPlanNodeType::SCHEMA_SCAN_NODE: {
1994
0
        op = std::make_shared<SchemaScanOperatorX>(pool, tnode, next_operator_id(), descs);
1995
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
1996
0
        break;
1997
0
    }
1998
0
    case TPlanNodeType::META_SCAN_NODE: {
1999
0
        op = std::make_shared<MetaScanOperatorX>(pool, tnode, next_operator_id(), descs);
2000
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
2001
0
        break;
2002
0
    }
2003
0
    case TPlanNodeType::SELECT_NODE: {
2004
0
        op = std::make_shared<SelectOperatorX>(pool, tnode, next_operator_id(), descs);
2005
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
2006
0
        break;
2007
0
    }
2008
0
    case TPlanNodeType::REC_CTE_NODE: {
2009
0
        op = std::make_shared<RecCTESourceOperatorX>(pool, tnode, next_operator_id(), descs);
2010
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
2011
2012
0
        const auto downstream_pipeline_id = cur_pipe->id();
2013
0
        if (!_dag.contains(downstream_pipeline_id)) {
2014
0
            _dag.insert({downstream_pipeline_id, {}});
2015
0
        }
2016
2017
0
        PipelinePtr anchor_side_pipe = add_pipeline(cur_pipe);
2018
0
        _dag[downstream_pipeline_id].push_back(anchor_side_pipe->id());
2019
2020
0
        DataSinkOperatorPtr anchor_sink;
2021
0
        anchor_sink = std::make_shared<RecCTEAnchorSinkOperatorX>(next_sink_operator_id(),
2022
0
                                                                  op->operator_id(), tnode, descs);
2023
0
        RETURN_IF_ERROR(anchor_side_pipe->set_sink(anchor_sink));
2024
0
        RETURN_IF_ERROR(anchor_side_pipe->sink()->init(tnode, _runtime_state.get()));
2025
0
        _pipeline_parent_map.push(op->node_id(), anchor_side_pipe);
2026
2027
0
        PipelinePtr rec_side_pipe = add_pipeline(cur_pipe);
2028
0
        _dag[downstream_pipeline_id].push_back(rec_side_pipe->id());
2029
2030
0
        DataSinkOperatorPtr rec_sink;
2031
0
        rec_sink = std::make_shared<RecCTESinkOperatorX>(next_sink_operator_id(), op->operator_id(),
2032
0
                                                         tnode, descs);
2033
0
        RETURN_IF_ERROR(rec_side_pipe->set_sink(rec_sink));
2034
0
        RETURN_IF_ERROR(rec_side_pipe->sink()->init(tnode, _runtime_state.get()));
2035
0
        _pipeline_parent_map.push(op->node_id(), rec_side_pipe);
2036
2037
0
        break;
2038
0
    }
2039
0
    case TPlanNodeType::REC_CTE_SCAN_NODE: {
2040
0
        op = std::make_shared<RecCTEScanOperatorX>(pool, tnode, next_operator_id(), descs);
2041
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
2042
0
        break;
2043
0
    }
2044
0
    case TPlanNodeType::LOCAL_EXCHANGE_NODE: {
2045
0
        op = std::make_shared<LocalExchangeSourceOperatorX>(pool, tnode, next_operator_id(), descs);
2046
        // The downstream pipeline (containing LocalExchangeSource) must have
2047
        // _num_instances tasks — matching BE-native _inherit_pipeline_properties
2048
        // which sets pipe_with_source.set_num_tasks(_num_instances).
2049
        // Without this, when the parent pipeline was reduced by a serial operator
2050
        // (e.g., serial Exchange with use_serial_exchange=true, or UNPARTITIONED
2051
        // Exchange), the downstream inherits the reduced num_tasks via
2052
        // add_pipeline(parent).  The deferred exchanger creates _num_instances
2053
        // channels but only fewer source tasks initialize mem_counters — the
2054
        // sink round-robins to all channels and crashes on uninitialized ones.
2055
0
        RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
2056
        // Restore downstream pipeline's num_tasks (mirroring _inherit_pipeline_properties:
2057
        // downstream keeps _num_instances, upstream gets the serial/reduced count)
2058
0
        cur_pipe->set_num_tasks(_num_instances);
2059
2060
0
        const auto downstream_pipeline_id = cur_pipe->id();
2061
0
        if (!_dag.contains(downstream_pipeline_id)) {
2062
0
            _dag.insert({downstream_pipeline_id, {}});
2063
0
        }
2064
0
        cur_pipe = add_pipeline(cur_pipe);
2065
        // If this local exchange was inserted because of a serial scan (is_serial_operator),
2066
        // the upstream pipeline (cur_pipe) should have num_tasks=1 (only 1 scan task).
2067
        // We set this now so the exchanger is created with the correct sender count.
2068
        // Child operators added later (serial scan) will also set num_tasks=1, which is
2069
        // consistent with this.
2070
0
        if (op->is_serial_operator() && _parallel_instances > 0) {
2071
0
            cur_pipe->set_num_tasks(_parallel_instances);
2072
0
        }
2073
0
        _dag[downstream_pipeline_id].push_back(cur_pipe->id());
2074
0
        int num_partitions = 0;
2075
0
        std::map<int, int> shuffle_id_to_instance_idx;
2076
0
        auto partition_type = tnode.local_exchange_node.partition_type;
2077
0
        switch (partition_type) {
2078
0
        case TLocalPartitionType::BUCKET_HASH_SHUFFLE:
2079
0
            num_partitions = _params.num_buckets;
2080
0
            shuffle_id_to_instance_idx = _params.bucket_seq_to_instance_idx;
2081
0
            break;
2082
0
        case TLocalPartitionType::LOCAL_EXECUTION_HASH_SHUFFLE:
2083
0
            for (int i = 0; i < _num_instances; i++) {
2084
0
                shuffle_id_to_instance_idx[i] = i;
2085
0
            }
2086
0
            num_partitions = _num_instances;
2087
0
            break;
2088
0
        case TLocalPartitionType::GLOBAL_EXECUTION_HASH_SHUFFLE:
2089
0
            num_partitions = _total_instances;
2090
0
            shuffle_id_to_instance_idx = _params.shuffle_idx_to_instance_idx;
2091
0
            break;
2092
0
        default:
2093
0
            break;
2094
0
        }
2095
0
        auto local_exchange_id = op->operator_id();
2096
0
        auto sink_id = next_sink_operator_id();
2097
0
        DataSinkOperatorPtr sink = std::make_shared<LocalExchangeSinkOperatorX>(
2098
0
                sink_id, local_exchange_id, tnode, num_partitions, shuffle_id_to_instance_idx);
2099
0
        sink_ops.push_back(sink);
2100
0
        RETURN_IF_ERROR(cur_pipe->set_sink(sink));
2101
0
        RETURN_IF_ERROR(cur_pipe->sink()->init(tnode, _runtime_state.get()));
2102
2103
        // For FE-planned local exchange, we need to:
2104
        // 1. Initialize the partitioner for hash shuffle types
2105
        // 2. Defer exchanger creation until after the full plan tree is built
2106
        //    (child operators like serial ExchangeNode may change cur_pipe->num_tasks())
2107
        // 3. Register shared state so pipeline tasks can find it
2108
0
        RETURN_IF_ERROR(static_cast<LocalExchangeSinkOperatorX*>(cur_pipe->sink())
2109
0
                                ->init_partitioner(_runtime_state.get()));
2110
2111
0
        int free_blocks_limit =
2112
0
                _runtime_state->query_options().__isset.local_exchange_free_blocks_limit
2113
0
                        ? cast_set<int>(
2114
0
                                  _runtime_state->query_options().local_exchange_free_blocks_limit)
2115
0
                        : 0;
2116
0
        auto shared_state = LocalExchangeSharedState::create_shared(_num_instances);
2117
0
        shared_state->create_source_dependencies(_num_instances, local_exchange_id,
2118
0
                                                 local_exchange_id, "LOCAL_EXCHANGE_OPERATOR");
2119
0
        shared_state->create_sink_dependency(sink_id, local_exchange_id, "LOCAL_EXCHANGE_SINK");
2120
0
        _op_id_to_shared_state.insert({local_exchange_id, {shared_state, shared_state->sink_deps}});
2121
        // Defer exchanger creation: sender count depends on final upstream num_tasks
2122
0
        _deferred_exchangers.push_back({shared_state, cur_pipe, partition_type, num_partitions,
2123
0
                                        free_blocks_limit, local_exchange_id, sink_id});
2124
0
        break;
2125
0
    }
2126
0
    default:
2127
0
        return Status::InternalError("Unsupported exec type in pipeline: {}",
2128
0
                                     print_plan_node_type(tnode.node_type));
2129
2
    }
2130
0
    if (_params.__isset.parallel_instances && fe_with_old_version) {
2131
0
        cur_pipe->set_num_tasks(_params.parallel_instances);
2132
0
        op->set_serial_operator();
2133
0
    }
2134
2135
0
    return Status::OK();
2136
2
}
2137
// NOLINTEND(readability-function-cognitive-complexity)
2138
// NOLINTEND(readability-function-size)
2139
2140
template <bool is_intersect>
2141
Status PipelineFragmentContext::_build_operators_for_set_operation_node(
2142
        ObjectPool* pool, const TPlanNode& tnode, const DescriptorTbl& descs, OperatorPtr& op,
2143
0
        PipelinePtr& cur_pipe, std::vector<DataSinkOperatorPtr>& sink_ops) {
2144
0
    op.reset(new SetSourceOperatorX<is_intersect>(pool, tnode, next_operator_id(), descs));
2145
0
    RETURN_IF_ERROR(cur_pipe->add_operator(op, _parallel_instances));
2146
2147
0
    const auto downstream_pipeline_id = cur_pipe->id();
2148
0
    if (!_dag.contains(downstream_pipeline_id)) {
2149
0
        _dag.insert({downstream_pipeline_id, {}});
2150
0
    }
2151
2152
0
    for (int child_id = 0; child_id < tnode.num_children; child_id++) {
2153
0
        PipelinePtr probe_side_pipe = add_pipeline(cur_pipe);
2154
0
        _dag[downstream_pipeline_id].push_back(probe_side_pipe->id());
2155
2156
0
        if (child_id == 0) {
2157
0
            sink_ops.push_back(std::make_shared<SetSinkOperatorX<is_intersect>>(
2158
0
                    child_id, next_sink_operator_id(), op->operator_id(), pool, tnode, descs));
2159
0
        } else {
2160
0
            sink_ops.push_back(std::make_shared<SetProbeSinkOperatorX<is_intersect>>(
2161
0
                    child_id, next_sink_operator_id(), op->operator_id(), pool, tnode, descs));
2162
0
        }
2163
0
        RETURN_IF_ERROR(probe_side_pipe->set_sink(sink_ops.back()));
2164
0
        RETURN_IF_ERROR(probe_side_pipe->sink()->init(tnode, _runtime_state.get()));
2165
        // prepare children pipelines. if any pipeline found this as its father, will use the prepared pipeline to build.
2166
0
        _pipeline_parent_map.push(op->node_id(), probe_side_pipe);
2167
0
    }
2168
2169
0
    return Status::OK();
2170
0
}
Unexecuted instantiation: _ZN5doris23PipelineFragmentContext39_build_operators_for_set_operation_nodeILb1EEENS_6StatusEPNS_10ObjectPoolERKNS_9TPlanNodeERKNS_13DescriptorTblERSt10shared_ptrINS_13OperatorXBaseEERSB_INS_8PipelineEERSt6vectorISB_INS_21DataSinkOperatorXBaseEESaISK_EE
Unexecuted instantiation: _ZN5doris23PipelineFragmentContext39_build_operators_for_set_operation_nodeILb0EEENS_6StatusEPNS_10ObjectPoolERKNS_9TPlanNodeERKNS_13DescriptorTblERSt10shared_ptrINS_13OperatorXBaseEERSB_INS_8PipelineEERSt6vectorISB_INS_21DataSinkOperatorXBaseEESaISK_EE
2171
2172
0
Status PipelineFragmentContext::submit() {
2173
0
    if (_submitted) {
2174
0
        return Status::InternalError("submitted");
2175
0
    }
2176
0
    _submitted = true;
2177
2178
0
    int submit_tasks = 0;
2179
0
    Status st;
2180
0
    auto* scheduler = _query_ctx->get_pipe_exec_scheduler();
2181
0
    for (auto& task : _tasks) {
2182
0
        for (auto& t : task) {
2183
0
            st = scheduler->submit(t.first);
2184
0
            DBUG_EXECUTE_IF("PipelineFragmentContext.submit.failed",
2185
0
                            { st = Status::Aborted("PipelineFragmentContext.submit.failed"); });
2186
0
            if (!st) {
2187
0
                cancel(Status::InternalError("submit context to executor fail"));
2188
0
                std::lock_guard<std::mutex> l(_task_mutex);
2189
0
                _total_tasks = submit_tasks;
2190
0
                break;
2191
0
            }
2192
0
            submit_tasks++;
2193
0
        }
2194
0
    }
2195
0
    if (!st.ok()) {
2196
0
        bool need_remove = false;
2197
0
        {
2198
0
            std::lock_guard<std::mutex> l(_task_mutex);
2199
0
            if (_closed_tasks >= _total_tasks) {
2200
0
                need_remove = _close_fragment_instance();
2201
0
            }
2202
0
        }
2203
        // Call remove_pipeline_context() outside _task_mutex to avoid ABBA deadlock.
2204
0
        if (need_remove) {
2205
0
            _exec_env->fragment_mgr()->remove_pipeline_context({_query_id, _fragment_id});
2206
0
        }
2207
0
        return Status::InternalError("Submit pipeline failed. err = {}, BE: {}", st.to_string(),
2208
0
                                     BackendOptions::get_localhost());
2209
0
    } else {
2210
0
        return st;
2211
0
    }
2212
0
}
2213
2214
0
void PipelineFragmentContext::print_profile(const std::string& extra_info) {
2215
0
    if (_runtime_state->enable_profile()) {
2216
0
        std::stringstream ss;
2217
0
        for (auto runtime_profile_ptr : _runtime_state->pipeline_id_to_profile()) {
2218
0
            runtime_profile_ptr->pretty_print(&ss);
2219
0
        }
2220
2221
0
        if (_runtime_state->load_channel_profile()) {
2222
0
            _runtime_state->load_channel_profile()->pretty_print(&ss);
2223
0
        }
2224
2225
0
        auto profile_str =
2226
0
                fmt::format("Query {} fragment {} {}, profile, {}", print_id(this->_query_id),
2227
0
                            this->_fragment_id, extra_info, ss.str());
2228
0
        LOG_LONG_STRING(INFO, profile_str);
2229
0
    }
2230
0
}
2231
// If all pipeline tasks binded to the fragment instance are finished, then we could
2232
// close the fragment instance.
2233
// Returns true if the caller should call remove_pipeline_context() **after** releasing
2234
// _task_mutex. We must not call remove_pipeline_context() here because it acquires
2235
// _pipeline_map's shard lock, and this function is called while _task_mutex is held.
2236
// Acquiring _pipeline_map while holding _task_mutex creates an ABBA deadlock with
2237
// dump_pipeline_tasks(), which acquires _pipeline_map first and then _task_mutex
2238
// (via debug_string()).
2239
0
bool PipelineFragmentContext::_close_fragment_instance() {
2240
0
    if (_is_fragment_instance_closed) {
2241
0
        return false;
2242
0
    }
2243
0
    Defer defer_op {[&]() { _is_fragment_instance_closed = true; }};
2244
0
    _fragment_level_profile->total_time_counter()->update(_fragment_watcher.elapsed_time());
2245
0
    if (!_need_notify_close) {
2246
0
        auto st = send_report(true);
2247
0
        if (!st) {
2248
0
            LOG(WARNING) << fmt::format("Failed to send report for query {}, fragment {}: {}",
2249
0
                                        print_id(_query_id), _fragment_id, st.to_string());
2250
0
        }
2251
0
    }
2252
    // Print profile content in info log is a tempoeray solution for stream load and external_connector.
2253
    // Since stream load does not have someting like coordinator on FE, so
2254
    // backend can not report profile to FE, ant its profile can not be shown
2255
    // in the same way with other query. So we print the profile content to info log.
2256
2257
0
    if (_runtime_state->enable_profile() &&
2258
0
        (_query_ctx->get_query_source() == QuerySource::STREAM_LOAD ||
2259
0
         _query_ctx->get_query_source() == QuerySource::EXTERNAL_CONNECTOR ||
2260
0
         _query_ctx->get_query_source() == QuerySource::GROUP_COMMIT_LOAD)) {
2261
0
        std::stringstream ss;
2262
        // Compute the _local_time_percent before pretty_print the runtime_profile
2263
        // Before add this operation, the print out like that:
2264
        // UNION_NODE (id=0):(Active: 56.720us, non-child: 00.00%)
2265
        // After add the operation, the print out like that:
2266
        // UNION_NODE (id=0):(Active: 56.720us, non-child: 82.53%)
2267
        // We can easily know the exec node execute time without child time consumed.
2268
0
        for (auto runtime_profile_ptr : _runtime_state->pipeline_id_to_profile()) {
2269
0
            runtime_profile_ptr->pretty_print(&ss);
2270
0
        }
2271
2272
0
        if (_runtime_state->load_channel_profile()) {
2273
0
            _runtime_state->load_channel_profile()->pretty_print(&ss);
2274
0
        }
2275
2276
0
        LOG_INFO("Query {} fragment {} profile:\n {}", print_id(_query_id), _fragment_id, ss.str());
2277
0
    }
2278
2279
0
    if (_query_ctx->enable_profile()) {
2280
0
        _query_ctx->add_fragment_profile(_fragment_id, collect_realtime_profile(),
2281
0
                                         collect_realtime_load_channel_profile());
2282
0
    }
2283
2284
    // Return whether the caller needs to remove from the pipeline map.
2285
    // The caller must do this after releasing _task_mutex.
2286
0
    return !_need_notify_close;
2287
0
}
2288
2289
1
void PipelineFragmentContext::decrement_running_task(PipelineId pipeline_id) {
2290
    // If all tasks of this pipeline has been closed, upstream tasks is never needed, and we just make those runnable here
2291
1
    DCHECK(_pip_id_to_pipeline.contains(pipeline_id));
2292
1
    if (_pip_id_to_pipeline[pipeline_id]->close_task()) {
2293
1
        if (_dag.contains(pipeline_id)) {
2294
0
            for (auto dep : _dag[pipeline_id]) {
2295
0
                _pip_id_to_pipeline[dep]->make_all_runnable(pipeline_id);
2296
0
            }
2297
0
        }
2298
1
    }
2299
1
    bool need_remove = false;
2300
1
    {
2301
1
        std::lock_guard<std::mutex> l(_task_mutex);
2302
1
        ++_closed_tasks;
2303
        // Update query-level finished task progress in real time.
2304
1
        _query_ctx->inc_finished_task_num();
2305
1
        if (_closed_tasks >= _total_tasks) {
2306
0
            need_remove = _close_fragment_instance();
2307
0
        }
2308
1
    }
2309
    // Call remove_pipeline_context() outside _task_mutex to avoid ABBA deadlock.
2310
1
    if (need_remove) {
2311
0
        _exec_env->fragment_mgr()->remove_pipeline_context({_query_id, _fragment_id});
2312
0
    }
2313
1
}
2314
2315
1
std::string PipelineFragmentContext::get_load_error_url() {
2316
1
    if (const auto& str = _runtime_state->get_error_log_file_path(); !str.empty()) {
2317
0
        return to_load_error_http_path(str);
2318
0
    }
2319
1
    for (auto& tasks : _tasks) {
2320
0
        for (auto& task : tasks) {
2321
0
            if (const auto& str = task.second->get_error_log_file_path(); !str.empty()) {
2322
0
                return to_load_error_http_path(str);
2323
0
            }
2324
0
        }
2325
0
    }
2326
1
    return "";
2327
1
}
2328
2329
1
std::string PipelineFragmentContext::get_first_error_msg() {
2330
1
    if (const auto& str = _runtime_state->get_first_error_msg(); !str.empty()) {
2331
0
        return str;
2332
0
    }
2333
1
    for (auto& tasks : _tasks) {
2334
0
        for (auto& task : tasks) {
2335
0
            if (const auto& str = task.second->get_first_error_msg(); !str.empty()) {
2336
0
                return str;
2337
0
            }
2338
0
        }
2339
0
    }
2340
1
    return "";
2341
1
}
2342
2343
0
std::string PipelineFragmentContext::_to_http_path(const std::string& file_name) const {
2344
0
    std::stringstream url;
2345
0
    url << "http://" << BackendOptions::get_localhost() << ":" << config::webserver_port
2346
0
        << "/api/_download_load?"
2347
0
        << "token=" << _exec_env->token() << "&file=" << file_name;
2348
0
    return url.str();
2349
0
}
2350
2351
0
void PipelineFragmentContext::_coordinator_callback(const ReportStatusRequest& req) {
2352
0
    DBUG_EXECUTE_IF("FragmentMgr::coordinator_callback.report_delay", {
2353
0
        int random_seconds = req.status.is<ErrorCode::DATA_QUALITY_ERROR>() ? 8 : 2;
2354
0
        LOG_INFO("sleep : ").tag("time", random_seconds).tag("query_id", print_id(req.query_id));
2355
0
        std::this_thread::sleep_for(std::chrono::seconds(random_seconds));
2356
0
        LOG_INFO("sleep done").tag("query_id", print_id(req.query_id));
2357
0
    });
2358
2359
0
    DCHECK(req.status.ok() || req.done); // if !status.ok() => done
2360
0
    if (req.coord_addr.hostname == "external") {
2361
        // External query (flink/spark read tablets) not need to report to FE.
2362
0
        return;
2363
0
    }
2364
0
    int callback_retries = 10;
2365
0
    const int sleep_ms = 1000;
2366
0
    Status exec_status = req.status;
2367
0
    Status coord_status;
2368
0
    std::unique_ptr<FrontendServiceConnection> coord = nullptr;
2369
0
    do {
2370
0
        coord = std::make_unique<FrontendServiceConnection>(_exec_env->frontend_client_cache(),
2371
0
                                                            req.coord_addr, &coord_status);
2372
0
        if (!coord_status.ok()) {
2373
0
            std::this_thread::sleep_for(std::chrono::milliseconds(sleep_ms));
2374
0
        }
2375
0
    } while (!coord_status.ok() && callback_retries-- > 0);
2376
2377
0
    if (!coord_status.ok()) {
2378
0
        UniqueId uid(req.query_id.hi, req.query_id.lo);
2379
0
        static_cast<void>(req.cancel_fn(Status::InternalError(
2380
0
                "query_id: {}, couldn't get a client for {}, reason is {}", uid.to_string(),
2381
0
                PrintThriftNetworkAddress(req.coord_addr), coord_status.to_string())));
2382
0
        return;
2383
0
    }
2384
2385
0
    TReportExecStatusParams params;
2386
0
    params.protocol_version = FrontendServiceVersion::V1;
2387
0
    params.__set_query_id(req.query_id);
2388
0
    params.__set_backend_num(req.backend_num);
2389
0
    params.__set_fragment_instance_id(req.fragment_instance_id);
2390
0
    params.__set_fragment_id(req.fragment_id);
2391
0
    params.__set_status(exec_status.to_thrift());
2392
0
    params.__set_done(req.done);
2393
0
    params.__set_query_type(req.runtime_state->query_type());
2394
0
    params.__isset.profile = false;
2395
2396
0
    DCHECK(req.runtime_state != nullptr);
2397
2398
0
    if (req.runtime_state->query_type() == TQueryType::LOAD) {
2399
0
        params.__set_loaded_rows(req.runtime_state->num_rows_load_total());
2400
0
        params.__set_loaded_bytes(req.runtime_state->num_bytes_load_total());
2401
0
    } else {
2402
0
        DCHECK(!req.runtime_states.empty());
2403
0
        if (!req.runtime_state->output_files().empty()) {
2404
0
            params.__isset.delta_urls = true;
2405
0
            for (auto& it : req.runtime_state->output_files()) {
2406
0
                params.delta_urls.push_back(_to_http_path(it));
2407
0
            }
2408
0
        }
2409
0
        if (!params.delta_urls.empty()) {
2410
0
            params.__isset.delta_urls = true;
2411
0
        }
2412
0
    }
2413
2414
0
    static std::string s_dpp_normal_all = "dpp.norm.ALL";
2415
0
    static std::string s_dpp_abnormal_all = "dpp.abnorm.ALL";
2416
0
    static std::string s_unselected_rows = "unselected.rows";
2417
0
    int64_t num_rows_load_success = 0;
2418
0
    int64_t num_rows_load_filtered = 0;
2419
0
    int64_t num_rows_load_unselected = 0;
2420
0
    if (req.runtime_state->num_rows_load_total() > 0 ||
2421
0
        req.runtime_state->num_rows_load_filtered() > 0 ||
2422
0
        req.runtime_state->num_finished_range() > 0) {
2423
0
        params.__isset.load_counters = true;
2424
2425
0
        num_rows_load_success = req.runtime_state->num_rows_load_success();
2426
0
        num_rows_load_filtered = req.runtime_state->num_rows_load_filtered();
2427
0
        num_rows_load_unselected = req.runtime_state->num_rows_load_unselected();
2428
0
        params.__isset.fragment_instance_reports = true;
2429
0
        TFragmentInstanceReport t;
2430
0
        t.__set_fragment_instance_id(req.runtime_state->fragment_instance_id());
2431
0
        t.__set_num_finished_range(cast_set<int>(req.runtime_state->num_finished_range()));
2432
0
        t.__set_loaded_rows(req.runtime_state->num_rows_load_total());
2433
0
        t.__set_loaded_bytes(req.runtime_state->num_bytes_load_total());
2434
0
        params.fragment_instance_reports.push_back(t);
2435
0
    } else if (!req.runtime_states.empty()) {
2436
0
        for (auto* rs : req.runtime_states) {
2437
0
            if (rs->num_rows_load_total() > 0 || rs->num_rows_load_filtered() > 0 ||
2438
0
                rs->num_finished_range() > 0) {
2439
0
                params.__isset.load_counters = true;
2440
0
                num_rows_load_success += rs->num_rows_load_success();
2441
0
                num_rows_load_filtered += rs->num_rows_load_filtered();
2442
0
                num_rows_load_unselected += rs->num_rows_load_unselected();
2443
0
                params.__isset.fragment_instance_reports = true;
2444
0
                TFragmentInstanceReport t;
2445
0
                t.__set_fragment_instance_id(rs->fragment_instance_id());
2446
0
                t.__set_num_finished_range(cast_set<int>(rs->num_finished_range()));
2447
0
                t.__set_loaded_rows(rs->num_rows_load_total());
2448
0
                t.__set_loaded_bytes(rs->num_bytes_load_total());
2449
0
                params.fragment_instance_reports.push_back(t);
2450
0
            }
2451
0
        }
2452
0
    }
2453
0
    params.load_counters.emplace(s_dpp_normal_all, std::to_string(num_rows_load_success));
2454
0
    params.load_counters.emplace(s_dpp_abnormal_all, std::to_string(num_rows_load_filtered));
2455
0
    params.load_counters.emplace(s_unselected_rows, std::to_string(num_rows_load_unselected));
2456
2457
0
    if (!req.load_error_url.empty()) {
2458
0
        params.__set_tracking_url(req.load_error_url);
2459
0
    }
2460
0
    if (!req.first_error_msg.empty()) {
2461
0
        params.__set_first_error_msg(req.first_error_msg);
2462
0
    }
2463
0
    for (auto* rs : req.runtime_states) {
2464
0
        if (rs->wal_id() > 0) {
2465
0
            params.__set_txn_id(rs->wal_id());
2466
0
            params.__set_label(rs->import_label());
2467
0
        }
2468
0
    }
2469
0
    if (!req.runtime_state->export_output_files().empty()) {
2470
0
        params.__isset.export_files = true;
2471
0
        params.export_files = req.runtime_state->export_output_files();
2472
0
    } else if (!req.runtime_states.empty()) {
2473
0
        for (auto* rs : req.runtime_states) {
2474
0
            if (!rs->export_output_files().empty()) {
2475
0
                params.__isset.export_files = true;
2476
0
                params.export_files.insert(params.export_files.end(),
2477
0
                                           rs->export_output_files().begin(),
2478
0
                                           rs->export_output_files().end());
2479
0
            }
2480
0
        }
2481
0
    }
2482
0
    if (auto tci = req.runtime_state->tablet_commit_infos(); !tci.empty()) {
2483
0
        params.__isset.commitInfos = true;
2484
0
        params.commitInfos.insert(params.commitInfos.end(), tci.begin(), tci.end());
2485
0
    } else if (!req.runtime_states.empty()) {
2486
0
        for (auto* rs : req.runtime_states) {
2487
0
            if (auto rs_tci = rs->tablet_commit_infos(); !rs_tci.empty()) {
2488
0
                params.__isset.commitInfos = true;
2489
0
                params.commitInfos.insert(params.commitInfos.end(), rs_tci.begin(), rs_tci.end());
2490
0
            }
2491
0
        }
2492
0
    }
2493
0
    if (auto eti = req.runtime_state->error_tablet_infos(); !eti.empty()) {
2494
0
        params.__isset.errorTabletInfos = true;
2495
0
        params.errorTabletInfos.insert(params.errorTabletInfos.end(), eti.begin(), eti.end());
2496
0
    } else if (!req.runtime_states.empty()) {
2497
0
        for (auto* rs : req.runtime_states) {
2498
0
            if (auto rs_eti = rs->error_tablet_infos(); !rs_eti.empty()) {
2499
0
                params.__isset.errorTabletInfos = true;
2500
0
                params.errorTabletInfos.insert(params.errorTabletInfos.end(), rs_eti.begin(),
2501
0
                                               rs_eti.end());
2502
0
            }
2503
0
        }
2504
0
    }
2505
0
    if (auto hpu = req.runtime_state->hive_partition_updates(); !hpu.empty()) {
2506
0
        params.__isset.hive_partition_updates = true;
2507
0
        params.hive_partition_updates.insert(params.hive_partition_updates.end(), hpu.begin(),
2508
0
                                             hpu.end());
2509
0
    } else if (!req.runtime_states.empty()) {
2510
0
        for (auto* rs : req.runtime_states) {
2511
0
            if (auto rs_hpu = rs->hive_partition_updates(); !rs_hpu.empty()) {
2512
0
                params.__isset.hive_partition_updates = true;
2513
0
                params.hive_partition_updates.insert(params.hive_partition_updates.end(),
2514
0
                                                     rs_hpu.begin(), rs_hpu.end());
2515
0
            }
2516
0
        }
2517
0
    }
2518
0
    if (auto icd = req.runtime_state->iceberg_commit_datas(); !icd.empty()) {
2519
0
        params.__isset.iceberg_commit_datas = true;
2520
0
        params.iceberg_commit_datas.insert(params.iceberg_commit_datas.end(), icd.begin(),
2521
0
                                           icd.end());
2522
0
    } else if (!req.runtime_states.empty()) {
2523
0
        for (auto* rs : req.runtime_states) {
2524
0
            if (auto rs_icd = rs->iceberg_commit_datas(); !rs_icd.empty()) {
2525
0
                params.__isset.iceberg_commit_datas = true;
2526
0
                params.iceberg_commit_datas.insert(params.iceberg_commit_datas.end(),
2527
0
                                                   rs_icd.begin(), rs_icd.end());
2528
0
            }
2529
0
        }
2530
0
    }
2531
2532
0
    if (auto mcd = req.runtime_state->mc_commit_datas(); !mcd.empty()) {
2533
0
        params.__isset.mc_commit_datas = true;
2534
0
        params.mc_commit_datas.insert(params.mc_commit_datas.end(), mcd.begin(), mcd.end());
2535
0
    } else if (!req.runtime_states.empty()) {
2536
0
        for (auto* rs : req.runtime_states) {
2537
0
            if (auto rs_mcd = rs->mc_commit_datas(); !rs_mcd.empty()) {
2538
0
                params.__isset.mc_commit_datas = true;
2539
0
                params.mc_commit_datas.insert(params.mc_commit_datas.end(), rs_mcd.begin(),
2540
0
                                              rs_mcd.end());
2541
0
            }
2542
0
        }
2543
0
    }
2544
2545
0
    req.runtime_state->get_unreported_errors(&(params.error_log));
2546
0
    params.__isset.error_log = (!params.error_log.empty());
2547
2548
0
    if (_exec_env->cluster_info()->backend_id != 0) {
2549
0
        params.__set_backend_id(_exec_env->cluster_info()->backend_id);
2550
0
    }
2551
2552
0
    TReportExecStatusResult res;
2553
0
    Status rpc_status;
2554
2555
0
    VLOG_DEBUG << "reportExecStatus params is "
2556
0
               << apache::thrift::ThriftDebugString(params).c_str();
2557
0
    if (!exec_status.ok()) {
2558
0
        LOG(WARNING) << "report error status: " << exec_status.msg()
2559
0
                     << " to coordinator: " << req.coord_addr
2560
0
                     << ", query id: " << print_id(req.query_id);
2561
0
    }
2562
0
    try {
2563
0
        try {
2564
0
            (*coord)->reportExecStatus(res, params);
2565
0
        } catch (apache::thrift::transport::TTransportException& e) {
2566
0
            LOG(WARNING) << "Retrying ReportExecStatus. query id: " << print_id(req.query_id)
2567
0
                         << ", instance id: " << print_id(req.fragment_instance_id) << " to "
2568
0
                         << req.coord_addr << ", err: " << e.what();
2569
0
            rpc_status = coord->reopen();
2570
2571
0
            if (!rpc_status.ok()) {
2572
0
                req.cancel_fn(rpc_status);
2573
0
                return;
2574
0
            }
2575
0
            (*coord)->reportExecStatus(res, params);
2576
0
        }
2577
2578
0
        rpc_status = Status::create<false>(res.status);
2579
0
    } catch (apache::thrift::TException& e) {
2580
0
        rpc_status = Status::InternalError("ReportExecStatus() to {} failed: {}",
2581
0
                                           PrintThriftNetworkAddress(req.coord_addr), e.what());
2582
0
    }
2583
2584
0
    if (!rpc_status.ok()) {
2585
0
        LOG_INFO("Going to cancel query {} since report exec status got rpc failed: {}",
2586
0
                 print_id(req.query_id), rpc_status.to_string());
2587
0
        req.cancel_fn(rpc_status);
2588
0
    }
2589
0
}
2590
2591
2
Status PipelineFragmentContext::send_report(bool done) {
2592
2
    Status exec_status = _query_ctx->exec_status();
2593
2594
2
    if (!_is_report_success) {
2595
        // _is_report_success means this is not a load job, do not need to report to fe periodically.
2596
2
        if (exec_status.is<ErrorCode::LIMIT_REACH>() || exec_status.is<ErrorCode::FINISHED>() ||
2597
2
            exec_status.ok()) {
2598
2
            return Status::OK();
2599
2
        } else {
2600
            // else it means there is some error in processing the query, and we need to send report to FE to let FE know the error.
2601
0
        }
2602
2
    } else {
2603
        // This is a load job, need report the process status to FE periodly, so that FE can know the process of the load job.
2604
0
    }
2605
2606
0
    std::vector<RuntimeState*> runtime_states;
2607
2608
0
    for (auto& tasks : _tasks) {
2609
0
        for (auto& task : tasks) {
2610
0
            runtime_states.push_back(task.second.get());
2611
0
        }
2612
0
    }
2613
2614
0
    std::string load_eror_url = _query_ctx->get_load_error_url().empty()
2615
0
                                        ? get_load_error_url()
2616
0
                                        : _query_ctx->get_load_error_url();
2617
0
    std::string first_error_msg = _query_ctx->get_first_error_msg().empty()
2618
0
                                          ? get_first_error_msg()
2619
0
                                          : _query_ctx->get_first_error_msg();
2620
2621
0
    ReportStatusRequest req {.status = exec_status,
2622
0
                             .runtime_states = runtime_states,
2623
0
                             .done = done || !exec_status.ok(),
2624
0
                             .coord_addr = _query_ctx->coord_addr,
2625
0
                             .query_id = _query_id,
2626
0
                             .fragment_id = _fragment_id,
2627
0
                             .fragment_instance_id = TUniqueId(),
2628
0
                             .backend_num = -1,
2629
0
                             .runtime_state = _runtime_state.get(),
2630
0
                             .load_error_url = load_eror_url,
2631
0
                             .first_error_msg = first_error_msg,
2632
0
                             .cancel_fn = [this](const Status& reason) { cancel(reason); }};
2633
0
    auto ctx = std::dynamic_pointer_cast<PipelineFragmentContext>(shared_from_this());
2634
0
    return _exec_env->fragment_mgr()->get_thread_pool()->submit_func([this, req, ctx]() {
2635
0
        SCOPED_ATTACH_TASK(ctx->get_query_ctx()->query_mem_tracker());
2636
0
        _coordinator_callback(req);
2637
0
        if (!req.done) {
2638
0
            ctx->refresh_next_report_time();
2639
0
        }
2640
0
    });
2641
2
}
2642
2643
0
size_t PipelineFragmentContext::get_revocable_size(bool* has_running_task) const {
2644
0
    size_t res = 0;
2645
    // _tasks will be cleared during ~PipelineFragmentContext, so that it's safe
2646
    // here to traverse the vector.
2647
0
    for (const auto& task_instances : _tasks) {
2648
0
        for (const auto& task : task_instances) {
2649
0
            if (task.first->is_running()) {
2650
0
                LOG_EVERY_N(INFO, 50) << "Query: " << print_id(_query_id)
2651
0
                                      << " is running, task: " << (void*)task.first.get()
2652
0
                                      << ", is_running: " << task.first->is_running();
2653
0
                *has_running_task = true;
2654
0
                return 0;
2655
0
            }
2656
2657
0
            size_t revocable_size = task.first->get_revocable_size();
2658
0
            if (revocable_size >= SpillFile::MIN_SPILL_WRITE_BATCH_MEM) {
2659
0
                res += revocable_size;
2660
0
            }
2661
0
        }
2662
0
    }
2663
0
    return res;
2664
0
}
2665
2666
0
std::vector<PipelineTask*> PipelineFragmentContext::get_revocable_tasks() const {
2667
0
    std::vector<PipelineTask*> revocable_tasks;
2668
0
    for (const auto& task_instances : _tasks) {
2669
0
        for (const auto& task : task_instances) {
2670
0
            size_t revocable_size_ = task.first->get_revocable_size();
2671
2672
0
            if (revocable_size_ >= SpillFile::MIN_SPILL_WRITE_BATCH_MEM) {
2673
0
                revocable_tasks.emplace_back(task.first.get());
2674
0
            }
2675
0
        }
2676
0
    }
2677
0
    return revocable_tasks;
2678
0
}
2679
2680
0
std::string PipelineFragmentContext::debug_string() {
2681
0
    std::lock_guard<std::mutex> l(_task_mutex);
2682
0
    fmt::memory_buffer debug_string_buffer;
2683
0
    fmt::format_to(debug_string_buffer,
2684
0
                   "PipelineFragmentContext Info: _closed_tasks={}, _total_tasks={}, "
2685
0
                   "need_notify_close={}, fragment_id={}, _rec_cte_stage={}\n",
2686
0
                   _closed_tasks, _total_tasks, _need_notify_close, _fragment_id, _rec_cte_stage);
2687
0
    for (size_t j = 0; j < _tasks.size(); j++) {
2688
0
        fmt::format_to(debug_string_buffer, "Tasks in instance {}:\n", j);
2689
0
        for (size_t i = 0; i < _tasks[j].size(); i++) {
2690
0
            fmt::format_to(debug_string_buffer, "Task {}: {}\n", i,
2691
0
                           _tasks[j][i].first->debug_string());
2692
0
        }
2693
0
    }
2694
2695
0
    return fmt::to_string(debug_string_buffer);
2696
0
}
2697
2698
std::vector<std::shared_ptr<TRuntimeProfileTree>>
2699
0
PipelineFragmentContext::collect_realtime_profile() const {
2700
0
    std::vector<std::shared_ptr<TRuntimeProfileTree>> res;
2701
2702
    // we do not have mutex to protect pipeline_id_to_profile
2703
    // so we need to make sure this funciton is invoked after fragment context
2704
    // has already been prepared.
2705
0
    if (!_prepared) {
2706
0
        std::string msg =
2707
0
                "Query " + print_id(_query_id) + " collecting profile, but its not prepared";
2708
0
        DCHECK(false) << msg;
2709
0
        LOG_ERROR(msg);
2710
0
        return res;
2711
0
    }
2712
2713
    // Make sure first profile is fragment level profile
2714
0
    auto fragment_profile = std::make_shared<TRuntimeProfileTree>();
2715
0
    _fragment_level_profile->to_thrift(fragment_profile.get(), _runtime_state->profile_level());
2716
0
    res.push_back(fragment_profile);
2717
2718
    // pipeline_id_to_profile is initialized in prepare stage
2719
0
    for (auto pipeline_profile : _runtime_state->pipeline_id_to_profile()) {
2720
0
        auto profile_ptr = std::make_shared<TRuntimeProfileTree>();
2721
0
        pipeline_profile->to_thrift(profile_ptr.get(), _runtime_state->profile_level());
2722
0
        res.push_back(profile_ptr);
2723
0
    }
2724
2725
0
    return res;
2726
0
}
2727
2728
std::shared_ptr<TRuntimeProfileTree>
2729
0
PipelineFragmentContext::collect_realtime_load_channel_profile() const {
2730
    // we do not have mutex to protect pipeline_id_to_profile
2731
    // so we need to make sure this funciton is invoked after fragment context
2732
    // has already been prepared.
2733
0
    if (!_prepared) {
2734
0
        std::string msg =
2735
0
                "Query " + print_id(_query_id) + " collecting profile, but its not prepared";
2736
0
        DCHECK(false) << msg;
2737
0
        LOG_ERROR(msg);
2738
0
        return nullptr;
2739
0
    }
2740
2741
0
    for (const auto& tasks : _tasks) {
2742
0
        for (const auto& task : tasks) {
2743
0
            if (task.second->load_channel_profile() == nullptr) {
2744
0
                continue;
2745
0
            }
2746
2747
0
            auto tmp_load_channel_profile = std::make_shared<TRuntimeProfileTree>();
2748
2749
0
            task.second->load_channel_profile()->to_thrift(tmp_load_channel_profile.get(),
2750
0
                                                           _runtime_state->profile_level());
2751
0
            _runtime_state->load_channel_profile()->update(*tmp_load_channel_profile);
2752
0
        }
2753
0
    }
2754
2755
0
    auto load_channel_profile = std::make_shared<TRuntimeProfileTree>();
2756
0
    _runtime_state->load_channel_profile()->to_thrift(load_channel_profile.get(),
2757
0
                                                      _runtime_state->profile_level());
2758
0
    return load_channel_profile;
2759
0
}
2760
2761
// Collect runtime filter IDs registered by all tasks in this PFC.
2762
// Used during recursive CTE stage transitions to know which filters to deregister
2763
// before creating the new PFC for the next recursion round.
2764
// Called from rerun_fragment(wait_for_destroy) while tasks are still closing.
2765
// Thread safety: safe because _tasks is structurally immutable after prepare() —
2766
// the vector sizes do not change, and individual RuntimeState filter sets are
2767
// written only during open() which has completed by the time we reach rerun.
2768
0
std::set<int> PipelineFragmentContext::get_deregister_runtime_filter() const {
2769
0
    std::set<int> result;
2770
0
    for (const auto& _task : _tasks) {
2771
0
        for (const auto& task : _task) {
2772
0
            auto set = task.first->runtime_state()->get_deregister_runtime_filter();
2773
0
            result.merge(set);
2774
0
        }
2775
0
    }
2776
0
    if (_runtime_state) {
2777
0
        auto set = _runtime_state->get_deregister_runtime_filter();
2778
0
        result.merge(set);
2779
0
    }
2780
0
    return result;
2781
0
}
2782
2783
34
void PipelineFragmentContext::_release_resource() {
2784
34
    std::lock_guard<std::mutex> l(_task_mutex);
2785
    // The memory released by the query end is recorded in the query mem tracker.
2786
34
    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(_query_ctx->query_mem_tracker());
2787
34
    auto st = _query_ctx->exec_status();
2788
34
    for (auto& _task : _tasks) {
2789
0
        if (!_task.empty()) {
2790
0
            _call_back(_task.front().first->runtime_state(), &st);
2791
0
        }
2792
0
    }
2793
34
    _tasks.clear();
2794
34
    _dag.clear();
2795
34
    _pip_id_to_pipeline.clear();
2796
34
    _pipelines.clear();
2797
34
    _sink.reset();
2798
34
    _root_op.reset();
2799
34
    _runtime_filter_mgr_map.clear();
2800
34
    _op_id_to_shared_state.clear();
2801
34
}
2802
2803
} // namespace doris