This is an automated email from the ASF dual-hosted git repository.

yiguolei pushed a commit to branch branch-4.2
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.2 by this push:
     new c35767f565c [improvement](be) Refactor thread-pool scan scheduling 
(branch-4.2) (#68571)
c35767f565c is described below

commit c35767f565c0d01a58c50d88a3eb0240342f1dd0
Author: Jerry Hu <[email protected]>
AuthorDate: Tue Sep 29 09:40:12 2026 +0800

    [improvement](be) Refactor thread-pool scan scheduling (branch-4.2) (#68571)
    
    ### What problem does this PR solve?
    
    Issue Number: None
    
    Related PR: #67070
    
    Problem Summary:
    
    This is a rewrite of #67070 for branch-4.2, not a cherry-pick. #67070
    builds on scheduler refactors that branch-4.2 does not have, and those
    are intentionally **not** picked here:
    
    - #61617 (ScanTask state machine, `_pending_tasks` / `_completed_tasks`)
    was reverted on this release line by #62191 and stays reverted.
    - #61271 (adaptive scanners and the scan memory limiter), #62222 (limit
    push down to the segment iterator) and #65814 (shared scan limit) are
    large behavior changes of their own.
    
    The admission logic is therefore written against the branch-4.2 model
    (`_pending_scanners`, `_tasks_queue`, `_num_scheduled_scanners`,
    `cached_blocks`). The parts of #67070 that depend on the shared scan
    limit or on adaptive scanners / the memory limiter have no counterpart
    here and are dropped.
    
    The problem is the same on branch-4.2: the ThreadPool scan scheduler
    submits one runnable per scanner and serializes scheduling of every
    Context through a scheduler-wide lock. With many scanners this inflates
    queue occupancy and couples ThreadPool admission to the TaskExecutor
    scheduling logic.
    
    With this change the ThreadPool scheduler:
    
    - queues at most one runnable per `ScannerContext`; the runnable admits
    one pending scanner under the Context transfer lock, submits its
    successor and then scans without the lock;
    - keeps submission failures local to the Context (`TOO_MANY_TASKS` or
    the shutdown error make only that Context terminal);
    - checks the terminal Context state in `get_block_from_queue()` before
    rescheduling a drained scanner;
    - fails the Context at its next scheduling attempt when the pool has
    been stopped while its runnable was still queued, instead of treating
    the runnable that shutdown dropped as still pending.
    
    Admission keeps the limits branch-4.2 applies through `_get_margin()`
    and `_pull_next_scan_task()`: a Context never exceeds
    `_max_scan_concurrency` (cached results count as occupied slots), low
    memory mode caps running scanners, and once the pool has no slack
    (`active + queued >= min_active_scan_threads`) a Context is held at
    `_min_scan_concurrency`, or at `_max_scan_concurrency` while the
    operator is starving. One scanner is always admitted when nothing is
    progressing so the operator can be woken. A worker admitting the task it
    runs itself does not count its own thread against the pool budget,
    matching `_get_margin()`, which runs on the operator thread.
    
    The Context runnable waits in the pool queue on behalf of the scanner it
    admits, so that wait is credited to the scanner's wait-worker time
    (`Scanner::add_wait_worker_time()`), keeping `ScannerWorkerWaitTime` /
    `PerScannerWaitTime` meaningful on the ThreadPool path.
    
    TaskExecutor remains the default scan scheduler and keeps its behavior.
    Generic ThreadPool behavior is unchanged. `debug_string()` now reports
    `_is_context_queued` and `_num_finished_scanners`, and the ThreadPool
    path logs admission, refusal and runnable submission at `VLOG_DEBUG`.
    
    Validation:
    
    - `./run-be-ut.sh --run
    --filter="ScannerContextTest.*:ThreadPoolTest.*:*ScanScheduler*:*Scanner*"`
    (132 tests passed)
    - The new ThreadPool tests repeated 50 times with `GTEST_REPEAT=50`
    without failure.
    - clang-format 16 check and `git diff --check`.
    
    ### Release note
    
    None
    
    ### Check List (For Author)
    
    - Test <!-- At least one of them must be included. -->
        - [ ] Regression test
        - [x] Unit Test
        - [ ] Manual test (add detailed scripts or steps below)
        - [ ] No need to test or manual test. Explain why:
    - [ ] This is a refactor/code format and no logic has been changed.
            - [ ] Previous test can cover this change.
            - [ ] No code files have been changed.
            - [ ] Other reason <!-- Add your reason?  -->
    
    - Behavior changed:
        - [ ] No.
    - [x] Yes. The opt-in ThreadPool scan scheduler queues and admits work
    per `ScannerContext`.
    
    - Does this need documentation?
        - [x] No.
    - [ ] Yes. <!-- Add document PR link here. eg:
    https://github.com/apache/doris-website/pull/1214 -->
    
    ### Check List (For Reviewer who merge this PR)
    
    - [ ] Confirm the release note
    - [ ] Confirm test cases
    - [ ] Confirm document
    - [ ] Add branch pick label <!-- Add branch pick label that this PR
    should merge into -->
---
 be/src/exec/scan/scanner.h                     |    4 +
 be/src/exec/scan/scanner_context.cpp           |  130 ++-
 be/src/exec/scan/scanner_context.h             |   50 +-
 be/src/exec/scan/scanner_scheduler.cpp         |   28 +-
 be/src/exec/scan/scanner_scheduler.h           |   14 +-
 be/src/exec/scan/simplified_scan_scheduler.cpp |  127 ++-
 be/test/exec/scan/scanner_context_test.cpp     | 1040 ++++++++++++++++++++++++
 7 files changed, 1370 insertions(+), 23 deletions(-)

diff --git a/be/src/exec/scan/scanner.h b/be/src/exec/scan/scanner.h
index a389e029a58..6aaa1ec6940 100644
--- a/be/src/exec/scan/scanner.h
+++ b/be/src/exec/scan/scanner.h
@@ -169,6 +169,10 @@ public:
 
     void update_wait_worker_timer() { _scanner_wait_worker_timer += 
_watch.elapsed_time(); }
 
+    // Credit wait time spent before start_wait_worker_timer() could be 
called, e.g. while the
+    // ThreadPool Context runnable that admits this scanner was queued for a 
worker.
+    void add_wait_worker_time(int64_t wait_ns) { _scanner_wait_worker_timer += 
wait_ns; }
+
     int64_t get_scanner_wait_worker_timer() const { return 
_scanner_wait_worker_timer; }
 
     void update_scan_cpu_timer();
diff --git a/be/src/exec/scan/scanner_context.cpp 
b/be/src/exec/scan/scanner_context.cpp
index 279040c5168..9d5a53005d3 100644
--- a/be/src/exec/scan/scanner_context.cpp
+++ b/be/src/exec/scan/scanner_context.cpp
@@ -22,6 +22,7 @@
 #include <glog/logging.h>
 #include <zconf.h>
 
+#include <algorithm>
 #include <cstdint>
 #include <ctime>
 #include <memory>
@@ -278,6 +279,8 @@ Status ScannerContext::get_block_from_queue(RuntimeState* 
state, Block* block, b
     }
 
     std::shared_ptr<ScanTask> scan_task = nullptr;
+    // Whether all cached blocks of scan_task were consumed, so the scanner 
needs scheduling again.
+    bool scan_task_drained = false;
 
     if (!_tasks_queue.empty() && !done()) {
         // https://en.cppreference.com/w/cpp/container/list/front
@@ -323,19 +326,21 @@ Status ScannerContext::get_block_from_queue(RuntimeState* 
state, Block* block, b
                 _scan_starving = _tasks_queue.empty() &&
                                  _num_finished_scanners < 
cast_set<int32_t>(_all_scanners.size()) &&
                                  (_num_scheduled_scanners > 0 || 
!_pending_scanners.empty());
-                RETURN_IF_ERROR(
-                        
_scanner_scheduler->schedule_scan_task(shared_from_this(), nullptr, l));
             } else {
                 _scan_starving = _tasks_queue.empty();
-                RETURN_IF_ERROR(
-                        
_scanner_scheduler->schedule_scan_task(shared_from_this(), scan_task, l));
             }
+            scan_task_drained = true;
         }
     }
 
     if (_num_finished_scanners == _all_scanners.size() && 
_tasks_queue.empty()) {
         _set_scanner_done();
         _is_finished = true;
+    } else if (scan_task_drained) {
+        // Check the terminal state before scheduling more work, so a 
completed Context never
+        // submits another runnable that could only fail on a full scanner 
pool.
+        RETURN_IF_ERROR(_scanner_scheduler->schedule_scan_task(
+                shared_from_this(), scan_task->is_eos() ? nullptr : scan_task, 
l));
     }
 
     *eos = done();
