mrhhsg commented on code in PR #68610:
URL: https://github.com/apache/doris/pull/68610#discussion_r4130224014
##########
be/src/common/config.cpp:
##########
@@ -335,9 +335,9 @@ DEFINE_Int32(task_executor_max_concurrency_per_task, "-1");
DEFINE_Int32(task_executor_initial_max_concurrency_per_task, "-1");
// Enable task executor in internal table scan.
-DEFINE_Bool(enable_task_executor_in_internal_table, "true");
+DEFINE_Bool(enable_task_executor_in_internal_table, "false");
// Enable task executor in external table scan.
-DEFINE_Bool(enable_task_executor_in_external_table, "true");
+DEFINE_Bool(enable_task_executor_in_external_table, "false");
Review Comment:
Fixed in bfea36ff571.
`ThreadPool::do_submit()` now withdraws its own queued task when the pool
has no worker and creating the first worker fails. It restores the token and
`_queue` bookkeeping, moves the token back to IDLE when it has no entries left,
and only then returns the error. A failed submit therefore never runs later, so
`multiget_data_v2` completing `done` itself on the error is safe. The fix is in
the shared `ThreadPool`, so every caller that relies on "failed submit ==
rejected" benefits, including the Context runnable of
`ThreadPoolSimplifiedScanScheduler`.
- A `Task::taken_by_worker` marker is set under the pool lock when a worker
takes the task. This tells "a worker started by a concurrent submit already ran
it" (the submit returns OK) apart from "`shutdown()` released it without
running it" (the error is still returned).
- The same failure branch now notifies `_no_threads_cond`. Previously a
`shutdown()` waiting for the last pending thread could hang forever when that
thread failed to start. The branch also logs a WARNING.
- New UT `ThreadPoolTest.TestFailedSubmitWithoutWorkersNeverRuns`. It
injects the thread-creation failure through the new debug point
`ThreadPool.create_thread.inject_failure`, for both the pool and a SERIAL
token. It checks that nothing stays queued, that the token goes back to idle,
and that the failed tasks never run after later successful submits. The old
code fails all of these checks.
##########
be/src/exec/scan/scanner_context.cpp:
##########
@@ -603,6 +603,14 @@ bool ScannerContext::can_admit_scan_task(const
std::unique_lock<std::mutex>& tra
if (done() || _pending_tasks.empty()) {
return false;
}
+ // Same rule as _pull_next_scan_task() on the TaskExecutor path: once the
shared LIMIT is
+ // exhausted, pending scanners would only open and immediately report EOS,
so do not admit
+ // them while a completed or in-flight task can still wake the operator.
If neither exists,
+ // admit one so it can report EOS and wake the pipeline task.
+ if (_is_shared_scan_limit_exhausted() &&
Review Comment:
Fixed in 30617c27748, following the "account for LIMIT on emitted rows"
option.
The shared scan LIMIT counter is no longer charged inside
`Scanner::get_block()`. It is charged in `Scanner::get_block_after_projects()`,
for the rows that actually leave the scanner, and that call reports EOS once
the counter is exhausted. Rows still buffered in `_padding_block` are no longer
counted. So an exhausted counter now always means at least LIMIT rows were
already emitted, and neither the new ThreadPool admission guard nor the
existing `get_block_from_queue()` EOS check can drop rows that were charged. In
your example, A's 10 buffered rows are not charged, so after B's 1586 rows the
counter is 10, not 0, and A is admitted again to flush them.
The trade-off is declared in the commit message. Rows a scanner is still
padding (fewer than half a batch) are not visible to peer scanners until they
are emitted, so a scanner may read up to that many extra rows. Results are
unaffected.
Tests:
- `ScannerProjectionTest.shared_limit_charges_only_emitted_rows` models this
scenario at the scanner level: a buffered row, a peer that consumes part of the
limit, then the buffered row is flushed. It fails on the old accounting.
-
`ScannerProjectionTest.shared_limit_reports_eos_when_emitted_rows_exhaust_it`
covers the new EOS path.
- The local regression suites
`correctness_p0/test_shared_scan_limit_pending_tasks` and the `query_p0/limit`
suites pass with the thread pool scheduler as default.
--
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]