mrhhsg commented on code in PR #68610:
URL: https://github.com/apache/doris/pull/68610#discussion_r4130925739
##########
be/src/util/threadpool.cpp:
##########
@@ -521,8 +528,14 @@ Status ThreadPool::do_submit(std::shared_ptr<Runnable> r,
ThreadPoolToken* token
l.lock();
_num_threads_pending_start--;
if (_num_threads + _num_threads_pending_start == 0) {
- // If we have no threads, we can't do any work.
- return status;
+ // shutdown() may be waiting for the last pending thread to go
away.
+ _no_threads_cond.notify_all();
+ // If we have no threads, we can't do any work. Callers treat
a failed submit as
+ // a task that will never run (for example, an RPC handler
completes its closure
+ // itself), so withdraw the task queued above before returning
the error.
+ LOG(WARNING) << "Thread pool " << _name
Review Comment:
Valid, fixed in 25c0e58885a and d29c2133dbb.
- `do_submit()` now waits for the outcome of pending thread starts while the
pool has no running worker (`_thread_start_cond`, notified when a worker
starts, when a creation fails, and on shutdown). A task is therefore only
queued when a worker is already running or when the submit starts a thread
itself, so B no longer queues behind A's unconfirmed start. This also covers a
second submit that cannot start its own thread because `max_threads` is reached.
- A submit whose own creation fails now withdraws its task whenever no
worker is running (`_num_threads == 0`), instead of returning OK because
another start is still pending.
- The pre-push review found one more way to strand an accepted task: after
`set_max_threads()` lowers the limit, the only running worker exited with
queued tasks while a start was pending. `dispatch_thread()` now keeps the only
running worker while tasks are queued.
Tests: `ThreadPoolTest.TestOverlappingFailedThreadStartsLeaveNoTask` (both
`max_threads = 2` and `max_threads = 1`) and
`ThreadPoolTest.TestOnlyWorkerKeepsQueuedTasksWhenMaxThreadsShrinks`; the whole
`ThreadPoolTest` suite passes locally.
##########
be/src/exec/scan/scanner.cpp:
##########
@@ -88,6 +88,20 @@ Status Scanner::init(RuntimeState* state, const
VExprContextSPtrs& conjuncts) {
}
Status Scanner::get_block_after_projects(RuntimeState* state, Block* block,
bool* eos) {
+ RETURN_IF_ERROR(_get_block_after_projects(state, block, eos));
+ // Publish progress to the shared counter so peer scanners can observe it.
Only rows that leave
+ // the scanner are charged: rows still held in _padding_block are not
charged, because once
+ // the counter is exhausted the context may finish without running this
scanner again. The
+ // counter may go negative when several scanners subtract concurrently;
that is harmless
+ // because the operator's reached_limit() makes the final cut.
+ if (_shared_scan_limit && block->rows() > 0) {
Review Comment:
Valid, fixed in 4887216d1d8.
The scan operator now also counts the rows that scanners hold in their
padding blocks (`_shared_scan_buffered_rows`, next to `_shared_scan_limit`). A
scanner publishes its padding rows after each merge and stops padding as soon
as it holds rows and the rows held by all scanners cover the remaining LIMIT.
The held rows are then emitted and charged as before, so rows that may still be
dropped are never charged, and the extra read is bounded by the `get_block()`
call in progress (`doris_scanner_row_num`).
In the reported case both scanners hold one row, the buffered rows (2) cover
the remaining LIMIT (2), and both emit without reading the filtered tail. A
scanner retired while holding padding rows leaves them in the counter; that
only makes peers emit smaller blocks, because EOS and admission still depend on
`_shared_scan_limit` alone.
Test:
`ScannerProjectionTest.shared_limit_stops_padding_when_buffered_rows_cover_it`
(one matching row followed by a long filtered tail; asserts the tail is not
read). The limit regression suites pass on a local cluster.
--
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]