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]

Reply via email to