mrhhsg commented on code in PR #68571:
URL: https://github.com/apache/doris/pull/68571#discussion_r4121646699


##########
be/src/exec/scan/simplified_scan_scheduler.cpp:
##########
@@ -34,7 +39,112 @@ Status 
TaskExecutorSimplifiedScanScheduler::schedule_scan_task(
 Status ThreadPoolSimplifiedScanScheduler::schedule_scan_task(
         std::shared_ptr<ScannerContext> scanner_ctx, std::shared_ptr<ScanTask> 
current_scan_task,
         std::unique_lock<std::mutex>& transfer_lock) {
-    std::unique_lock<std::shared_mutex> wl(_lock);
-    return scanner_ctx->schedule_scan_task(current_scan_task, transfer_lock, 
wl);
+    // Unlike TaskExecutor, ThreadPool queues a Context runnable. It later 
admits one pending task
+    // under transfer_lock. This bounds queue entries to one per Context even 
when many scanners
+    // become runnable together.
+    DORIS_CHECK(transfer_lock.owns_lock());
+    if (current_scan_task != nullptr) {
+        // The operator has consumed all blocks of a non-EOS scanner, making 
it eligible for
+        // another scan attempt. Queue the scanner first; the Context runnable 
chooses it later.
+        scanner_ctx->push_pending_scan_task(std::move(current_scan_task), 
transfer_lock);
+    }
+    const bool context_queued = scanner_ctx->is_context_queued(transfer_lock);
+    if (context_queued && !_is_stop) {
+        // A queued runnable will see all pending scanners added before it 
obtains transfer_lock.
+        // Submitting another runnable would only duplicate work. Once the 
pool is stopped, its
+        // shutdown drops queued runnables without running them, so the marker 
may never be
+        // cleared; fail below instead of waiting for that runnable forever.
+        return Status::OK();
+    }
+    if (!context_queued && !scanner_ctx->can_admit_scan_task(transfer_lock, 
false)) {

Review Comment:
   Fixed in d783450098d. `ThreadPoolSimplifiedScanScheduler` now has a 
scheduler-wide `_submit_lock` held from the pool-budget check 
(`can_admit_scan_task`) through `submit_func`, so concurrent Contexts are 
serialized across that decision the way the removed `_lock` did: the second 
Context sees the runnable the first one just submitted and defers at its 
minimum concurrency instead of failing on a full pool. The queued-runnable fast 
path returns before taking the lock. Lock order is transfer_lock -> 
`_submit_lock` -> ThreadPool lock and is never inverted (`submit_func` never 
runs the task inline). In addition, `can_admit_scan_task` now reads 
`get_queue_size()` before `get_active_threads()`: each read takes the pool lock 
separately while a worker moves a task from queued to active atomically, so a 
dequeue between the reads now over-counts (defer) instead of under-counting 
(overbooking the last slot). New UTs: 
`thread_pool_budget_check_and_submit_are_atomic_across_contexts` (3 workers, 
 zero queue, 2 parked; a debug point parks the first Context between check and 
submit; without the lock it fails with TOO_MANY_TASKS, verified locally) and 
`thread_pool_admission_reads_queue_before_active_threads`.



##########
be/src/exec/scan/scanner_context.cpp:
##########
@@ -448,17 +452,119 @@ std::string ScannerContext::debug_string() {
     return fmt::format(
             "id: {}, total scanners: {}, pending tasks: {},"
             " _should_stop: {}, _is_finished: {}, free blocks: {},"
-            " limit: {}, _num_running_scanners: {}, _max_thread_num: {},"
+            " limit: {}, _num_running_scanners: {}, _is_context_queued: {},"
+            " _num_finished_scanners: {}, _max_thread_num: {},"
             " _max_bytes_in_queue: {}, query_id: {}",
             ctx_id, _all_scanners.size(), _tasks_queue.size(), _should_stop, 
_is_finished,
-            _free_blocks.size_approx(), limit, _num_scheduled_scanners, 
_max_scan_concurrency,
-            _max_bytes_in_queue, print_id(_query_id));
+            _free_blocks.size_approx(), limit, _num_scheduled_scanners, 
_is_context_queued,
+            _num_finished_scanners, _max_scan_concurrency, _max_bytes_in_queue,
+            print_id(_query_id));
 }
 
 void ScannerContext::_set_scanner_done() {
     _dependency->set_always_ready();
 }
 
+bool ScannerContext::is_context_queued(const std::unique_lock<std::mutex>& 
transfer_lock) const {
+    DORIS_CHECK(transfer_lock.owns_lock());
+    return _is_context_queued;
+}
+
+void ScannerContext::set_context_queued(bool queued,
+                                        const std::unique_lock<std::mutex>& 
transfer_lock) {
+    DORIS_CHECK(transfer_lock.owns_lock());
+    DORIS_CHECK(_is_context_queued != queued);
+    _is_context_queued = queued;
+}
+
+void ScannerContext::set_context_failure(const Status& failure,
+                                         const std::unique_lock<std::mutex>& 
transfer_lock) {
+    DORIS_CHECK(transfer_lock.owns_lock());
+    DORIS_CHECK(!failure.ok());
+    _process_status = failure;
+    _is_finished = true;
+    _set_scanner_done();
+}
+
+void ScannerContext::push_pending_scan_task(std::shared_ptr<ScanTask> 
scan_task,
+                                            const 
std::unique_lock<std::mutex>& transfer_lock) {
+    DORIS_CHECK(transfer_lock.owns_lock());
+    DORIS_CHECK(scan_task != nullptr);
+    DORIS_CHECK(scan_task->cached_blocks.empty());
+    DORIS_CHECK(!scan_task->is_eos());
+    _pending_scanners.push(std::move(scan_task));
+}
+
+bool ScannerContext::can_admit_scan_task(const std::unique_lock<std::mutex>& 
transfer_lock,
+                                         bool admitting_on_worker) const {
+    DORIS_CHECK(transfer_lock.owns_lock());
+    if (done() || _pending_scanners.empty()) {
+        return false;
+    }
+
+    // Blocks waiting in _tasks_queue still occupy a concurrency slot until 
the operator consumes
+    // them. Counting both prevents a fast producer from exceeding the 
per-Context scanner limit.
+    const int32_t current_concurrency =
+            cast_set<int32_t>(_tasks_queue.size()) + _num_scheduled_scanners;
+    // Keep one task progressing whatever the limits are. Otherwise no worker 
can publish a result
+    // and wake the operator to make another scheduling decision.
+    if (current_concurrency == 0) {
+        return true;
+    }
+    // The per-Context ceiling applied by _pull_next_scan_task().
+    if (current_concurrency >= _max_scan_concurrency) {
+        return false;
+    }
+    // In low memory mode _get_margin() limits the number of running scanners.
+    if (low_memory_mode() && _num_scheduled_scanners >= 
low_memory_mode_scanners()) {
+        return false;
+    }
+    // Mirror the scheduler-wide budget of _get_margin(): while the pool has 
slack a Context may
+    // ramp to its maximum. Once it has none, a Context is held at its target 
concurrency, which is
+    // its minimum unless the operator is starving. Both counters are read 
here under
+    // _transfer_lock exactly as the TaskExecutor path reads them. A worker 
admitting the task it
+    // runs itself is already counted as active; like the task _get_margin() 
is about to submit, it
+    // must not count against the budget, otherwise the pool would stop one 
slot short of it.
+    const int32_t busy_scan_slots = _scanner_scheduler->get_active_threads() +
+                                    _scanner_scheduler->get_queue_size() -
+                                    (admitting_on_worker ? 1 : 0);
+    if (busy_scan_slots < _min_scan_concurrency_of_scan_scheduler) {
+        return true;
+    }
+    const int32_t target_scan_concurrency =
+            _scan_starving && _tasks_queue.empty() ? _max_scan_concurrency : 
_min_scan_concurrency;
+    return current_concurrency < target_scan_concurrency;
+}
+
+std::shared_ptr<ScanTask> ScannerContext::try_get_next_scan_task(
+        const std::unique_lock<std::mutex>& transfer_lock, int64_t 
context_queue_wait_ns) {
+    if (!can_admit_scan_task(transfer_lock, true)) {
+        VLOG_DEBUG << fmt::format(
+                "[{}|{}] refuse admission, pending: {}, task queue: {}, 
scheduled: {}, done: {}",
+                print_id(_query_id), ctx_id, _pending_scanners.size(), 
_tasks_queue.size(),
+                _num_scheduled_scanners, done());
+        return nullptr;
+    }
+
+    // Pop and count as scheduled while holding the same lock used by 
completion and consumption.
+    // Thus concurrent Context workers cannot admit the same task or both pass 
the limit check.
+    auto scan_task = _pending_scanners.top();
+    _pending_scanners.pop();
+    // ThreadPool admission bypasses ScannerScheduler::submit(); restart the 
per-scanner wait
+    // timer here so it measures admission-to-execution instead of everything 
since the previous
+    // attempt paused, which would include time the cached blocks waited for 
the operator. The
+    // Context runnable waited for a worker on behalf of this scanner, so 
credit that wait too.
+    if (auto scanner_delegate = scan_task->scanner.lock()) {
+        scanner_delegate->_scanner->start_wait_worker_timer();
+        
scanner_delegate->_scanner->add_wait_worker_time(context_queue_wait_ns);

Review Comment:
   Fixed in d783450098d. `ScanTask` now records `pending_since_ns` when 
`push_pending_scan_task()` returns it to `_pending_scanners` (the only place 
that does so on the ThreadPool path; initial scanners keep 0). `_run_context()` 
passes the runnable's submit and start times, and `try_get_next_scan_task()` 
credits only `max(0, start - max(submit, pending_since_ns))`, which is the 
overlap of the runnable's queue wait with the selected scanner's own pending 
interval. A LIFO-requeued scanner therefore no longer gets its previous 
execution counted as `ScannerWorkerWaitTime` / `PerScannerWaitTime`. 
`thread_pool_admission_state` now covers both cases: a scanner that became 
pending after the runnable was queued (credited 500 instead of 2000) and one 
that was pending before (credited the whole 2000).



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to