Coverage Report

Created: 2026-10-10 03:14

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/scan/scanner_context.h
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
#pragma once
19
20
#include <bthread/types.h>
21
#include <stdint.h>
22
23
#include <atomic>
24
#include <cstdint>
25
#include <list>
26
#include <memory>
27
#include <mutex>
28
#include <stack>
29
#include <string>
30
#include <utility>
31
#include <vector>
32
33
#include "common/config.h"
34
#include "common/factory_creator.h"
35
#include "common/metrics/doris_metrics.h"
36
#include "common/status.h"
37
#include "concurrentqueue.h"
38
#include "core/block/block.h"
39
#include "exec/common/memory.h"
40
#include "exec/scan/scanner.h"
41
#include "exec/scan/task_executor/split_runner.h"
42
#include "runtime/runtime_profile.h"
43
44
namespace doris {
45
46
class ResourceContext;
47
class RuntimeState;
48
class TupleDescriptor;
49
class WorkloadGroup;
50
51
class ScanLocalStateBase;
52
class Dependency;
53
54
class Scanner;
55
class ScannerDelegate;
56
class ScannerScheduler;
57
class TaskExecutor;
58
class TaskHandle;
59
60
// Adaptive processor for dynamic scanner concurrency adjustment
61
struct ScannerAdaptiveProcessor {
62
    ENABLE_FACTORY_CREATOR(ScannerAdaptiveProcessor)
63
40
    ScannerAdaptiveProcessor() = default;
64
    ~ScannerAdaptiveProcessor() = default;
65
    // Expected scanners in this cycle
66
67
    int expected_scanners = 0;
68
    // Timing metrics
69
    // int64_t context_start_time = 0;
70
    // int64_t scanner_total_halt_time = 0;
71
    // int64_t scanner_gen_blocks_time = 0;
72
    // std::atomic_int64_t scanner_total_io_time = 0;
73
    // std::atomic_int64_t scanner_total_running_time = 0;
74
    // std::atomic_int64_t scanner_total_scan_bytes = 0;
75
76
    // Timestamps
77
    // std::atomic_int64_t last_scanner_finish_timestamp = 0;
78
    // int64_t check_all_scanners_last_timestamp = 0;
79
    // int64_t last_driver_output_full_timestamp = 0;
80
    int64_t adjust_scanners_last_timestamp = 0;
81
82
    // Adjustment strategy fields
83
    // bool try_add_scanners = false;
84
    // double expected_speedup_ratio = 0;
85
    // double last_scanner_scan_speed = 0;
86
    // int64_t last_scanner_total_scan_bytes = 0;
87
    // int try_add_scanners_fail_count = 0;
88
    // int check_slow_io = 0;
89
    // int32_t slow_io_latency_ms = 100; // Default from config
90
};
91
92
class ScanTask {
93
public:
94
    enum class State : int {
95
        PENDING,   // not scheduled yet
96
        IN_FLIGHT, // scheduled and running
97
        COMPLETED, // finished with result or error, waiting to be collected by scan node
98
        EOS,       // finished and no more data, waiting to be collected by scan node
99
        PARKED,    // waiting, off the workers, for what its scanner cannot read on without
100
    };
101
    ScanTask(std::weak_ptr<ScannerDelegate> delegate_scanner);
102
103
    ~ScanTask();
104
105
private:
106
    // whether current scanner is finished
107
    Status status = Status::OK();
108
    std::shared_ptr<ResourceContext> _resource_ctx;
109
    State _state = State::PENDING;
110
111
public:
112
    std::weak_ptr<ScannerDelegate> scanner;
113
    BlockUPtr cached_block = nullptr;
114
    bool is_first_schedule = true;
115
    // Use weak_ptr to avoid circular references and potential memory leaks with SplitRunner.
116
    // ScannerContext only needs to observe the lifetime of SplitRunner without owning it.
117
    // When SplitRunner is destroyed, split_runner.lock() will return nullptr, ensuring safe access.
118
    std::weak_ptr<SplitRunner> split_runner;
119
120
3
    void set_status(Status _status) {
121
3
        if (_status.is<ErrorCode::END_OF_FILE>()) {
122
            // set `eos` if `END_OF_FILE`, don't take `END_OF_FILE` as error
123
0
            _state = State::EOS;
124
0
        }
125
3
        status = _status;
126
3
    }
127
3
    Status get_status() const { return status; }
128
209
    bool status_ok() { return status.ok() || status.is<ErrorCode::END_OF_FILE>(); }
129
271
    bool is_eos() const { return _state == State::EOS; }
130
220
    void set_state(State state) {
131
220
        switch (state) {
132
51
        case State::PENDING:
133
            // A task returns to PENDING after the operator consumes its non-EOS cached block.
134
            // For example, one scanner may produce several blocks, so COMPLETED is not terminal.
135
            // A parked task returns to it once what it waited for is done.
136
51
            DCHECK(_state == State::PENDING || _state == State::IN_FLIGHT ||
137
0
                   _state == State::COMPLETED || _state == State::PARKED)
138
0
                    << (int)_state;
139
51
            DCHECK(cached_block == nullptr);
140
51
            break;
141
88
        case State::IN_FLIGHT:
142
88
            DCHECK(_state == State::COMPLETED || _state == State::PENDING ||
143
0
                   _state == State::IN_FLIGHT)
144
0
                    << (int)_state;
145
88
            DCHECK(cached_block == nullptr);
146
88
            break;
147
52
        case State::COMPLETED:
148
52
            DCHECK(_state == State::IN_FLIGHT) << (int)_state;
149
52
            DCHECK(cached_block != nullptr);
150
52
            break;
151
20
        case State::EOS:
152
20
            DCHECK(_state == State::IN_FLIGHT || status.is<ErrorCode::END_OF_FILE>())
153
0
                    << (int)_state;
154
20
            break;
155
9
        case State::PARKED:
156
9
            DCHECK(_state == State::IN_FLIGHT) << (int)_state;
157
9
            DCHECK(cached_block == nullptr);
158
9
            break;
159
0
        default:
160
0
            break;
161
220
        }
162
163
220
        _state = state;
164
220
    }
165
};
166
167
// ScannerContext is responsible for recording the execution status
168
// of a group of Scanners corresponding to a ScanNode.
169
// Including how many scanners are being scheduled, and maintaining
170
// a producer-consumer blocks queue between scanners and scan nodes.
171
//
172
// ScannerContext is also the scheduling unit of ScannerScheduler.
173
// ScannerScheduler schedules a ScannerContext at a time,
174
// and submits the Scanners to the scanner thread pool for data scanning.
175
class ScannerContext : public std::enable_shared_from_this<ScannerContext>,
176
                       public HasTaskExecutionCtx {
177
    ENABLE_FACTORY_CREATOR(ScannerContext);
178
    friend class ScannerScheduler;
179
180
public:
181
    ScannerContext(RuntimeState* state, ScanLocalStateBase* local_state,
182
                   const TupleDescriptor* output_tuple_desc, bool has_projection,
183
                   const std::list<std::shared_ptr<ScannerDelegate>>& scanners, int64_t limit_,
184
                   std::shared_ptr<Dependency> dependency, std::atomic<int64_t>* shared_scan_limit,
185
                   std::shared_ptr<MemShareArbitrator> arb, std::shared_ptr<MemLimiter> limiter,
186
                   int ins_idx, bool enable_adaptive_scan
187
#ifdef BE_TEST
188
                   ,
189
                   int num_parallel_instances
190
#endif
191
    );
192
193
    ~ScannerContext() override;
194
    Status init();
195
196
    // TODO(gabriel): we can also consider to return a list of blocks to reduce the scheduling overhead, but it may cause larger memory usage and more complex logic of block management.
197
    BlockUPtr get_free_block(bool force);
198
    void return_free_block(BlockUPtr block);
199
    void clear_free_blocks();
200
62
    inline void inc_block_usage(size_t usage) { _block_memory_usage += usage; }
201
202
0
    int64_t block_memory_usage() { return _block_memory_usage; }
203
204
    // Caller should make sure the pipeline task is still running when calling this function
205
    void update_peak_running_scanner(int num);
206
    void reestimated_block_mem_bytes(int64_t num);
207
208
    // Get next block from blocks queue. Called by ScanNode/ScanOperator
209
    // Set eos to true if there is no more data to read.
210
    Status get_block_from_queue(RuntimeState* state, Block* block, bool* eos, int id);
211
212
    [[nodiscard]] Status validate_block_schema(Block* block);
213
214
    // submit the running scanner to thread pool in `ScannerScheduler`
215
    // set the next scanned block to `ScanTask::current_block`
216
    // set the error state to `ScanTask::status`
217
    // set the `eos` to `ScanTask::eos` if there is no more data in current scanner
218
    Status submit_scan_task(std::shared_ptr<ScanTask> scan_task, std::unique_lock<std::mutex>&);
219
220
    // Publish a task whose current scan attempt has completed. The operator consumes its cached
221
    // block and returns a non-EOS task to PENDING for its next scan attempt.
222
    void push_completed_scan_task(std::shared_ptr<ScanTask> scan_task);
223
224
    // Park a task whose scan attempt ended without a block because its scanner cannot read on until
225
    // `waiting_for` is done (Scanner::take_waiting_for()). A parked task holds no worker and no
226
    // concurrency slot - another scanner of this context, perhaps one holding what it waits for,
227
    // is admitted in its place - and returns to scheduling, as if the operator had just consumed
228
    // its block, once the future is done. The caller must not touch the task afterwards: by then
229
    // another worker may already be running it.
230
    void park_scan_task(std::shared_ptr<ScanTask> scan_task,
231
                        SharedListenableFuture<Void> waiting_for);
232
233
    // Return true if this ScannerContext need no more process
234
498
    bool done() const { return _is_finished || _should_stop; }
235
236
    std::string debug_string();
237
238
17
    std::shared_ptr<TaskHandle> task_handle() const { return _task_handle; }
239
240
0
    std::shared_ptr<ResourceContext> resource_ctx() const { return _resource_ctx; }
241
242
150
    RuntimeState* state() { return _state; }
243
244
    void stop_scanners(RuntimeState* state);
245
246
50
    int batch_size() const { return _batch_size; }
247
248
    // During low memory mode, there will be at most 4 scanners running and every scanner will
249
    // cache at most 1MB data. So that every instance will keep 8MB buffer.
250
    bool low_memory_mode() const;
251
252
    // TODO(yiguolei) add this as session variable
253
0
    int32_t low_memory_mode_scan_bytes_per_scanner() const {
254
0
        return 1 * 1024 * 1024; // 1MB
255
0
    }
256
257
0
    int32_t low_memory_mode_scanners() const { return 4; }
258
259
0
    ScanLocalStateBase* local_state() const { return _local_state; }
260
261
    // the unique id of this context
262
    std::string ctx_id;
263
    TUniqueId _query_id;
264
265
    bool _should_reset_thread_name = true;
266
267
0
    int32_t num_scheduled_scanners() {
268
0
        std::lock_guard<std::mutex> l(_transfer_lock);
269
0
        return _in_flight_tasks_num;
270
0
    }
271
272
    Status schedule_scan_task(std::shared_ptr<ScanTask> current_scan_task,
273
                              std::unique_lock<std::mutex>& transfer_lock,
274
                              std::unique_lock<std::shared_mutex>& scheduler_lock);
275
276
    // Context scheduling and operator consumption share this lock so queue-state changes and task
277
    // admission form one atomic decision. For example, two worker callbacks cannot both admit the
278
    // last available concurrency slot.
279
79
    std::mutex& transfer_lock() { return _transfer_lock; }
280
281
    // One Context submission represents many pending scanners in the ThreadPool scheduler.
282
    // Keeping this separate from scanner execution prevents duplicate runnables from accumulating.
283
    bool is_context_queued(const std::unique_lock<std::mutex>& transfer_lock) const;
284
    // Transition the Context runnable's queue state. The caller must hold _transfer_lock.
285
    void set_context_queued(bool queued, const std::unique_lock<std::mutex>& transfer_lock);
286
287
    // Publish a scheduler failure and make the Context terminal. The caller must hold
288
    // _transfer_lock so a retained ThreadPool callback cannot admit another scanner concurrently.
289
    void set_context_failure(const Status& failure,
290
                             const std::unique_lock<std::mutex>& transfer_lock);
291
292
    // Return a scanner to the admission queue after its block is consumed. It may not own a cached
293
    // block and may not be EOS: EOS scanners are terminal and must not run again.
294
    void push_pending_scan_task(std::shared_ptr<ScanTask> scan_task,
295
                                const std::unique_lock<std::mutex>& transfer_lock);
296
297
    // Return whether a Context worker can currently admit one pending scanner. This check has no
298
    // side effects, so the scheduler can avoid submitting a runnable that would immediately exit.
299
    // It always admits one scanner when nothing is progressing so the operator can be woken, and it
300
    // holds the Context at max(1, _min_scan_concurrency) while the scheduler pool has no slack,
301
    // like _get_margin() on the TaskExecutor path. The caller must hold _transfer_lock.
302
    bool can_admit_scan_task(const std::unique_lock<std::mutex>& transfer_lock) const;
303
304
    // Atomically check whether this context can start another scan task, move one task from
305
    // pending to in-flight, and return it. The caller must hold _transfer_lock.
306
    std::shared_ptr<ScanTask> try_get_next_scan_task(
307
            const std::unique_lock<std::mutex>& transfer_lock);
308
309
protected:
310
    /// Four criteria to determine whether to increase the parallelism of the scanners
311
    /// 1. It ran for at least `SCALE_UP_DURATION` ms after last scale up
312
    /// 2. Half(`WAIT_BLOCK_DURATION_RATIO`) of the duration is waiting to get blocks
313
    /// 3. `_free_blocks_memory_usage` < `_max_bytes_in_queue`, remains enough memory to scale up
314
    /// 4. At most scale up `MAX_SCALE_UP_RATIO` times to `_max_thread_num`
315
    void _set_scanner_done();
316
    bool _is_shared_scan_limit_exhausted() const;
317
    // The callback of a parked task's future; runs on whichever thread completed it.
318
    void _resume_parked_task(const std::shared_ptr<ScanTask>& scan_task);
319
320
    RuntimeState* _state = nullptr;
321
    ScanLocalStateBase* _local_state = nullptr;
322
323
    // the comment of same fields in VScanNode
324
    const TupleDescriptor* _output_tuple_desc = nullptr;
325
326
    Status _process_status = Status::OK();
327
    std::atomic_bool _should_stop = false;
328
    std::atomic_bool _is_finished = false;
329
330
    // Lazy-allocated blocks for all scanners to share, for memory reuse.
331
    moodycamel::ConcurrentQueue<BlockUPtr> _free_blocks;
332
333
    int _batch_size;
334
    // The limit from SQL's limit clause
335
    int64_t limit;
336
    // Points to the shared remaining limit on ScanOperatorX, shared across all
337
    // parallel instances and their scanners. -1 means no limit.
338
    std::atomic<int64_t>* _shared_scan_limit = nullptr;
339
340
    int64_t _max_bytes_in_queue = 0;
341
    // _transfer_lock protects _completed_tasks, _pending_tasks, and all other shared state
342
    // accessed by both the scanner thread pool and the operator (get_block_from_queue).
343
    std::mutex _transfer_lock;
344
345
    // Together, _completed_tasks and _in_flight_tasks_num represent all "occupied" concurrency
346
    // slots.  The scheduler uses their sum as the current concurrency:
347
    //
348
    //   current_concurrency = _completed_tasks.size() + _in_flight_tasks_num
349
    //
350
    // Lifecycle of a ScanTask:
351
    //   _pending_tasks  --(submit_scan_task on the TaskExecutor path,
352
    //                      try_get_next_scan_task on the ThreadPool path)--> [thread pool]
353
    //   --(push_completed_scan_task)--> _completed_tasks  --(get_block_from_queue)--> operator
354
    //   After consumption: non-EOS task goes back to _pending_tasks; EOS increments
355
    //   _num_finished_scanners.
356
357
    // Completed scan tasks whose cached_block is ready for the operator to consume.
358
    // Protected by _transfer_lock.  Written by push_completed_scan_task() (scanner thread),
359
    // read/popped by get_block_from_queue() (operator thread).
360
    std::list<std::shared_ptr<ScanTask>> _completed_tasks;
361
362
    // Scanners waiting to be admitted for execution. Stored as a stack (LIFO) so that
363
    // recently-used scanners are re-scheduled first, which is more likely to be cache-friendly.
364
    // Protected by _transfer_lock. Populated in the constructor and when an operator returns a
365
    // non-EOS task; drained by try_get_next_scan_task() or the TaskExecutor scheduler.
366
    std::stack<std::shared_ptr<ScanTask>> _pending_tasks;
367
368
    // True from the start of one Context submission until its runnable starts. The marker may
369
    // remain true when no runnable was retained: a failed submission makes the Context terminal,
370
    // and a submission that threw inside _run_context() publishes the error through the task that
371
    // was already admitted. In both cases the operator observes _process_status, so no further
372
    // submission is attempted. It does not describe scanners executing on workers. Protected by
373
    // _transfer_lock.
374
    bool _is_context_queued = false;
375
376
    // Number of scan tasks currently submitted to the scanner scheduler thread pool
377
    // (i.e. in-flight). Incremented before a task is submitted or directly admitted for
378
    // thread-pool execution, and decremented by push_completed_scan_task() when the worker
379
    // returns it.
380
    // Declared atomic so it can be read without _transfer_lock in non-critical paths,
381
    // but must be read under _transfer_lock whenever combined with _completed_tasks.size()
382
    // to form a consistent concurrency snapshot.
383
    std::atomic_int _in_flight_tasks_num = 0;
384
    // Tasks parked by park_scan_task() until what their scanners wait for is done. They are not in
385
    // flight and occupy no concurrency slot. Protected by _transfer_lock.
386
    int32_t _parked_tasks_num = 0;
387
    // Scanner that is eos or error.
388
    int32_t _num_finished_scanners = 0;
389
    // weak pointer for _scanners, used in stop function
390
    std::vector<std::weak_ptr<ScannerDelegate>> _all_scanners;
391
    std::shared_ptr<RuntimeProfile> _scanner_profile;
392
    // This counter refers to scan operator's local state
393
    RuntimeProfile::Counter* _scanner_memory_used_counter = nullptr;
394
    RuntimeProfile::Counter* _newly_create_free_blocks_num = nullptr;
395
    RuntimeProfile::Counter* _scale_up_scanners_counter = nullptr;
396
    std::shared_ptr<ResourceContext> _resource_ctx;
397
    std::shared_ptr<Dependency> _dependency = nullptr;
398
    std::shared_ptr<doris::TaskHandle> _task_handle;
399
    std::weak_ptr<doris::TaskExecutor> _task_executor;
400
401
    std::atomic<int64_t> _block_memory_usage = 0;
402
403
    // adaptive scan concurrency related
404
405
    ScannerScheduler* _scanner_scheduler = nullptr;
406
    MOCK_REMOVE(const) int32_t _min_scan_concurrency_of_scan_scheduler = 0;
407
    // The overall target of our system is to make full utilization of the resources.
408
    // At the same time, we dont want too many tasks are queued by scheduler, that is not necessary.
409
    // Each scan operator can submit _max_scan_concurrency scanner to scheduelr if scheduler has enough resource.
410
    // So that for a single query, we can make sure it could make full utilization of the resource.
411
    int32_t _max_scan_concurrency = 0;
412
    MOCK_REMOVE(const) int32_t _min_scan_concurrency = 1;
413
414
    std::shared_ptr<ScanTask> _pull_next_scan_task(std::shared_ptr<ScanTask> current_scan_task,
415
                                                   int32_t current_concurrency);
416
417
    int32_t _get_margin(std::unique_lock<std::mutex>& transfer_lock,
418
                        std::unique_lock<std::shared_mutex>& scheduler_lock);
419
420
    // Memory-aware adaptive scheduling
421
    std::shared_ptr<MemLimiter> _scanner_mem_limiter = nullptr;
422
    std::shared_ptr<MemShareArbitrator> _mem_share_arb = nullptr;
423
    std::shared_ptr<ScannerAdaptiveProcessor> _adaptive_processor = nullptr;
424
    const int _ins_idx;
425
    const bool _enable_adaptive_scanners = false;
426
427
    // Adjust scan memory limit based on arbitrator feedback
428
    void _adjust_scan_mem_limit(int64_t old_scanner_mem_bytes, int64_t new_scanner_mem_bytes);
429
430
    // Calculate available scanner count for adaptive scheduling
431
    int _available_pickup_scanner_count();
432
433
    // TODO: Add implementation of runtime_info_feed_back
434
    // adaptive scan concurrency related end
435
};
436
} // namespace doris