mrhhsg commented on code in PR #68571:
URL: https://github.com/apache/doris/pull/68571#discussion_r4122075922
##########
be/src/exec/scan/scanner_context.cpp:
##########
@@ -448,17 +453,129 @@ 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(); restart the
per-scanner wait
+ // timer here so it measures admission-to-execution instead of everything
since the previous
+ // attempt paused, which would include time the cached blocks waited for
the operator. The
+ // Context runnable waited for a worker on behalf of this scanner 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()) {
+ scanner_delegate->_scanner->start_wait_worker_timer();
Review Comment:
Fixed in 3ebd9f87a74. `try_get_next_scan_task()` no longer starts the
per-scanner stopwatch; it only credits the Context runnable's queue wait
(clipped to the scanner's own pending interval, unchanged). `_run_context()`
now calls `start_wait_worker_timer()` right before `execute_scan_task()`, after
the successor runnable has been submitted and `transfer_lock` released, so a
slow successor `submit_func()` (for example synchronous thread creation in the
remote pool) is no longer charged to `ScannerWorkerWaitTime` /
`PerScannerWaitTime` of a scanner that is already on a worker. The admission
credit and the stopwatch cover disjoint intervals, and the TaskExecutor path
(timer started in `ScannerScheduler::submit`) is unchanged.
New UT `ScannerContextTest.successor_submission_is_not_scanner_wait_time`:
one pool worker, two scanners; a debug point delays only the successor
submission issued from the pool worker by 1000 ms, and the test asserts neither
scanner is charged that delay. On d783450 it fails (the first scanner is
charged the full delay); with the fix `ScannerContextTest.*` passes (25 tests).
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]