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]