Coverage Report

Created: 2026-08-07 00:14

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/scan/scanner_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/scan/scanner_context.h"
19
20
#include <fmt/format.h>
21
#include <gen_cpp/Metrics_types.h>
22
#include <glog/logging.h>
23
#include <zconf.h>
24
25
#include <cstdint>
26
#include <ctime>
27
#include <memory>
28
#include <mutex>
29
#include <ostream>
30
#include <shared_mutex>
31
#include <tuple>
32
#include <utility>
33
34
#include "common/config.h"
35
#include "common/exception.h"
36
#include "common/logging.h"
37
#include "common/metrics/doris_metrics.h"
38
#include "common/status.h"
39
#include "core/block/block.h"
40
#include "exec/operator/scan_operator.h"
41
#include "exec/scan/scan_node.h"
42
#include "exec/scan/scanner_scheduler.h"
43
#include "exec/scan/task_executor/task_executor.h"
44
#include "runtime/descriptors.h"
45
#include "runtime/exec_env.h"
46
#include "runtime/runtime_profile.h"
47
#include "runtime/runtime_state.h"
48
#include "runtime/thread_context.h"
49
#include "runtime/workload_management/resource_context.h"
50
#include "storage/tablet/tablet.h"
51
#include "util/time.h"
52
#include "util/uid_util.h"
53
54
namespace doris {
55
56
using namespace std::chrono_literals;
57
58
// ==================== ScanTask ====================
59
181
ScanTask::ScanTask(std::weak_ptr<ScannerDelegate> delegate_scanner) : scanner(delegate_scanner) {
60
181
    _resource_ctx = thread_context()->resource_ctx();
61
181
    DorisMetrics::instance()->scanner_task_cnt->increment(1);
62
181
}
63
64
181
ScanTask::~ScanTask() {
65
181
    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(_resource_ctx->memory_context()->mem_tracker());
66
181
    DorisMetrics::instance()->scanner_task_cnt->increment(-1);
67
181
    cached_block.reset();
68
181
}
69
70
// ==================== ScannerContext ====================
71
ScannerContext::ScannerContext(RuntimeState* state, ScanLocalStateBase* local_state,
72
                               const TupleDescriptor* output_tuple_desc,
73
                               const RowDescriptor* output_row_descriptor,
74
                               const std::list<std::shared_ptr<ScannerDelegate>>& scanners,
75
                               int64_t limit_, std::shared_ptr<Dependency> dependency,
76
                               std::atomic<int64_t>* shared_scan_limit,
77
                               std::shared_ptr<MemShareArbitrator> arb,
78
                               std::shared_ptr<MemLimiter> limiter, int ins_idx,
79
                               bool enable_adaptive_scan
80
#ifdef BE_TEST
81
                               ,
82
                               int num_parallel_instances
83
#endif
84
                               )
85
18
        : HasTaskExecutionCtx(state),
86
18
          _state(state),
87
18
          _local_state(local_state),
88
18
          _output_tuple_desc(output_row_descriptor
89
18
                                     ? output_row_descriptor->tuple_descriptors().front()
90
18
                                     : output_tuple_desc),
91
18
          _output_row_descriptor(output_row_descriptor),
92
18
          _batch_size(state->batch_size()),
93
18
          limit(limit_),
94
18
          _shared_scan_limit(shared_scan_limit),
95
18
          _all_scanners(scanners.begin(), scanners.end()),
96
#ifndef BE_TEST
97
          _scanner_scheduler(local_state->scan_scheduler(state)),
98
          _min_scan_concurrency_of_scan_scheduler(
99
                  _scanner_scheduler->get_min_active_scan_threads()),
100
          _max_scan_concurrency(std::min(local_state->max_scanners_concurrency(state),
101
                                         cast_set<int>(scanners.size()))),
102
#else
103
18
          _scanner_scheduler(state->get_query_ctx()->get_scan_scheduler()),
104
18
          _min_scan_concurrency_of_scan_scheduler(0),
105
18
          _max_scan_concurrency(num_parallel_instances),
106
#endif
107
18
          _min_scan_concurrency(local_state->min_scanners_concurrency(state)),
108
18
          _scanner_mem_limiter(limiter),
109
18
          _mem_share_arb(arb),
110
18
          _ins_idx(ins_idx),
111
18
          _enable_adaptive_scanners(enable_adaptive_scan) {
112
18
    DCHECK(_state != nullptr);
113
18
    DCHECK(_output_row_descriptor == nullptr ||
114
18
           _output_row_descriptor->tuple_descriptors().size() == 1);
115
18
    _query_id = _state->get_query_ctx()->query_id();
116
18
    _resource_ctx = _state->get_query_ctx()->resource_ctx();
117
18
    ctx_id = UniqueId::gen_uid().to_string();
118
171
    for (auto& scanner : _all_scanners) {
119
171
        _pending_tasks.push(std::make_shared<ScanTask>(scanner));
120
171
    }
121
18
    if (limit < 0) {
122
0
        limit = -1;
123
0
    }
124
18
    _dependency = dependency;
125
    // Initialize adaptive processor
126
18
    _adaptive_processor = ScannerAdaptiveProcessor::create_shared();
127
18
    DorisMetrics::instance()->scanner_ctx_cnt->increment(1);
128
18
}
129
130
5
void ScannerContext::_adjust_scan_mem_limit(int64_t old_value, int64_t new_value) {
131
5
    if (!_enable_adaptive_scanners) {
132
0
        return;
133
0
    }
134
135
5
    int64_t new_scan_mem_limit = _mem_share_arb->update_mem_bytes(old_value, new_value);
136
5
    _scanner_mem_limiter->update_mem_limit(new_scan_mem_limit);
137
5
    _scanner_mem_limiter->update_arb_mem_bytes(new_value);
138
139
5
    VLOG_DEBUG << fmt::format(
140
0
            "adjust_scan_mem_limit. context = {}, new mem scan limit = {}, scanner mem bytes = {} "
141
0
            "-> {}",
142
0
            debug_string(), new_scan_mem_limit, old_value, new_value);
143
5
}
144
145
12
int ScannerContext::_available_pickup_scanner_count() {
146
12
    if (!_enable_adaptive_scanners) {
147
12
        return _max_scan_concurrency;
148
12
    }
149
150
0
    int min_scanners = std::max(1, _min_scan_concurrency);
151
0
    int max_scanners = _scanner_mem_limiter->available_scanner_count(_ins_idx);
152
0
    max_scanners = std::min(max_scanners, _max_scan_concurrency);
153
0
    min_scanners = std::min(min_scanners, max_scanners);
154
0
    if (_ins_idx == 0) {
155
        // Adjust memory limit via memory share arbitrator
156
0
        _adjust_scan_mem_limit(_scanner_mem_limiter->get_arb_scanner_mem_bytes(),
157
0
                               _scanner_mem_limiter->get_estimated_block_mem_bytes());
158
0
    }
159
160
0
    ScannerAdaptiveProcessor& P = *_adaptive_processor;
161
0
    int& scanners = P.expected_scanners;
162
0
    int64_t now = UnixMillis();
163
    // Avoid frequent adjustment - only adjust every 100ms
164
0
    if (now - P.adjust_scanners_last_timestamp <= config::doris_scanner_dynamic_interval_ms) {
165
0
        return scanners;
166
0
    }
167
0
    P.adjust_scanners_last_timestamp = now;
168
0
    auto old_scanners = P.expected_scanners;
169
170
0
    scanners = std::max(min_scanners, scanners);
171
0
    scanners = std::min(max_scanners, scanners);
172
0
    VLOG_DEBUG << fmt::format(
173
0
            "_available_pickup_scanner_count. context = {}, old_scanners = {}, scanners = {} "
174
0
            ", min_scanners: {}, max_scanners: {}",
175
0
            debug_string(), old_scanners, scanners, min_scanners, max_scanners);
176
177
    // TODO(gabriel): Scanners are scheduled adaptively based on the memory usage now.
178
0
    return scanners;
179
0
}
180
181
// After init function call, should not access _parent
182
6
Status ScannerContext::init() {
183
#ifndef BE_TEST
184
    _scanner_profile = _local_state->_scanner_profile;
185
    _newly_create_free_blocks_num = _local_state->_newly_create_free_blocks_num;
186
    _scanner_memory_used_counter = _local_state->_memory_used_counter;
187
188
    // 3. get thread token
189
    if (!_state->get_query_ctx()) {
190
        return Status::InternalError("Query context of {} is not set",
191
                                     print_id(_state->query_id()));
192
    }
193
194
    if (_state->get_query_ctx()->get_scan_scheduler()) {
195
        _should_reset_thread_name = false;
196
    }
197
198
    auto scanner = _all_scanners.front().lock();
199
    DCHECK(scanner != nullptr);
200
201
    if (auto* task_executor_scheduler =
202
                dynamic_cast<TaskExecutorSimplifiedScanScheduler*>(_scanner_scheduler)) {
203
        std::shared_ptr<TaskExecutor> task_executor = task_executor_scheduler->task_executor();
204
        _task_executor = task_executor;
205
        TaskId task_id(fmt::format("{}-{}", print_id(_state->query_id()), ctx_id));
206
        _task_handle = DORIS_TRY(task_executor->create_task(
207
                task_id, []() { return 0.0; },
208
                config::task_executor_initial_max_concurrency_per_task > 0
209
                        ? config::task_executor_initial_max_concurrency_per_task
210
                        : std::max(48, CpuInfo::num_cores() * 2),
211
                std::chrono::milliseconds(100), std::nullopt));
212
    }
213
#endif
214
    // _max_bytes_in_queue controls the maximum memory that can be used by a single scan operator.
215
    // scan_queue_mem_limit on FE is 100MB by default, on backend we will make sure its actual value
216
    // is larger than 10MB.
217
6
    _max_bytes_in_queue = std::max(_state->scan_queue_mem_limit(), (int64_t)1024 * 1024 * 10);
218
219
    // Provide more memory for wide tables, increase proportionally by multiples of 300
220
6
    _max_bytes_in_queue *= _output_tuple_desc->slots().size() / 300 + 1;
221
222
6
    if (_all_scanners.empty()) {
223
0
        _is_finished = true;
224
0
        _set_scanner_done();
225
0
    }
226
227
    // Initialize memory limiter if memory-aware scheduling is enabled
228
6
    if (_enable_adaptive_scanners) {
229
0
        DCHECK(_scanner_mem_limiter && _mem_share_arb);
230
0
        int64_t c = _scanner_mem_limiter->update_open_tasks_count(1);
231
        // TODO(gabriel): set estimated block size
232
0
        _scanner_mem_limiter->reestimated_block_mem_bytes(DEFAULT_SCANNER_MEM_BYTES);
233
0
        _scanner_mem_limiter->update_arb_mem_bytes(DEFAULT_SCANNER_MEM_BYTES);
234
0
        if (c == 0) {
235
            // First scanner context to open, adjust scan memory limit
236
0
            _adjust_scan_mem_limit(DEFAULT_SCANNER_MEM_BYTES,
237
0
                                   _scanner_mem_limiter->get_arb_scanner_mem_bytes());
238
0
        }
239
0
    }
240
241
    // when user not specify scan_thread_num, so we can try downgrade _max_thread_num.
242
    // becaue we found in a table with 5k columns, column reader may ocuppy too much memory.
243
    // you can refer https://github.com/apache/doris/issues/35340 for details.
244
6
    const int32_t max_column_reader_num = _state->max_column_reader_num();
245
246
6
    if (_max_scan_concurrency != 1 && max_column_reader_num > 0) {
247
0
        int32_t scan_column_num = cast_set<int32_t>(_output_tuple_desc->slots().size());
248
0
        int32_t current_column_num = scan_column_num * _max_scan_concurrency;
249
0
        if (current_column_num > max_column_reader_num) {
250
0
            int32_t new_max_thread_num = max_column_reader_num / scan_column_num;
251
0
            new_max_thread_num = new_max_thread_num <= 0 ? 1 : new_max_thread_num;
252
0
            if (new_max_thread_num < _max_scan_concurrency) {
253
0
                int32_t origin_max_thread_num = _max_scan_concurrency;
254
0
                _max_scan_concurrency = new_max_thread_num;
255
0
                LOG(INFO) << "downgrade query:" << print_id(_state->query_id())
256
0
                          << " scan's max_thread_num from " << origin_max_thread_num << " to "
257
0
                          << _max_scan_concurrency << ",column num: " << scan_column_num
258
0
                          << ", max_column_reader_num: " << max_column_reader_num;
259
0
            }
260
0
        }
261
0
    }
262
263
6
    COUNTER_SET(_local_state->_max_scan_concurrency, (int64_t)_max_scan_concurrency);
264
6
    COUNTER_SET(_local_state->_min_scan_concurrency, (int64_t)_min_scan_concurrency);
265
266
6
    std::unique_lock<std::mutex> l(_transfer_lock);
267
6
    RETURN_IF_ERROR(_scanner_scheduler->schedule_scan_task(shared_from_this(), nullptr, l));
268
269
6
    return Status::OK();
270
6
}
271
272
18
ScannerContext::~ScannerContext() {
273
18
    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(_resource_ctx->memory_context()->mem_tracker());
274
18
    _completed_tasks.clear();
275
18
    BlockUPtr block;
276
19
    while (_free_blocks.try_dequeue(block)) {
277
        // do nothing
278
1
    }
279
18
    block.reset();
280
18
    DorisMetrics::instance()->scanner_ctx_cnt->increment(-1);
281
282
    // Cleanup memory limiter if last context closing
283
18
    if (_enable_adaptive_scanners) {
284
4
        if (_scanner_mem_limiter->update_open_tasks_count(-1) == 1) {
285
            // Last scanner context to close, reset scan memory limit
286
4
            _adjust_scan_mem_limit(_scanner_mem_limiter->get_arb_scanner_mem_bytes(), 0);
287
4
        }
288
4
    }
289
290
18
    if (_task_handle) {
291
0
        if (auto task_executor = _task_executor.lock()) {
292
0
            static_cast<void>(task_executor->remove_task(_task_handle));
293
0
        }
294
0
        _task_handle = nullptr;
295
0
        _task_executor.reset();
296
0
    }
297
18
}
298
299
3
BlockUPtr ScannerContext::get_free_block(bool force) {
300
3
    BlockUPtr block = nullptr;
301
3
    if (_free_blocks.try_dequeue(block)) {
302
1
        DCHECK(block->mem_reuse());
303
1
        _block_memory_usage -= block->allocated_bytes();
304
1
        _scanner_memory_used_counter->set(_block_memory_usage);
305
        // A free block is reused, so the memory usage should be decreased
306
        // The caller of get_free_block will increase the memory usage
307
2
    } else if (_block_memory_usage < _max_bytes_in_queue || force) {
308
2
        _newly_create_free_blocks_num->update(1);
309
2
        block = Block::create_unique(_output_tuple_desc->slots(), 0);
310
2
    }
311
3
    return block;
312
3
}
313
314
1
void ScannerContext::return_free_block(BlockUPtr block) {
315
    // If under low memory mode, should not return the freeblock, it will occupy too much memory.
316
1
    if (!_local_state->low_memory_mode() && block->mem_reuse() &&
317
1
        _block_memory_usage < _max_bytes_in_queue) {
318
1
        size_t block_size_to_reuse = block->allocated_bytes();
319
1
        _block_memory_usage += block_size_to_reuse;
320
1
        _scanner_memory_used_counter->set(_block_memory_usage);
321
1
        block->clear_column_data();
322
        // Free blocks is used to improve memory efficiency. Failure during pushing back
323
        // free block will not incur any bad result so just ignore the return value.
324
1
        _free_blocks.enqueue(std::move(block));
325
1
    }
326
1
}
327
328
Status ScannerContext::submit_scan_task(std::shared_ptr<ScanTask> scan_task,
329
18
                                        std::unique_lock<std::mutex>& /*transfer_lock*/) {
330
    // increase _num_finished_scanners no matter the scan_task is submitted successfully or not.
331
    // since if submit failed, it will be added back by ScannerContext::push_back_scan_task
332
    // and _num_finished_scanners will be reduced.
333
    // if submit succeed, it will be also added back by ScannerContext::push_back_scan_task
334
    // see ScannerScheduler::_scanner_scan.
335
18
    _in_flight_tasks_num++;
336
18
    return _scanner_scheduler->submit(shared_from_this(), scan_task);
337
18
}
338
339
0
void ScannerContext::clear_free_blocks() {
340
0
    clear_blocks(_free_blocks);
341
0
}
342
343
5
void ScannerContext::push_back_scan_task(std::shared_ptr<ScanTask> scan_task) {
344
5
    if (scan_task->status_ok()) {
345
5
        if (scan_task->cached_block && scan_task->cached_block->rows() > 0) {
346
0
            Status st = validate_block_schema(scan_task->cached_block.get());
347
0
            if (!st.ok()) {
348
0
                scan_task->set_status(st);
349
0
            }
350
0
        }
351
5
    }
352
353
5
    std::lock_guard<std::mutex> l(_transfer_lock);
354
5
    if (!scan_task->status_ok()) {
355
0
        _process_status = scan_task->get_status();
356
0
    }
357
5
    _completed_tasks.push_back(scan_task);
358
5
    _in_flight_tasks_num--;
359
360
5
    _dependency->set_ready();
361
5
}
362
363
3
Status ScannerContext::get_block_from_queue(RuntimeState* state, Block* block, bool* eos, int id) {
364
3
    if (state->is_cancelled()) {
365
1
        _set_scanner_done();
366
1
        return state->cancel_reason();
367
1
    }
368
2
    std::unique_lock l(_transfer_lock);
369
370
2
    if (!_process_status.ok()) {
371
1
        _set_scanner_done();
372
1
        return _process_status;
373
1
    }
374
375
1
    std::shared_ptr<ScanTask> scan_task = nullptr;
376
377
1
    if (!_completed_tasks.empty() && !done()) {
378
        // https://en.cppreference.com/w/cpp/container/list/front
379
        // The behavior is undefined if the list is empty.
380
1
        scan_task = _completed_tasks.front();
381
1
        _completed_tasks.pop_front();
382
1
    }
383
384
1
    if (scan_task != nullptr) {
385
        // The abnormal status of scanner may come from the execution of the scanner itself,
386
        // or come from the scanner scheduler, such as TooManyTasks.
387
1
        if (!scan_task->status_ok()) {
388
            // TODO: If the scanner status is TooManyTasks, maybe we can retry the scanner after a while.
389
0
            _process_status = scan_task->get_status();
390
0
            _set_scanner_done();
391
0
            return _process_status;
392
0
        }
393
394
1
        if (scan_task->cached_block) {
395
            // No need to worry about small block, block is merged together when they are appended to cached_blocks.
396
0
            auto current_block = std::move(scan_task->cached_block);
397
0
            auto block_size = current_block->allocated_bytes();
398
0
            scan_task->cached_block.reset();
399
0
            _block_memory_usage -= block_size;
400
            // consume current block
401
0
            block->swap(*current_block);
402
0
            return_free_block(std::move(current_block));
403
0
        }
404
1
        VLOG_DEBUG << fmt::format(
405
0
                "ScannerContext {} get block from queue, current scan "
406
0
                "task remaing cached_block size {}, eos {}, scheduled tasks {}",
407
0
                ctx_id, _completed_tasks.size(), scan_task->is_eos(), _in_flight_tasks_num);
408
1
        if (scan_task->is_eos()) {
409
            // 1. if eos, record a finished scanner.
410
1
            _num_finished_scanners++;
411
1
            RETURN_IF_ERROR(_scanner_scheduler->schedule_scan_task(shared_from_this(), nullptr, l));
412
1
        } else {
413
0
            scan_task->set_state(ScanTask::State::IN_FLIGHT);
414
0
            RETURN_IF_ERROR(
415
0
                    _scanner_scheduler->schedule_scan_task(shared_from_this(), scan_task, l));
416
0
        }
417
1
    }
418
419
1
    if (_completed_tasks.empty() &&
420
1
        (_num_finished_scanners == _all_scanners.size() ||
421
1
         (_is_shared_scan_limit_exhausted() && _in_flight_tasks_num == 0))) {
422
1
        _set_scanner_done();
423
1
        _is_finished = true;
424
1
    }
425
426
1
    *eos = done();
427
428
1
    if (_completed_tasks.empty()) {
429
1
        _dependency->block();
430
1
    }
431
432
1
    return Status::OK();
433
1
}
434
435
0
Status ScannerContext::validate_block_schema(Block* block) {
436
0
    size_t index = 0;
437
0
    for (auto& slot : _output_tuple_desc->slots()) {
438
0
        auto& data = block->get_by_position(index++);
439
0
        if (data.column->is_nullable() != data.type->is_nullable()) {
440
0
            return Status::Error<ErrorCode::INVALID_SCHEMA>(
441
0
                    "column(name: {}) nullable({}) does not match type nullable({}), slot(id: "
442
0
                    "{}, "
443
0
                    "name:{})",
444
0
                    data.name, data.column->is_nullable(), data.type->is_nullable(), slot->id(),
445
0
                    slot->col_name());
446
0
        }
447
448
0
        if (data.column->is_nullable() != slot->is_nullable()) {
449
0
            return Status::Error<ErrorCode::INVALID_SCHEMA>(
450
0
                    "column(name: {}) nullable({}) does not match slot(id: {}, name: {}) "
451
0
                    "nullable({})",
452
0
                    data.name, data.column->is_nullable(), slot->id(), slot->col_name(),
453
0
                    slot->is_nullable());
454
0
        }
455
0
    }
456
0
    return Status::OK();
457
0
}
458
459
0
void ScannerContext::stop_scanners(RuntimeState* state) {
460
0
    std::lock_guard<std::mutex> l(_transfer_lock);
461
0
    if (_should_stop) {
462
0
        return;
463
0
    }
464
0
    _should_stop = true;
465
0
    _set_scanner_done();
466
0
    for (const std::weak_ptr<ScannerDelegate>& scanner : _all_scanners) {
467
0
        if (std::shared_ptr<ScannerDelegate> sc = scanner.lock()) {
468
0
            sc->_scanner->try_stop();
469
0
        }
470
0
    }
471
0
    _completed_tasks.clear();
472
0
    if (_task_handle) {
473
0
        if (auto task_executor = _task_executor.lock()) {
474
0
            static_cast<void>(task_executor->remove_task(_task_handle));
475
0
        }
476
0
        _task_handle = nullptr;
477
0
        _task_executor.reset();
478
0
    }
479
    // TODO yiguolei, call mark close to scanners
480
0
    if (state->enable_profile()) {
481
0
        std::stringstream scanner_statistics;
482
0
        std::stringstream scanner_rows_read;
483
0
        std::stringstream scanner_wait_worker_time;
484
0
        std::stringstream scanner_projection;
485
0
        std::stringstream scanner_prepare_time;
486
0
        std::stringstream scanner_open_time;
487
0
        scanner_statistics << "[";
488
0
        scanner_rows_read << "[";
489
0
        scanner_wait_worker_time << "[";
490
0
        scanner_projection << "[";
491
0
        scanner_prepare_time << "[";
492
0
        scanner_open_time << "[";
493
        // Scanners can in 3 state
494
        //  state 1: in scanner context, not scheduled
495
        //  state 2: in scanner worker pool's queue, scheduled but not running
496
        //  state 3: scanner is running.
497
0
        for (auto& scanner_ref : _all_scanners) {
498
0
            auto scanner = scanner_ref.lock();
499
0
            if (scanner == nullptr) {
500
0
                continue;
501
0
            }
502
            // Add per scanner running time before close them
503
0
            scanner_statistics << PrettyPrinter::print(scanner->_scanner->get_time_cost_ns(),
504
0
                                                       TUnit::TIME_NS)
505
0
                               << ", ";
506
0
            scanner_projection << PrettyPrinter::print(scanner->_scanner->projection_time(),
507
0
                                                       TUnit::TIME_NS)
508
0
                               << ", ";
509
0
            scanner_rows_read << PrettyPrinter::print(scanner->_scanner->get_rows_read(),
510
0
                                                      TUnit::UNIT)
511
0
                              << ", ";
512
0
            scanner_wait_worker_time
513
0
                    << PrettyPrinter::print(scanner->_scanner->get_scanner_wait_worker_timer(),
514
0
                                            TUnit::TIME_NS)
515
0
                    << ", ";
516
0
            scanner_prepare_time << PrettyPrinter::print(
517
0
                                            scanner->_scanner->get_prepare_time_cost_ns(),
518
0
                                            TUnit::TIME_NS)
519
0
                                 << ", ";
520
0
            scanner_open_time << PrettyPrinter::print(scanner->_scanner->get_open_time_cost_ns(),
521
0
                                                      TUnit::TIME_NS)
522
0
                              << ", ";
523
            // since there are all scanners, some scanners is running, so that could not call scanner
524
            // close here.
525
0
        }
526
0
        scanner_statistics << "]";
527
0
        scanner_rows_read << "]";
528
0
        scanner_wait_worker_time << "]";
529
0
        scanner_projection << "]";
530
0
        scanner_prepare_time << "]";
531
0
        scanner_open_time << "]";
532
0
        _scanner_profile->add_info_string("PerScannerRunningTime", scanner_statistics.str());
533
0
        _scanner_profile->add_info_string("PerScannerRowsRead", scanner_rows_read.str());
534
0
        _scanner_profile->add_info_string("PerScannerWaitTime", scanner_wait_worker_time.str());
535
0
        _scanner_profile->add_info_string("PerScannerProjectionTime", scanner_projection.str());
536
0
        _scanner_profile->add_info_string("PerScannerPrepareTime", scanner_prepare_time.str());
537
0
        _scanner_profile->add_info_string("PerScannerOpenTime", scanner_open_time.str());
538
0
    }
539
0
}
540
541
18
std::string ScannerContext::debug_string() {
542
18
    return fmt::format(
543
18
            "_query_id: {}, id: {}, total scanners: {}, pending tasks: {}, completed tasks: {},"
544
18
            " _should_stop: {}, _is_finished: {}, free blocks: {},"
545
18
            " limit: {}, _in_flight_tasks_num: {}, remaining_limit: {}, _num_running_scanners: {}, "
546
18
            "_max_thread_num: {},"
547
18
            " _max_bytes_in_queue: {}, _ins_idx: {}, _enable_adaptive_scanners: {}, "
548
18
            "_mem_share_arb: {}, _scanner_mem_limiter: {}",
549
18
            print_id(_query_id), ctx_id, _all_scanners.size(), _pending_tasks.size(),
550
18
            _completed_tasks.size(), _should_stop, _is_finished, _free_blocks.size_approx(), limit,
551
18
            _shared_scan_limit->load(std::memory_order_relaxed), _in_flight_tasks_num,
552
18
            _num_finished_scanners, _max_scan_concurrency, _max_bytes_in_queue, _ins_idx,
553
18
            _enable_adaptive_scanners,
554
18
            _enable_adaptive_scanners ? _mem_share_arb->debug_string() : "NULL",
555
18
            _enable_adaptive_scanners ? _scanner_mem_limiter->debug_string() : "NULL");
556
18
}
557
558
3
void ScannerContext::_set_scanner_done() {
559
3
    _dependency->set_always_ready();
560
3
}
561
562
20
bool ScannerContext::_is_shared_scan_limit_exhausted() const {
563
20
    return limit >= 0 && _shared_scan_limit->load(std::memory_order_acquire) <= 0;
564
20
}
565
566
1
void ScannerContext::update_peak_running_scanner(int num) {
567
#ifndef BE_TEST
568
    _local_state->_peak_running_scanner->add(num);
569
#endif
570
1
    if (_enable_adaptive_scanners) {
571
1
        _scanner_mem_limiter->update_running_tasks_count(num);
572
1
    }
573
1
}
574
575
1
void ScannerContext::reestimated_block_mem_bytes(int64_t num) {
576
1
    if (_enable_adaptive_scanners) {
577
1
        _scanner_mem_limiter->reestimated_block_mem_bytes(num);
578
1
    }
579
1
}
580
581
int32_t ScannerContext::_get_margin(std::unique_lock<std::mutex>& transfer_lock,
582
12
                                    std::unique_lock<std::shared_mutex>& scheduler_lock) {
583
    // Get effective max concurrency considering adaptive scheduling
584
12
    int32_t effective_max_concurrency = _available_pickup_scanner_count();
585
12
    DCHECK_LE(effective_max_concurrency, _max_scan_concurrency);
586
587
    // margin_1 is used to ensure each scan operator could have at least _min_scan_concurrency scan tasks.
588
12
    int32_t margin_1 = _min_scan_concurrency -
589
12
                       (cast_set<int32_t>(_completed_tasks.size()) + _in_flight_tasks_num);
590
591
    // margin_2 is used to ensure the scan scheduler could have at least _min_scan_concurrency_of_scan_scheduler scan tasks.
592
12
    int32_t margin_2 =
593
12
            _min_scan_concurrency_of_scan_scheduler -
594
12
            (_scanner_scheduler->get_active_threads() + _scanner_scheduler->get_queue_size());
595
596
    // margin_3 is used to respect adaptive max concurrency limit
597
12
    int32_t margin_3 =
598
12
            std::max(effective_max_concurrency -
599
12
                             (cast_set<int32_t>(_completed_tasks.size()) + _in_flight_tasks_num),
600
12
                     1);
601
602
12
    if (margin_1 <= 0 && margin_2 <= 0) {
603
1
        return 0;
604
1
    }
605
606
11
    int32_t margin = std::max(margin_1, margin_2);
607
11
    if (_enable_adaptive_scanners) {
608
0
        margin = std::min(margin, margin_3); // Cap by adaptive limit
609
0
    }
610
611
11
    if (low_memory_mode()) {
612
        // In low memory mode, we will limit the number of running scanners to `low_memory_mode_scanners()`.
613
        // So that we will not submit too many scan tasks to scheduler.
614
0
        margin = std::min(low_memory_mode_scanners() - _in_flight_tasks_num, margin);
615
0
    }
616
617
11
    VLOG_DEBUG << fmt::format(
618
0
            "[{}|{}] schedule scan task, margin_1: {} = {} - ({} + {}), margin_2: {} = {} - "
619
0
            "({} + {}), margin_3: {} = {} - ({} + {}), margin: {}, adaptive: {}",
620
0
            print_id(_query_id), ctx_id, margin_1, _min_scan_concurrency, _completed_tasks.size(),
621
0
            _in_flight_tasks_num, margin_2, _min_scan_concurrency_of_scan_scheduler,
622
0
            _scanner_scheduler->get_active_threads(), _scanner_scheduler->get_queue_size(),
623
0
            margin_3, effective_max_concurrency, _completed_tasks.size(), _in_flight_tasks_num,
624
0
            margin, _enable_adaptive_scanners);
625
626
11
    return margin;
627
12
}
628
629
// This function must be called with:
630
// 1. _transfer_lock held.
631
// 2. ScannerScheduler::_lock held.
632
Status ScannerContext::schedule_scan_task(std::shared_ptr<ScanTask> current_scan_task,
633
                                          std::unique_lock<std::mutex>& transfer_lock,
634
7
                                          std::unique_lock<std::shared_mutex>& scheduler_lock) {
635
7
    if (current_scan_task &&
636
7
        (current_scan_task->cached_block != nullptr || current_scan_task->is_eos())) {
637
1
        throw doris::Exception(ErrorCode::INTERNAL_ERROR, "Scanner scheduler logical error.");
638
1
    }
639
640
6
    std::list<std::shared_ptr<ScanTask>> tasks_to_submit;
641
642
6
    int32_t margin = _get_margin(transfer_lock, scheduler_lock);
643
644
    // margin is less than zero. Means this scan operator could not submit any scan task for now.
645
6
    if (margin <= 0) {
646
        // Be careful with current scan task.
647
        // We need to add it back to task queue to make sure it could be resubmitted.
648
0
        if (current_scan_task) {
649
            // This usually happens when we should downgrade the concurrency.
650
0
            current_scan_task->set_state(ScanTask::State::PENDING);
651
0
            _pending_tasks.push(current_scan_task);
652
0
            VLOG_DEBUG << fmt::format(
653
0
                    "{} push back scanner to task queue, because diff <= 0, _completed_tasks size "
654
0
                    "{}, _in_flight_tasks_num {}",
655
0
                    ctx_id, _completed_tasks.size(), _in_flight_tasks_num);
656
0
        }
657
658
0
#ifndef NDEBUG
659
        // This DCHECK is necessary.
660
        // We need to make sure each scan operator could have at least 1 scan tasks.
661
        // Or this scan operator will not be re-scheduled.
662
0
        if (!_pending_tasks.empty() && _in_flight_tasks_num == 0 && _completed_tasks.empty()) {
663
0
            throw doris::Exception(ErrorCode::INTERNAL_ERROR, "Scanner scheduler logical error.");
664
0
        }
665
0
#endif
666
667
0
        return Status::OK();
668
0
    }
669
670
6
    bool first_pull = true;
671
672
24
    while (margin-- > 0) {
673
24
        std::shared_ptr<ScanTask> task_to_run;
674
24
        const int32_t current_concurrency = cast_set<int32_t>(
675
24
                _completed_tasks.size() + _in_flight_tasks_num + tasks_to_submit.size());
676
24
        VLOG_DEBUG << fmt::format("{} currenct concurrency: {} = {} + {} + {}", ctx_id,
677
0
                                  current_concurrency, _completed_tasks.size(),
678
0
                                  _in_flight_tasks_num, tasks_to_submit.size());
679
24
        if (first_pull) {
680
6
            task_to_run = _pull_next_scan_task(current_scan_task, current_concurrency);
681
6
            if (task_to_run == nullptr) {
682
                // In three situations we will get nullptr.
683
                // 1. current_concurrency already reached _max_scan_concurrency.
684
                // 2. all scanners are finished.
685
                // 3. The shared LIMIT is exhausted while completed or in-flight tasks can still
686
                //    make progress.
687
2
                if (current_scan_task) {
688
1
                    DCHECK(current_scan_task->cached_block == nullptr);
689
1
                    DCHECK(!current_scan_task->is_eos());
690
1
                    if (current_scan_task->cached_block != nullptr || current_scan_task->is_eos()) {
691
                        // This should not happen.
692
0
                        throw doris::Exception(ErrorCode::INTERNAL_ERROR,
693
0
                                               "Scanner scheduler logical error.");
694
0
                    }
695
                    // Current scan task is not scheduled, we need to add it back to task queue to make sure it could be resubmitted.
696
1
                    current_scan_task->set_state(ScanTask::State::PENDING);
697
1
                    _pending_tasks.push(current_scan_task);
698
1
                }
699
2
            }
700
6
            first_pull = false;
701
18
        } else {
702
18
            task_to_run = _pull_next_scan_task(nullptr, current_concurrency);
703
18
        }
704
705
24
        if (task_to_run) {
706
18
            tasks_to_submit.push_back(task_to_run);
707
18
        } else {
708
6
            break;
709
6
        }
710
24
    }
711
712
6
    if (tasks_to_submit.empty()) {
713
2
        return Status::OK();
714
2
    }
715
716
4
    VLOG_DEBUG << fmt::format("[{}:{}] submit {} scan tasks to scheduler, remaining scanner: {}",
717
0
                              print_id(_query_id), ctx_id, tasks_to_submit.size(),
718
0
                              _pending_tasks.size());
719
720
18
    for (auto& scan_task_iter : tasks_to_submit) {
721
18
        Status submit_status = submit_scan_task(scan_task_iter, transfer_lock);
722
18
        if (!submit_status.ok()) {
723
0
            _process_status = submit_status;
724
0
            _set_scanner_done();
725
0
            return _process_status;
726
0
        }
727
18
    }
728
729
4
    return Status::OK();
730
4
}
731
732
std::shared_ptr<ScanTask> ScannerContext::_pull_next_scan_task(
733
31
        std::shared_ptr<ScanTask> current_scan_task, int32_t current_concurrency) {
734
31
    int32_t effective_max_concurrency = _max_scan_concurrency;
735
31
    if (_enable_adaptive_scanners) {
736
0
        effective_max_concurrency = _adaptive_processor->expected_scanners > 0
737
0
                                            ? _adaptive_processor->expected_scanners
738
0
                                            : _max_scan_concurrency;
739
0
    }
740
741
31
    if (current_concurrency >= effective_max_concurrency) {
742
7
        VLOG_DEBUG << fmt::format(
743
0
                "ScannerContext {} current concurrency {} >= effective_max_concurrency {}, skip "
744
0
                "pull",
745
0
                ctx_id, current_concurrency, effective_max_concurrency);
746
7
        return nullptr;
747
7
    }
748
749
24
    if (current_scan_task != nullptr) {
750
3
        if (current_scan_task->cached_block != nullptr || current_scan_task->is_eos()) {
751
            // This should not happen.
752
2
            throw doris::Exception(ErrorCode::INTERNAL_ERROR, "Scanner scheduler logical error.");
753
2
        }
754
1
        return current_scan_task;
755
3
    }
756
757
21
    if (!_pending_tasks.empty()) {
758
        // Do not submit more pending scanners after the shared LIMIT is exhausted while
759
        // completed or in-flight tasks can still make progress. If neither exists, allow pending
760
        // scanners to be submitted so they can report EOS and wake the pipeline task.
761
19
        if (_is_shared_scan_limit_exhausted() &&
762
19
            (_in_flight_tasks_num != 0 || !_completed_tasks.empty())) {
763
0
            return nullptr;
764
0
        }
765
19
        std::shared_ptr<ScanTask> next_scan_task;
766
19
        next_scan_task = _pending_tasks.top();
767
19
        _pending_tasks.pop();
768
19
        return next_scan_task;
769
19
    } else {
770
2
        return nullptr;
771
2
    }
772
21
}
773
774
11
bool ScannerContext::low_memory_mode() const {
775
11
    return _local_state->low_memory_mode();
776
11
}
777
} // namespace doris