@@ -448,17 +453,128 @@ 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());
+    scan_task->pending_since_ns = MonotonicNanos();
+    _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.
+    // The two counters are read under separate pool locks, and a worker moves 
a task from queued
+    // to active atomically. Reading the queue first makes a concurrent 
dequeue count that task
+    // twice rather than not at all, so a torn read defers instead of 
overbooking the last slot.
+    const int32_t queued_scan_tasks = _scanner_scheduler->get_queue_size();
+    const int32_t active_scan_threads = 
_scanner_scheduler->get_active_threads();
+    const int32_t busy_scan_slots =
+            active_scan_threads + queued_scan_tasks - (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_submit_time_ns,
+        int64_t context_start_time_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(), so credit the 
time the Context
+    // runnable waited for a worker here; the caller restarts the per-scanner 
wait timer right
+    // before executing the task. The runnable waited on behalf of this 
scanner only while the
+    // scanner was pending: with LIFO re-admission the runnable may have been 
queued before the
+    // scanner's previous attempt ran, so do not credit the part of the 
runnable's wait before it
+    // was pending.
+    if (auto scanner_delegate = scan_task->scanner.lock()) {
+        const int64_t wait_start_ns = std::max(context_submit_time_ns, 
scan_task->pending_since_ns);
+        scanner_delegate->_scanner->add_wait_worker_time(
+                std::max<int64_t>(0, context_start_time_ns - wait_start_ns));
+    }
+    ++_num_scheduled_scanners;
+    VLOG_DEBUG << fmt::format("[{}|{}] admit scanner, pending: {}, task queue: 
{}, scheduled: {}",
+                              print_id(_query_id), ctx_id, 
_pending_scanners.size(),
+                              _tasks_queue.size(), _num_scheduled_scanners);
+    return scan_task;
+}
+
 void ScannerContext::update_peak_running_scanner(int num) {
     _local_state->_peak_running_scanner->add(num);
 }
diff --git a/be/src/exec/scan/scanner_context.h 
b/be/src/exec/scan/scanner_context.h
index ba1c3235660..cdb0ca6ff79 100644
--- a/be/src/exec/scan/scanner_context.h
+++ b/be/src/exec/scan/scanner_context.h
@@ -79,6 +79,9 @@ public:
     std::weak_ptr<ScannerDelegate> scanner;
     std::list<std::pair<BlockUPtr, size_t>> cached_blocks;
     bool is_first_schedule = true;
+    // MonotonicNanos() at which the ThreadPool scheduler returned this task 
to _pending_scanners,
+    // or 0 if it has been pending since the Context was created. Protected by 
_transfer_lock.
+    int64_t pending_since_ns = 0;
     // Use weak_ptr to avoid circular references and potential memory leaks 
with SplitRunner.
     // ScannerContext only needs to observe the lifetime of SplitRunner 
without owning it.
     // When SplitRunner is destroyed, split_runner.lock() will return nullptr, 
ensuring safe access.
@@ -193,6 +196,44 @@ public:
                               std::unique_lock<std::mutex>& transfer_lock,
                               std::unique_lock<std::shared_mutex>& 
scheduler_lock);
 
+    // Context scheduling and operator consumption share this lock so 
queue-state changes and task
+    // admission form one atomic decision. For example, two worker callbacks 
cannot both admit the
+    // last available concurrency slot.
+    std::mutex& transfer_lock() { return _transfer_lock; }
+
+    // One Context submission represents many pending scanners in the 
ThreadPool scheduler.
+    // Keeping this separate from scanner execution prevents duplicate 
runnables from accumulating.
+    bool is_context_queued(const std::unique_lock<std::mutex>& transfer_lock) 
const;
+    // Transition the Context runnable's queue state. The caller must hold 
_transfer_lock.
+    void set_context_queued(bool queued, const std::unique_lock<std::mutex>& 
transfer_lock);
+
+    // Publish a scheduler failure and make the Context terminal. The caller 
must hold
+    // _transfer_lock so a retained ThreadPool callback cannot admit another 
scanner concurrently.
+    void set_context_failure(const Status& failure,
+                             const std::unique_lock<std::mutex>& 
transfer_lock);
+
+    // Return a scanner to the admission queue after its blocks are consumed. 
It may not own cached
+    // blocks and may not be EOS: EOS scanners are terminal and must not run 
again.
+    void push_pending_scan_task(std::shared_ptr<ScanTask> scan_task,
+                                const std::unique_lock<std::mutex>& 
transfer_lock);
+
+    // Return whether a Context worker can currently admit one pending 
scanner. This check has no
+    // side effects, so the scheduler can avoid submitting a runnable that 
would immediately exit.
+    // It always admits one scanner when nothing is progressing so the 
operator can be woken, and
+    // otherwise applies the same limits as _get_margin() and 
_pull_next_scan_task() on the
+    // TaskExecutor path. `admitting_on_worker` is true when the caller is the 
pool worker that will
+    // run the admitted task itself. The caller must hold _transfer_lock.
+    bool can_admit_scan_task(const std::unique_lock<std::mutex>& transfer_lock,
+                             bool admitting_on_worker) const;
+
+    // Atomically check whether this context can start another scan task, move 
one task from
+    // pending to scheduled, and return it. It is called by the pool worker 
that will run the task;
+    // the worker's Context runnable was submitted at `context_submit_time_ns` 
and started at
+    // `context_start_time_ns` (both MonotonicNanos()). The caller must hold 
_transfer_lock.
+    std::shared_ptr<ScanTask> try_get_next_scan_task(
+            const std::unique_lock<std::mutex>& transfer_lock, int64_t 
context_submit_time_ns,
+            int64_t context_start_time_ns);
+
 protected:
     /// Four criteria to determine whether to increase the parallelism of the 
scanners
     /// 1. It ran for at least `SCALE_UP_DURATION` ms after last scale up
@@ -225,7 +266,14 @@ protected:
     int64_t _max_bytes_in_queue = 0;
     // Using stack so that we can resubmit scanner in a LIFO order, maybe more 
cache friendly
     std::stack<std::shared_ptr<ScanTask>> _pending_scanners;
-    // Scanner that is submitted to the scheduler.
+    // True from the start of one Context submission until its runnable 
starts. The marker may
+    // remain true when no runnable was retained: a failed submission makes 
the Context terminal,
+    // and a submission that threw inside _run_context() publishes the error 
through the task that
+    // was already admitted. In both cases the operator observes 
_process_status, so no further
+    // submission is attempted. It does not describe scanners executing on 
workers. Only used by
+    // the ThreadPool scheduler. Protected by _transfer_lock.
+    bool _is_context_queued = false;
+    // Scanner that is submitted to the scheduler, or directly admitted for 
thread-pool execution.
     std::atomic_int _num_scheduled_scanners = 0;
     // Scanner that is eos or error.
     int32_t _num_finished_scanners = 0;
diff --git a/be/src/exec/scan/scanner_scheduler.cpp 
b/be/src/exec/scan/scanner_scheduler.cpp
index a6d35f3a0b7..71f97e318e7 100644
--- a/be/src/exec/scan/scanner_scheduler.cpp
+++ b/be/src/exec/scan/scanner_scheduler.cpp
@@ -71,17 +71,7 @@ Status 
ScannerScheduler::submit(std::shared_ptr<ScannerContext> ctx,
     TabletStorageType type = scanner_delegate->_scanner->get_storage_type();
     auto sumbit_task = [&]() {
         auto work_func = [scanner_ref = scan_task, ctx]() {
-            auto status = [&] {
-                RETURN_IF_CATCH_EXCEPTION(_scanner_scan(ctx, scanner_ref));
-                return Status::OK();
-            }();
-
-            if (!status.ok()) {
-                scanner_ref->set_status(status);
-                ctx->push_back_scan_task(scanner_ref);
-                return true;
-            }
-            return scanner_ref->is_eos();
+            return execute_scan_task(ctx, scanner_ref);
         };
         SimplifiedScanTask simple_scan_task = {work_func, ctx, scan_task};
         return this->submit_scan_task(simple_scan_task);
@@ -100,6 +90,22 @@ Status 
ScannerScheduler::submit(std::shared_ptr<ScannerContext> ctx,
     return Status::OK();
 }
 
+bool ScannerScheduler::execute_scan_task(const 
std::shared_ptr<ScannerContext>& ctx,
+                                         const std::shared_ptr<ScanTask>& 
scan_task) {
+    // Both schedulers admit tasks differently, but exceptions must always 
become a completed task
+    // so the operator observes the error and releases the task's scheduled 
concurrency slot.
+    auto status = [&] {
+        RETURN_IF_CATCH_EXCEPTION(_scanner_scan(ctx, scan_task));
+        return Status::OK();
+    }();
+    if (!status.ok()) {
+        scan_task->set_status(status);
+        ctx->push_back_scan_task(scan_task);
+        return true;
+    }
+    return scan_task->is_eos();
+}
+
 void handle_reserve_memory_failure(RuntimeState* state, 
std::shared_ptr<ScannerContext> ctx,
                                    const Status& st, size_t reserve_size) {
     ctx->clear_free_blocks();
diff --git a/be/src/exec/scan/scanner_scheduler.h 
b/be/src/exec/scan/scanner_scheduler.h
index 66d5fd55f65..f5dbc65d0e9 100644
--- a/be/src/exec/scan/scanner_scheduler.h
+++ b/be/src/exec/scan/scanner_scheduler.h
@@ -19,6 +19,7 @@
 
 #include <atomic>
 #include <memory>
+#include <mutex>
 
 #include "common/be_mock_util.h"
 #include "common/status.h"
@@ -138,7 +139,11 @@ public:
 protected:
     int _min_active_scan_threads;
 
-private:
+    // Execute one admitted task for both scheduler implementations. The 
return value is consumed
+    // by TaskExecutor to distinguish terminal EOS/error tasks from scanners 
that remain runnable.
+    static bool execute_scan_task(const std::shared_ptr<ScannerContext>& ctx,
+                                  const std::shared_ptr<ScanTask>& scan_task);
+
     static void _scanner_scan(std::shared_ptr<ScannerContext> ctx,
                               std::shared_ptr<ScanTask> scan_task);
 
@@ -240,12 +245,17 @@ public:
                               std::unique_lock<std::mutex>& transfer_lock) 
override;
 
 private:
+    // `submit_time_ns` is the MonotonicNanos() at which the runnable was 
submitted.
+    void _run_context(std::shared_ptr<ScannerContext> scanner_ctx, int64_t 
submit_time_ns);
+
     std::unique_ptr<ThreadPool> _scan_thread_pool;
     std::atomic<bool> _is_stop;
     std::weak_ptr<CgroupCpuCtl> _cgroup_cpu_ctl;
     std::string _sched_name;
     std::string _workload_group;
-    std::shared_mutex _lock;
+    // Serializes the scheduler-wide budget check and the Context submission 
across Contexts.
+    // Acquired while holding a Context's _transfer_lock; never the other way 
around.
+    std::mutex _submit_lock;
 };
 
 class TaskExecutorSimplifiedScanScheduler final : public ScannerScheduler {
diff --git a/be/src/exec/scan/simplified_scan_scheduler.cpp 
b/be/src/exec/scan/simplified_scan_scheduler.cpp
index 275461ff1dd..3c3fdedb217 100644
--- a/be/src/exec/scan/simplified_scan_scheduler.cpp
+++ b/be/src/exec/scan/simplified_scan_scheduler.cpp
@@ -17,8 +17,15 @@
 
 #include <memory>
 
+#include "common/exception.h"
+#include "common/logging.h"
+#include "exec/scan/scan_node.h"
+#include "exec/scan/scanner.h"
 #include "exec/scan/scanner_context.h"
 #include "exec/scan/scanner_scheduler.h"
+#include "runtime/thread_context.h"
+#include "util/debug_points.h"
+#include "util/time.h"
 
 namespace doris {
 class ScannerDelegate;
@@ -34,7 +41,123 @@ 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();
+    }
+    // The pool budget read by can_admit_scan_task() is shared by every 
Context. Hold the
+    // scheduler-wide lock from that check through submission, otherwise two 
Contexts can both see
+    // the last free slot and the loser fails its query on a full pool instead 
of deferring.
+    std::lock_guard<std::mutex> submit_lock(_submit_lock);
+    if (!context_queued && !scanner_ctx->can_admit_scan_task(transfer_lock, 
false)) {
+        // No runnable is needed when the Context has no pending scanner or 
its concurrency slots
+        // are occupied. A completion only wakes the operator; the operator's 
next consumption in
+        // get_block_from_queue(), or the successor submit in _run_context(), 
retries scheduling.
+        return Status::OK();
+    }
+
+    if (_is_stop) {
+        Status failure = Status::InternalError<false>("scanner pool {} is 
shutdown.", _sched_name);
+        scanner_ctx->set_context_failure(failure, transfer_lock);
+        return failure;
+    }
+
+    
DBUG_EXECUTE_IF("ThreadPoolSimplifiedScanScheduler.schedule_scan_task.before_submit",
+                    DBUG_RUN_CALLBACK());
+    // ThreadPool::submit_func() may return an error after retaining the 
runnable. Set the marker
+    // before submission so either outcome is safe: a retained callback clears 
it, while a truly
+    // rejected callback leaves a terminal Context that no longer needs 
rescheduling.
+    scanner_ctx->set_context_queued(true, transfer_lock);
+    Status status =
+            _scan_thread_pool->submit_func([this, scanner_ctx, submit_time_ns 
= MonotonicNanos()] {
+                _run_context(scanner_ctx, submit_time_ns);
+            });
+    if (status.ok()) {
+        VLOG_DEBUG << "submit context runnable to scanner pool " << 
_sched_name << ", "
+                   << scanner_ctx->debug_string();
+        return Status::OK();
+    }
+    Status failure =
+            Status::TooManyTasks("Failed to submit scanner context {} to 
scanner pool, reason: {}",
+                                 scanner_ctx->ctx_id, status.msg());
+    scanner_ctx->set_context_failure(failure, transfer_lock);
+    return failure;
+}
+
+void 
ThreadPoolSimplifiedScanScheduler::_run_context(std::shared_ptr<ScannerContext> 
scanner_ctx,
+                                                     int64_t submit_time_ns) {
+    const int64_t start_time_ns = MonotonicNanos();
+    std::shared_ptr<ScanTask> scan_task;
+    Status admission_status = [&]() -> Status {
+        std::unique_lock<std::mutex> 
transfer_lock(scanner_ctx->transfer_lock());
+        scanner_ctx->set_context_queued(false, transfer_lock);
+
+        auto task_execution_lock = scanner_ctx->task_exec_ctx();
+        if (task_execution_lock == nullptr) {
+            return Status::OK();
+        }
+#ifndef BE_TEST
+        // Attach before admission: allocations below (for example the 
resubmitted
+        // FunctionRunnable) must charge the query rather than the orphan 
tracker. Scoped to this
+        // lambda so it detaches before execute_scan_task(), whose 
_scanner_scan() attaches again.
+        SCOPED_ATTACH_TASK(scanner_ctx->state());
+#endif
+        Status status = [&]() -> Status {
+            RETURN_IF_CATCH_EXCEPTION({
+                // Admission checks cached results, scheduled tasks, and pool 
slack while holding
+                // transfer_lock.
+                scan_task = scanner_ctx->try_get_next_scan_task(transfer_lock, 
submit_time_ns,
+                                                                start_time_ns);
+                if (scan_task != nullptr) {
+                    
DBUG_EXECUTE_IF("ThreadPoolSimplifiedScanScheduler._run_context.inject_failure",
+                                    {
+                                        throw 
Exception(ErrorCode::INTERNAL_ERROR,
+                                                        "injected admission 
failure");
+                                    });
+                    // Queue the next Context runnable before executing this 
task. Holding
+                    // transfer_lock keeps the admission decision atomic.
+                    RETURN_IF_ERROR(schedule_scan_task(scanner_ctx, nullptr, 
transfer_lock));
+                }
+            });
+            return Status::OK();
+        }();
+        if (!status.ok() && scan_task == nullptr) {
+            scanner_ctx->set_context_failure(status, transfer_lock);
+        }
+        return status;
+    }();
+    if (!admission_status.ok()) [[unlikely]] {
+        if (scan_task != nullptr) {
+            // The scanner was admitted before the failure. Publish it to 
release the scheduled
+            // slot and make the error observable by the operator.
+            scan_task->set_status(admission_status);
+            scanner_ctx->push_back_scan_task(scan_task);
+        }
+        return;
+    }
+    if (scan_task == nullptr) {
+        return;
+    }
+    // The scanner already runs on this worker. Start its wait timer only now 
so it does not count
+    // the successor submission above, during which the pool may synchronously 
create a thread.
+    if (auto scanner_delegate = scan_task->scanner.lock()) {
+        scanner_delegate->_scanner->start_wait_worker_timer();
+    }
+    // The scan runs without transfer_lock so the operator and other Context 
workers can continue
+    // consuming results and admitting work. Completion reacquires the lock 
before publishing.
+    execute_scan_task(scanner_ctx, scan_task);
 }
 } // namespace doris
diff --git a/be/test/exec/scan/scanner_context_test.cpp 
b/be/test/exec/scan/scanner_context_test.cpp
index 8b2097716de..1c77e0abf33 100644
--- a/be/test/exec/scan/scanner_context_test.cpp
+++ b/be/test/exec/scan/scanner_context_test.cpp
@@ -23,11 +23,17 @@
 #include <gen_cpp/Types_types.h>
 #include <gtest/gtest.h>
 
+#include <atomic>
+#include <chrono>
+#include <functional>
 #include <list>
 #include <memory>
 #include <mutex>
+#include <thread>
 #include <tuple>
+#include <vector>
 
+#include "common/config.h"
 #include "common/object_pool.h"
 #include "core/block/block.h"
 #include "exec/operator/olap_scan_operator.h"
@@ -35,12 +41,67 @@
 #include "exec/scan/mock_simplified_scan_scheduler.h"
 #include "exec/scan/olap_scanner.h"
 #include "exec/scan/scan_node.h"
+#include "exec/scan/scanner.h"
 #include "exec/scan/scanner_scheduler.h"
 #include "runtime/descriptors.h"
 #include "runtime/query_context.h"
+#include "runtime/task_execution_context.h"
 #include "testutil/mock/mock_runtime_state.h"
+#include "util/countdown_latch.h"
+#include "util/debug_points.h"
+#include "util/defer_op.h"
 
 namespace doris {
+// A scanner that produces `blocks_per_scanner` one-row blocks and then 
reports EOS, without any
+// tablet or file behind it. It lets the ThreadPool scheduler chain (admit -> 
execute -> publish ->
+// consume -> re-admit) run end to end in a unit test. `overlap` is counted 
down the first time two
+// attempts run concurrently; the first attempt waits for it so the peak 
concurrency observed by the
+// test does not depend on timing.
+class ChainMockScanner : public Scanner {
+public:
+    ChainMockScanner(RuntimeState* state, ScanLocalStateBase* local_state, 
RuntimeProfile* profile,
+                     int blocks_per_scanner, std::atomic<int>* running,
+                     std::atomic<int>* peak_running, CountDownLatch* overlap)
+            : Scanner(state, local_state, -1, profile),
+              _blocks_left(blocks_per_scanner),
+              _running(running),
+              _peak_running(peak_running),
+              _overlap(overlap) {}
+
+protected:
+    Status _get_block_impl(RuntimeState* /*state*/, Block* block, bool* eof) 
override {
+        const int running = ++*_running;
+        Defer done([&] { --*_running; });
+        int peak = _peak_running->load();
+        while (running > peak && !_peak_running->compare_exchange_weak(peak, 
running)) {
+        }
+        if (running >= 2) {
+            _overlap->count_down();
+        } else {
+            // Bounded so a scheduler that never admits a second scanner fails 
the test instead of
+            // hanging it.
+            static_cast<void>(_overlap->wait_for(std::chrono::seconds(5)));
+        }
+        if (_blocks_left == 0) {
+            *eof = true;
+            return Status::OK();
+        }
+        --_blocks_left;
+        block->get_by_position(0).column->assert_mutable()->insert_default();
+        *eof = false;
+        return Status::OK();
+    }
+
+    // The local state in these tests has no profile counters.
+    void _collect_profile_before_close() override {}
+
+private:
+    int _blocks_left;
+    std::atomic<int>* _running;
+    std::atomic<int>* _peak_running;
+    CountDownLatch* _overlap;
+};
+
 class ScannerContextTest : public testing::Test {
 public:
     void SetUp() override {
@@ -860,4 +921,983 @@ TEST_F(ScannerContextTest, get_block_from_queue) {
     EXPECT_EQ(scanner_context->_num_finished_scanners, 1);
 }
 
+TEST_F(ScannerContextTest, thread_pool_admission_state) {
+    const int parallel_tasks = 1;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+    std::list<std::shared_ptr<ScannerDelegate>> scanners {
+            std::make_shared<ScannerDelegate>(scanner)};
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+
+    // An idle pool: admission is bounded only by the per-Context limit.
+    std::unique_ptr<MockSimplifiedScanScheduler> scheduler =
+            std::make_unique<MockSimplifiedScanScheduler>(cgroup_cpu_ctl);
+    EXPECT_CALL(*scheduler, 
get_active_threads()).WillRepeatedly(testing::Return(0));
+    EXPECT_CALL(*scheduler, 
get_queue_size()).WillRepeatedly(testing::Return(0));
+    scanner_context->_scanner_scheduler = scheduler.get();
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 20;
+
+    std::unique_lock<std::mutex> 
context_transfer_lock(scanner_context->transfer_lock());
+    scanner_context->_pending_scanners = 
std::stack<std::shared_ptr<ScanTask>>();
+    scanner_context->_tasks_queue.clear();
+    scanner_context->_num_scheduled_scanners = 0;
+    // Even if the effective limit is temporarily zero, one pending task must 
run so it can publish
+    // a block or EOS and prevent the Context from stalling.
+    scanner_context->_max_scan_concurrency = 0;
+
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(context_transfer_lock, 
false));
+
+    // A consumed non-EOS scanner must be eligible for another Context 
admission.
+    auto consumed_task = std::make_shared<ScanTask>(scanners.front());
+    scanner_context->push_pending_scan_task(consumed_task, 
context_transfer_lock);
+    EXPECT_TRUE(scanner_context->can_admit_scan_task(context_transfer_lock, 
false));
+
+    EXPECT_FALSE(scanner_context->is_context_queued(context_transfer_lock));
+    scanner_context->set_context_queued(true, context_transfer_lock);
+    EXPECT_TRUE(scanner_context->is_context_queued(context_transfer_lock));
+    scanner_context->set_context_queued(false, context_transfer_lock);
+
+    // The Context admits exactly one scanner and counts it as scheduled. Its 
runnable was queued
+    // (at 1000) before the scanner was returned to pending (at 2500), as with 
LIFO re-admission
+    // of a scanner whose previous attempt ran while the runnable waited. Only 
the part of the
+    // runnable's wait during which the scanner was pending is its wait-worker 
time.
+    EXPECT_GT(consumed_task->pending_since_ns, 0);
+    consumed_task->pending_since_ns = 2500;
+    int64_t wait_worker_time = scanner->get_scanner_wait_worker_timer();
+    auto admitted_task = 
scanner_context->try_get_next_scan_task(context_transfer_lock, 1000, 3000);
+    EXPECT_EQ(admitted_task, consumed_task);
+    EXPECT_EQ(scanner->get_scanner_wait_worker_timer(), wait_worker_time + 
500);
+    EXPECT_EQ(scanner_context->_num_scheduled_scanners, 1);
+    EXPECT_TRUE(scanner_context->_pending_scanners.empty());
+
+    // A scanner already pending when the runnable was queued waited for the 
whole queue time.
+    scanner_context->_num_scheduled_scanners = 0;
+    auto early_task = std::make_shared<ScanTask>(scanners.front());
+    scanner_context->push_pending_scan_task(early_task, context_transfer_lock);
+    early_task->pending_since_ns = 500;
+    wait_worker_time = scanner->get_scanner_wait_worker_timer();
+    EXPECT_EQ(scanner_context->try_get_next_scan_task(context_transfer_lock, 
1000, 3000),
+              early_task);
+    EXPECT_EQ(scanner->get_scanner_wait_worker_timer(), wait_worker_time + 
2000);
+    EXPECT_EQ(scanner_context->_num_scheduled_scanners, 1);
+
+    auto blocked_task = 
std::make_shared<ScanTask>(std::make_shared<ScannerDelegate>(scanner));
+    scanner_context->push_pending_scan_task(blocked_task, 
context_transfer_lock);
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(context_transfer_lock, 
false));
+    EXPECT_EQ(scanner_context->try_get_next_scan_task(context_transfer_lock, 
0, 0), nullptr);
+
+    // A stopped Context admits nothing even when nothing is progressing.
+    scanner_context->_num_scheduled_scanners = 0;
+    EXPECT_TRUE(scanner_context->can_admit_scan_task(context_transfer_lock, 
false));
+    scanner_context->_should_stop = true;
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(context_transfer_lock, 
false));
+}
+
+TEST_F(ScannerContextTest, 
thread_pool_admission_reads_queue_before_active_threads) {
+    const int parallel_tasks = 2;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < 2; ++i) {
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+
+    // A worker moves a task from queued to active under one pool lock, while 
admission reads the
+    // two counters under separate locks. Reading the queue first means a 
dequeue in between is
+    // counted twice (defer) rather than missed (overbook the last slot and 
fail the submission).
+    std::unique_ptr<MockSimplifiedScanScheduler> scheduler =
+            std::make_unique<MockSimplifiedScanScheduler>(cgroup_cpu_ctl);
+    {
+        testing::InSequence sequence;
+        EXPECT_CALL(*scheduler, get_queue_size()).WillOnce(testing::Return(1));
+        EXPECT_CALL(*scheduler, 
get_active_threads()).WillOnce(testing::Return(3));
+    }
+    scanner_context->_scanner_scheduler = scheduler.get();
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 4;
+    scanner_context->_max_scan_concurrency = parallel_tasks;
+    scanner_context->_min_scan_concurrency = 1;
+    scanner_context->_scan_starving = false;
+    scanner_context->_num_scheduled_scanners = 1;
+
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    ASSERT_FALSE(scanner_context->_pending_scanners.empty());
+    // The over-counted pool looks full, so the Context stays at its minimum 
concurrency.
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(transfer_lock, false));
+}
+
+TEST_F(ScannerContextTest, thread_pool_admission_follows_margin_limits) {
+    const int parallel_tasks = 4;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    // A single-worker pool whose worker is parked: active + queued == 1, i.e. 
the pool has no
+    // slack once the scheduler-wide budget is 1.
+    ThreadPoolSimplifiedScanScheduler scheduler("saturated_pool_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(1, 1, 4, 1).ok());
+    CountDownLatch task_started(1);
+    CountDownLatch release_task(1);
+    Defer cleanup = [&] {
+        release_task.count_down();
+        scheduler.stop();
+    };
+    ASSERT_TRUE(scheduler
+                        .submit_scan_task(SimplifiedScanTask(
+                                [&] {
+                                    task_started.count_down();
+                                    release_task.wait();
+                                    return true;
+                                },
+                                nullptr, nullptr))
+                        .ok());
+    ASSERT_TRUE(task_started.wait_for(std::chrono::seconds(5)));
+    ASSERT_EQ(scheduler.get_active_threads(), 1);
+    ASSERT_EQ(scheduler.get_queue_size(), 0);
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < 4; ++i) {
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+    scanner_context->_scanner_scheduler = &scheduler;
+    scanner_context->_min_scan_concurrency = 1;
+
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    ASSERT_EQ(scanner_context->_max_scan_concurrency, parallel_tasks);
+
+    // Saturated and not starving: the minimum concurrency is the ceiling, as 
_get_margin()
+    // enforces on the TaskExecutor path.
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 1;
+    scanner_context->_scan_starving = false;
+    scanner_context->_num_scheduled_scanners = 0;
+    EXPECT_TRUE(scanner_context->can_admit_scan_task(transfer_lock, false));
+    scanner_context->_num_scheduled_scanners = 1;
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(transfer_lock, false));
+    // A pool worker admitting the task it runs itself does not count its own 
thread, so the
+    // parked worker alone leaves one slot of slack, just as it does for the 
operator thread when
+    // the budget is two.
+    EXPECT_TRUE(scanner_context->can_admit_scan_task(transfer_lock, true));
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 2;
+    EXPECT_TRUE(scanner_context->can_admit_scan_task(transfer_lock, false));
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 1;
+
+    // A larger minimum raises the saturated ceiling accordingly.
+    scanner_context->_min_scan_concurrency = 2;
+    EXPECT_TRUE(scanner_context->can_admit_scan_task(transfer_lock, false));
+    scanner_context->_num_scheduled_scanners = 2;
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(transfer_lock, false));
+
+    // A starving operator with an empty result queue lets the Context ramp to 
its maximum.
+    scanner_context->_scan_starving = true;
+    EXPECT_TRUE(scanner_context->can_admit_scan_task(transfer_lock, false));
+    scanner_context->_num_scheduled_scanners = parallel_tasks;
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(transfer_lock, false));
+
+    // Cached results end starvation for the margin, and they occupy a 
concurrency slot.
+    scanner_context->_num_scheduled_scanners = 1;
+    scanner_context->_tasks_queue.push_back(
+            
std::make_shared<ScanTask>(std::make_shared<ScannerDelegate>(scanner)));
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(transfer_lock, false));
+    scanner_context->_tasks_queue.clear();
+
+    // With slack the Context may ramp to its maximum again, but never beyond 
it.
+    scanner_context->_scan_starving = false;
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 20;
+    scanner_context->_num_scheduled_scanners = parallel_tasks - 1;
+    EXPECT_TRUE(scanner_context->can_admit_scan_task(transfer_lock, false));
+    scanner_context->_num_scheduled_scanners = parallel_tasks;
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(transfer_lock, false));
+}
+
+TEST_F(ScannerContextTest, debug_string_reports_context_queue_state) {
+    const int parallel_tasks = 3;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+    std::list<std::shared_ptr<ScannerDelegate>> scanners {
+            std::make_shared<ScannerDelegate>(scanner)};
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, 7, scan_dependency, parallel_tasks);
+
+    // Every value is distinct so a misplaced placeholder is visible in the 
output.
+    scanner_context->_num_scheduled_scanners = 2;
+    scanner_context->_is_context_queued = true;
+    scanner_context->_num_finished_scanners = 5;
+
+    const std::string debug = scanner_context->debug_string();
+    EXPECT_NE(debug.find("limit: 7, _num_running_scanners: 2, 
_is_context_queued: true, "
+                         "_num_finished_scanners: 5, _max_thread_num: 3,"),
+              std::string::npos)
+            << debug;
+}
+
+TEST_F(ScannerContextTest, thread_pool_submit_failure_policy) {
+    const int parallel_tasks = 2;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < 2; ++i) {
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+
+    // One worker, zero queue capacity, worker occupied: every submit_func() 
is rejected.
+    ThreadPoolSimplifiedScanScheduler scheduler("submit_failure_policy_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(1, 1, 0, 1).ok());
+    CountDownLatch task_started(1);
+    CountDownLatch release_task(1);
+    Defer cleanup = [&] {
+        release_task.count_down();
+        scheduler.stop();
+    };
+    ASSERT_TRUE(scheduler
+                        .submit_scan_task(SimplifiedScanTask(
+                                [&] {
+                                    task_started.count_down();
+                                    release_task.wait();
+                                    return true;
+                                },
+                                nullptr, nullptr))
+                        .ok());
+    ASSERT_TRUE(task_started.wait_for(std::chrono::seconds(5)));
+    scanner_context->_scanner_scheduler = &scheduler;
+
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    ASSERT_FALSE(scanner_context->_pending_scanners.empty());
+
+    // Context submission is fail-fast. Retrying here would couple the scanner 
scheduler to
+    // ThreadPool's internal rejection/retention behavior.
+    Status surfaced = scheduler.schedule_scan_task(scanner_context, nullptr, 
transfer_lock);
+    EXPECT_TRUE(surfaced.is<ErrorCode::TOO_MANY_TASKS>()) << 
surfaced.to_string();
+    EXPECT_TRUE(scanner_context->done());
+    EXPECT_FALSE(scanner_context->_process_status.ok());
+    EXPECT_TRUE(scan_dependency->ready());
+    // The marker is set before submit_func(). This rejected runnable was not 
retained, but the
+    // terminal Context no longer needs the marker cleared or another 
submission attempted.
+    EXPECT_TRUE(scanner_context->is_context_queued(transfer_lock));
+}
+
+TEST_F(ScannerContextTest, 
thread_pool_budget_check_and_submit_are_atomic_across_contexts) {
+    // Three workers and no queue capacity. Two workers are parked, as if each 
ran a scanner of
+    // another Context, so the pool can accept exactly one more runnable.
+    ThreadPoolSimplifiedScanScheduler scheduler("budget_race_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(3, 3, 0, 3).ok());
+    CountDownLatch tasks_started(2);
+    CountDownLatch release_tasks(1);
+    Defer cleanup = [&] {
+        release_tasks.count_down();
+        scheduler.stop();
+    };
+    for (int i = 0; i < 2; ++i) {
+        ASSERT_TRUE(scheduler
+                            .submit_scan_task(SimplifiedScanTask(
+                                    [&] {
+                                        tasks_started.count_down();
+                                        release_tasks.wait();
+                                        return true;
+                                    },
+                                    nullptr, nullptr))
+                            .ok());
+    }
+    ASSERT_TRUE(tasks_started.wait_for(std::chrono::seconds(5)));
+    ASSERT_EQ(scheduler.get_active_threads(), 2);
+
+    const int parallel_tasks = 4;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    // Each Context has one scanner running and another pending, and has 
reached its minimum
+    // concurrency. It may only ramp up while the pool has slack.
+    std::vector<std::shared_ptr<ScannerContext>> contexts;
+    for (int i = 0; i < 2; ++i) {
+        std::list<std::shared_ptr<ScannerDelegate>> scanners;
+        for (int j = 0; j < 2; ++j) {
+            scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+        }
+        auto scanner_context = ScannerContext::create_shared(
+                state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+                scanners, -1, scan_dependency, parallel_tasks);
+        scanner_context->_scanner_scheduler = &scheduler;
+        scanner_context->_min_scan_concurrency_of_scan_scheduler = 3;
+        scanner_context->_max_scan_concurrency = parallel_tasks;
+        scanner_context->_min_scan_concurrency = 1;
+        scanner_context->_scan_starving = false;
+        scanner_context->_num_scheduled_scanners = 1;
+        ASSERT_FALSE(scanner_context->_pending_scanners.empty());
+        contexts.push_back(std::move(scanner_context));
+    }
+
+    // Park the first Context between its budget check and its submission 
until the second one
+    // reaches the same point, or for a bounded time if the scheduler 
serializes them. Without a
+    // scheduler-wide lock both pass the check on the last free slot and one 
submission fails.
+    std::atomic<int> entered {0};
+    CountDownLatch second_entered(1);
+    const bool old_enable_debug_points = config::enable_debug_points;
+    config::enable_debug_points = true;
+    DebugPoints::instance()->add_with_callback(
+            
"ThreadPoolSimplifiedScanScheduler.schedule_scan_task.before_submit",
+            std::function<void()>([&] {
+                if (entered.fetch_add(1) == 0) {
+                    
static_cast<void>(second_entered.wait_for(std::chrono::milliseconds(500)));
+                } else {
+                    second_entered.count_down();
+                }
+            }));
+    Defer cleanup_debug_point = [&] {
+        DebugPoints::instance()->remove(
+                
"ThreadPoolSimplifiedScanScheduler.schedule_scan_task.before_submit");
+        config::enable_debug_points = old_enable_debug_points;
+    };
+
+    Status statuses[2];
+    std::vector<std::thread> threads;
+    for (int i = 0; i < 2; ++i) {
+        threads.emplace_back([&, i] {
+            std::unique_lock<std::mutex> 
transfer_lock(contexts[i]->transfer_lock());
+            statuses[i] = scheduler.schedule_scan_task(contexts[i], nullptr, 
transfer_lock);
+        });
+    }
+    for (auto& thread : threads) {
+        thread.join();
+    }
+
+    // The Context that lost the last slot defers instead of failing its query.
+    for (int i = 0; i < 2; ++i) {
+        EXPECT_TRUE(statuses[i].ok()) << statuses[i].to_string();
+        EXPECT_FALSE(contexts[i]->done());
+        std::unique_lock<std::mutex> 
transfer_lock(contexts[i]->transfer_lock());
+        EXPECT_TRUE(contexts[i]->_process_status.ok());
+    }
+}
+
+TEST_F(ScannerContextTest, run_context_publishes_admission_failure) {
+    const bool old_enable_debug_points = config::enable_debug_points;
+    config::enable_debug_points = true;
+    
DebugPoints::instance()->add("ThreadPoolSimplifiedScanScheduler._run_context.inject_failure");
+    Defer cleanup_debug_point = [&] {
+        DebugPoints::instance()->remove(
+                
"ThreadPoolSimplifiedScanScheduler._run_context.inject_failure");
+        config::enable_debug_points = old_enable_debug_points;
+    };
+
+    const int parallel_tasks = 2;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < 2; ++i) {
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+    // The worker's task_exec_ctx() must resolve, otherwise _run_context() 
exits before admission.
+    // HasTaskExecutionCtx snapshots the weak_ptr at construction, so set it 
before create_shared.
+    auto task_execution_context = std::make_shared<TaskExecutionContext>();
+    state->set_task_execution_context(task_execution_context);
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+
+    ThreadPoolSimplifiedScanScheduler scheduler("run_context_failure_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(1, 1, 1, 1).ok());
+    Defer cleanup = [&] { scheduler.stop(); };
+    scanner_context->_scanner_scheduler = &scheduler;
+
+    {
+        std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+        ASSERT_TRUE(scheduler.schedule_scan_task(scanner_context, nullptr, 
transfer_lock).ok());
+        ASSERT_TRUE(scanner_context->is_context_queued(transfer_lock));
+    }
+
+    // The worker admits a scanner and hits the injected exception. It must 
publish the failure
+    // as a completed task instead of terminating the process or leaking the 
scheduled slot.
+    bool published = false;
+    for (int i = 0; i < 10000; ++i) {
+        std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+        if (!scanner_context->_tasks_queue.empty()) {
+            published = true;
+            break;
+        }
+        transfer_lock.unlock();
+        std::this_thread::sleep_for(std::chrono::milliseconds(1));
+    }
+    ASSERT_TRUE(published);
+
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    ASSERT_EQ(scanner_context->_tasks_queue.size(), 1);
+    EXPECT_FALSE(scanner_context->_tasks_queue.front()->status_ok());
+    EXPECT_FALSE(scanner_context->_process_status.ok());
+    EXPECT_EQ(scanner_context->_num_scheduled_scanners, 0);
+    EXPECT_FALSE(scanner_context->is_context_queued(transfer_lock));
+}
+
+TEST_F(ScannerContextTest, thread_pool_context_chain_runs_all_scanners) {
+    const int parallel_tasks = 2;
+    const int scanner_count = 3;
+    const int blocks_per_scanner = 4;
+    // Return after every block so each scanner needs several admissions 
before it reaches EOS.
+    const int32_t old_doris_scanner_row_bytes = 
config::doris_scanner_row_bytes;
+    config::doris_scanner_row_bytes = 1;
+    Defer restore_config = [&] { config::doris_scanner_row_bytes = 
old_doris_scanner_row_bytes; };
+
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+    olap_scan_local_state->_parent = scan_operator.get();
+    olap_scan_local_state->_max_scan_concurrency = 
max_concurrency_counter.get();
+    olap_scan_local_state->_min_scan_concurrency = 
min_concurrency_counter.get();
+    RuntimeProfile::HighWaterMarkCounter peak_running_scanner(TUnit::UNIT, 0, 
"");
+    olap_scan_local_state->_peak_running_scanner = &peak_running_scanner;
+    scan_operator->_should_run_serial = false;
+    TQueryOptions query_options;
+    query_options.__set_max_column_reader_num(0);
+    state->set_query_options(query_options);
+
+    std::atomic<int> running {0};
+    std::atomic<int> peak_running {0};
+    CountDownLatch overlap(1);
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < scanner_count; ++i) {
+        std::shared_ptr<Scanner> scanner = std::make_shared<ChainMockScanner>(
+                state.get(), olap_scan_local_state.get(), profile.get(), 
blocks_per_scanner,
+                &running, &peak_running, &overlap);
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+
+    // The worker's task_exec_ctx() must resolve, otherwise _run_context() 
exits before admission.
+    auto task_execution_context = std::make_shared<TaskExecutionContext>();
+    state->set_task_execution_context(task_execution_context);
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+    scanner_context->_newly_create_free_blocks_num = 
newly_create_free_blocks_num.get();
+    scanner_context->_scanner_memory_used_counter = 
scanner_memory_used_counter.get();
+
+    // Two workers so the successor runnable can overlap with the executing 
scanner.
+    ThreadPoolSimplifiedScanScheduler scheduler("context_chain_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(2, 2, 16, 1).ok());
+    Defer cleanup = [&] { scheduler.stop(); };
+    scanner_context->_scanner_scheduler = &scheduler;
+    // The two-thread pool never reaches this budget, so the Context may ramp 
to its maximum.
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 20;
+    ASSERT_EQ(scanner_context->_max_scan_concurrency, parallel_tasks);
+
+    // init() performs the bootstrap submission of the first Context runnable.
+    ASSERT_TRUE(scanner_context->init().ok());
+
+    int64_t rows = 0;
+    bool eos = false;
+    const auto deadline = std::chrono::steady_clock::now() + 
std::chrono::seconds(20);
+    while (!eos) {
+        ASSERT_LT(std::chrono::steady_clock::now(), deadline) << 
scanner_context->debug_string();
+        // One Context submission represents all pending scanners, so the pool 
never holds more
+        // than one runnable for this Context.
+        EXPECT_LE(scheduler.get_queue_size(), 1);
+        Block block;
+        Status st = scanner_context->get_block_from_queue(state.get(), &block, 
&eos, 0);
+        ASSERT_TRUE(st.ok()) << st.to_string();
+        rows += block.rows();
+        if (!eos && block.rows() == 0) {
+            std::this_thread::sleep_for(std::chrono::milliseconds(1));
+        }
+    }
+
+    // Every consumed non-EOS scanner was re-admitted until it reported EOS.
+    EXPECT_EQ(rows, scanner_count * blocks_per_scanner);
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    EXPECT_EQ(scanner_context->_num_finished_scanners, scanner_count);
+    EXPECT_EQ(scanner_context->_num_scheduled_scanners, 0);
+    EXPECT_TRUE(scanner_context->_pending_scanners.empty());
+    EXPECT_FALSE(scanner_context->is_context_queued(transfer_lock));
+    EXPECT_TRUE(scanner_context->_process_status.ok());
+    // The successor runnable ramped concurrency to the per-Context limit, and 
never beyond it.
+    EXPECT_EQ(peak_running.load(), parallel_tasks);
+}
+
+TEST_F(ScannerContextTest, successor_submission_is_not_scanner_wait_time) {
+    const int parallel_tasks = 2;
+    const int scanner_count = 2;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+    olap_scan_local_state->_parent = scan_operator.get();
+    olap_scan_local_state->_max_scan_concurrency = 
max_concurrency_counter.get();
+    olap_scan_local_state->_min_scan_concurrency = 
min_concurrency_counter.get();
+    RuntimeProfile::HighWaterMarkCounter peak_running_scanner(TUnit::UNIT, 0, 
"");
+    olap_scan_local_state->_peak_running_scanner = &peak_running_scanner;
+    scan_operator->_should_run_serial = false;
+    TQueryOptions query_options;
+    query_options.__set_max_column_reader_num(0);
+    state->set_query_options(query_options);
+
+    // The scanners never wait for each other: the latch is already open.
+    std::atomic<int> running {0};
+    std::atomic<int> peak_running {0};
+    CountDownLatch overlap(0);
+    std::vector<std::shared_ptr<Scanner>> scanner_ptrs;
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < scanner_count; ++i) {
+        std::shared_ptr<Scanner> scanner = std::make_shared<ChainMockScanner>(
+                state.get(), olap_scan_local_state.get(), profile.get(), 1, 
&running, &peak_running,
+                &overlap);
+        scanner_ptrs.push_back(scanner);
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+
+    // The worker's task_exec_ctx() must resolve, otherwise _run_context() 
exits before admission.
+    auto task_execution_context = std::make_shared<TaskExecutionContext>();
+    state->set_task_execution_context(task_execution_context);
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+    scanner_context->_newly_create_free_blocks_num = 
newly_create_free_blocks_num.get();
+    scanner_context->_scanner_memory_used_counter = 
scanner_memory_used_counter.get();
+
+    ThreadPoolSimplifiedScanScheduler scheduler("successor_wait_time_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(1, 1, 16, 1).ok());
+    Defer cleanup = [&] { scheduler.stop(); };
+    scanner_context->_scanner_scheduler = &scheduler;
+    // The pool never reaches this budget, so the runnable submits a successor 
for the second
+    // pending scanner after admitting the first one.
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 20;
+
+    // Make the successor submission from the pool worker slow, as when 
ThreadPool synchronously
+    // creates a thread. The bootstrap submission from this thread is not 
delayed.
+    const int64_t submit_delay_ms = 1000;
+    const auto test_thread_id = std::this_thread::get_id();
+    const bool old_enable_debug_points = config::enable_debug_points;
+    config::enable_debug_points = true;
+    DebugPoints::instance()->add_with_callback(
+            
"ThreadPoolSimplifiedScanScheduler.schedule_scan_task.before_submit",
+            std::function<void()>([&] {
+                if (std::this_thread::get_id() != test_thread_id) {
+                    
std::this_thread::sleep_for(std::chrono::milliseconds(submit_delay_ms));
+                }
+            }));
+    Defer cleanup_debug_point = [&] {
+        DebugPoints::instance()->remove(
+                
"ThreadPoolSimplifiedScanScheduler.schedule_scan_task.before_submit");
+        config::enable_debug_points = old_enable_debug_points;
+    };
+
+    // init() performs the bootstrap submission of the first Context runnable.
+    ASSERT_TRUE(scanner_context->init().ok());
+
+    bool published = false;
+    for (int i = 0; i < 20000; ++i) {
+        std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+        if (scanner_context->_tasks_queue.size() == scanner_count) {
+            published = true;
+            break;
+        }
+        transfer_lock.unlock();
+        std::this_thread::sleep_for(std::chrono::milliseconds(1));
+    }
+    ASSERT_TRUE(published) << scanner_context->debug_string();
+
+    // Both scanners ran on a worker right after admission. Neither may be 
charged the delayed
+    // successor submission as time spent waiting for a worker; charging it 
would add the whole
+    // delay on top of the runnable's own queue wait, which only covers the 
worker wake-up.
+    for (const auto& scanner : scanner_ptrs) {
+        EXPECT_LT(scanner->get_scanner_wait_worker_timer(), submit_delay_ms * 
1000 * 1000);
+    }
+}
+
+TEST_F(ScannerContextTest, thread_pool_context_runnable_is_deduplicated) {
+    const int parallel_tasks = 2;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < 3; ++i) {
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+    scanner_context->_newly_create_free_blocks_num = 
newly_create_free_blocks_num.get();
+    scanner_context->_scanner_memory_used_counter = 
scanner_memory_used_counter.get();
+
+    // One worker, parked, with queue capacity: a submitted Context runnable 
stays observable in
+    // the queue instead of being executed or rejected.
+    ThreadPoolSimplifiedScanScheduler scheduler("context_dedup_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(1, 1, 4, 1).ok());
+    CountDownLatch task_started(1);
+    CountDownLatch release_task(1);
+    Defer cleanup = [&] {
+        release_task.count_down();
+        scheduler.stop();
+    };
+    ASSERT_TRUE(scheduler
+                        .submit_scan_task(SimplifiedScanTask(
+                                [&] {
+                                    task_started.count_down();
+                                    release_task.wait();
+                                    return true;
+                                },
+                                nullptr, nullptr))
+                        .ok());
+    ASSERT_TRUE(task_started.wait_for(std::chrono::seconds(5)));
+    scanner_context->_scanner_scheduler = &scheduler;
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 20;
+
+    auto completed_task = std::make_shared<ScanTask>(scanners.front());
+    {
+        std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+        ASSERT_TRUE(scheduler.schedule_scan_task(scanner_context, nullptr, 
transfer_lock).ok());
+        EXPECT_TRUE(scanner_context->is_context_queued(transfer_lock));
+        EXPECT_EQ(scheduler.get_queue_size(), 1);
+
+        // A second scheduling attempt while a runnable is queued must not add 
another runnable.
+        ASSERT_TRUE(scheduler.schedule_scan_task(scanner_context, nullptr, 
transfer_lock).ok());
+        EXPECT_TRUE(scanner_context->is_context_queued(transfer_lock));
+        EXPECT_EQ(scheduler.get_queue_size(), 1);
+
+        // Publish a completed non-EOS result so the operator can consume it 
below.
+        completed_task->cached_blocks.emplace_back(Block::create_unique(), 0);
+        scanner_context->_tasks_queue.push_back(completed_task);
+        scanner_context->_num_scheduled_scanners = 1;
+    }
+
+    // Consuming a non-EOS result returns the scanner to the admission queue. 
The queued runnable
+    // will see it, so no additional runnable is submitted.
+    Block block;
+    bool eos = false;
+    Status st = scanner_context->get_block_from_queue(state.get(), &block, 
&eos, 0);
+    ASSERT_TRUE(st.ok()) << st.to_string();
+    EXPECT_FALSE(eos);
+    {
+        std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+        EXPECT_TRUE(completed_task->cached_blocks.empty());
+        EXPECT_TRUE(scanner_context->_tasks_queue.empty());
+        ASSERT_FALSE(scanner_context->_pending_scanners.empty());
+        EXPECT_EQ(scanner_context->_pending_scanners.top(), completed_task);
+        EXPECT_TRUE(scanner_context->is_context_queued(transfer_lock));
+        EXPECT_EQ(scheduler.get_queue_size(), 1);
+        // Cancel the query before the parked worker runs the queued runnable, 
so it exits without
+        // touching the OlapScanner that has no tablet behind it.
+        scanner_context->_should_stop = true;
+    }
+}
+
+TEST_F(ScannerContextTest, thread_pool_stopped_scheduler_fails_context) {
+    const int parallel_tasks = 2;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners {
+            std::make_shared<ScannerDelegate>(scanner)};
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+
+    ThreadPoolSimplifiedScanScheduler scheduler("stopped_scheduler_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(1, 1, 1, 1).ok());
+    scheduler.stop();
+    scanner_context->_scanner_scheduler = &scheduler;
+
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    ASSERT_FALSE(scanner_context->_pending_scanners.empty());
+    Status surfaced = scheduler.schedule_scan_task(scanner_context, nullptr, 
transfer_lock);
+    EXPECT_TRUE(surfaced.is<ErrorCode::INTERNAL_ERROR>()) << 
surfaced.to_string();
+    // The Context is terminal and the operator is woken to observe the 
failure. No runnable was
+    // submitted, so the marker stays clear.
+    EXPECT_TRUE(scanner_context->done());
+    EXPECT_FALSE(scanner_context->_process_status.ok());
+    EXPECT_TRUE(scan_dependency->ready());
+    EXPECT_FALSE(scanner_context->is_context_queued(transfer_lock));
+}
+
+TEST_F(ScannerContextTest, thread_pool_stopped_scheduler_fails_queued_context) 
{
+    const int parallel_tasks = 2;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners {
+            std::make_shared<ScannerDelegate>(scanner)};
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+
+    // A single worker parked with queue capacity keeps the submitted Context 
runnable queued.
+    ThreadPoolSimplifiedScanScheduler 
scheduler("stopped_queued_scheduler_test", cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(1, 1, 4, 1).ok());
+    CountDownLatch task_started(1);
+    CountDownLatch release_task(1);
+    Defer cleanup = [&] {
+        release_task.count_down();
+        // The test sets only the stop flag below; clear it so stop() really 
shuts the pool down.
+        scheduler._is_stop = false;
+        scheduler.stop();
+    };
+    ASSERT_TRUE(scheduler
+                        .submit_scan_task(SimplifiedScanTask(
+                                [&] {
+                                    task_started.count_down();
+                                    release_task.wait();
+                                    return true;
+                                },
+                                nullptr, nullptr))
+                        .ok());
+    ASSERT_TRUE(task_started.wait_for(std::chrono::seconds(5)));
+    scanner_context->_scanner_scheduler = &scheduler;
+
+    {
+        std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+        ASSERT_TRUE(scheduler.schedule_scan_task(scanner_context, nullptr, 
transfer_lock).ok());
+        ASSERT_TRUE(scanner_context->is_context_queued(transfer_lock));
+    }
+
+    // Stopping the pool drops the queued runnable, so the marker is never 
cleared by it. The next
+    // scheduling attempt must fail the Context instead of waiting for that 
runnable forever. Only
+    // the stop flag is set here; the pool itself is shut down by the cleanup 
above.
+    scheduler._is_stop = true;
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    Status surfaced = scheduler.schedule_scan_task(scanner_context, nullptr, 
transfer_lock);
+    EXPECT_TRUE(surfaced.is<ErrorCode::INTERNAL_ERROR>()) << 
surfaced.to_string();
+    EXPECT_TRUE(scanner_context->done());
+    EXPECT_FALSE(scanner_context->_process_status.ok());
+    EXPECT_TRUE(scan_dependency->ready());
+}
+
+TEST_F(ScannerContextTest, run_context_publishes_successor_submit_failure) {
+    const int parallel_tasks = 2;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < 2; ++i) {
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+    // The worker's task_exec_ctx() must resolve, otherwise _run_context() 
exits before admission.
+    auto task_execution_context = std::make_shared<TaskExecutionContext>();
+    state->set_task_execution_context(task_execution_context);
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, -1, scan_dependency, parallel_tasks);
+
+    // One worker and no queue capacity: the first Context runnable is 
accepted by the idle worker,
+    // but the successor it submits while running is rejected.
+    ThreadPoolSimplifiedScanScheduler scheduler("successor_failure_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(1, 1, 0, 1).ok());
+    Defer cleanup = [&] { scheduler.stop(); };
+    scanner_context->_scanner_scheduler = &scheduler;
+    // The pool never reaches this budget, so the second pending scanner is 
admissible and the
+    // runnable tries to submit its successor.
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 20;
+
+    {
+        std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+        ASSERT_TRUE(scheduler.schedule_scan_task(scanner_context, nullptr, 
transfer_lock).ok());
+    }
+
+    // The admitted scanner is published with the submission failure instead 
of being executed,
+    // so its scheduled slot is released and the operator observes the error.
+    bool published = false;
+    for (int i = 0; i < 10000; ++i) {
+        std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+        if (!scanner_context->_tasks_queue.empty()) {
+            published = true;
+            break;
+        }
+        transfer_lock.unlock();
+        std::this_thread::sleep_for(std::chrono::milliseconds(1));
+    }
+    ASSERT_TRUE(published);
+
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    ASSERT_EQ(scanner_context->_tasks_queue.size(), 1);
+    
EXPECT_TRUE(scanner_context->_tasks_queue.front()->get_status().is<ErrorCode::TOO_MANY_TASKS>())
+            << scanner_context->_tasks_queue.front()->get_status().to_string();
+    
EXPECT_TRUE(scanner_context->_process_status.is<ErrorCode::TOO_MANY_TASKS>());
+    EXPECT_TRUE(scanner_context->done());
+    EXPECT_EQ(scanner_context->_num_scheduled_scanners, 0);
+    // The second scanner was never admitted.
+    EXPECT_EQ(scanner_context->_pending_scanners.size(), 1);
+}
+
+TEST_F(ScannerContextTest, terminal_eos_skips_context_submission) {
+    // A full pool: one parked worker and no queue capacity, so any submission 
would fail.
+    ThreadPoolSimplifiedScanScheduler scheduler("terminal_eos_test", 
cgroup_cpu_ctl);
+    ASSERT_TRUE(scheduler.start(1, 1, 0, 1).ok());
+    CountDownLatch task_started(1);
+    CountDownLatch release_task(1);
+    Defer cleanup = [&] {
+        release_task.count_down();
+        scheduler.stop();
+    };
+    ASSERT_TRUE(scheduler
+                        .submit_scan_task(SimplifiedScanTask(
+                                [&] {
+                                    task_started.count_down();
+                                    release_task.wait();
+                                    return true;
+                                },
+                                nullptr, nullptr))
+                        .ok());
+    ASSERT_TRUE(task_started.wait_for(std::chrono::seconds(5)));
+    ASSERT_EQ(scheduler.get_active_threads(), 1);
+
+    const int parallel_tasks = 1;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = 100;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners {
+            std::make_shared<ScannerDelegate>(scanner)};
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
output_row_descriptor,
+            scanners, 100, scan_dependency, parallel_tasks);
+    scanner_context->_scanner_scheduler = &scheduler;
+
+    // The only scanner has reported EOS and its result waits for the operator.
+    scanner_context->_pending_scanners = 
std::stack<std::shared_ptr<ScanTask>>();
+    auto eos_task = std::make_shared<ScanTask>(scanners.front());
+    eos_task->set_eos(true);
+    scanner_context->_tasks_queue.push_back(eos_task);
+    scanner_context->_num_scheduled_scanners = 0;
+
+    MockRuntimeStateLocal mock_runtime_state;
+    EXPECT_CALL(mock_runtime_state, 
is_cancelled()).WillRepeatedly(testing::Return(false));
+    Block block;
+    bool eos = false;
+    Status status = scanner_context->get_block_from_queue(&mock_runtime_state, 
&block, &eos, 0);
+
+    // All scanners completed: the Context finishes without submitting a 
runnable, even though
+    // the pool could not accept one.
+    EXPECT_TRUE(status.ok()) << status.to_string();
+    EXPECT_TRUE(eos);
+    EXPECT_EQ(scanner_context->_num_finished_scanners, 1);
+    EXPECT_EQ(scheduler.get_queue_size(), 0);
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    EXPECT_FALSE(scanner_context->is_context_queued(transfer_lock));
+}
+
 } // namespace doris


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

Reply via email to