mrhhsg commented on code in PR #68610:
URL: https://github.com/apache/doris/pull/68610#discussion_r4227061468


##########
be/src/exec/scan/scanner.cpp:
##########
@@ -92,7 +92,10 @@ Status Scanner::get_block_after_projects(RuntimeState* 
state, Block* block, bool
     const auto& row_descriptor = 
_local_state->_parent->operator_row_desc_before_projection();
     if (_has_projection) {
         _origin_block.clear_column_data(row_descriptor.num_slots());
-        if (!_can_merge_padding_blocks(_padding_block, _origin_block)) {
+        // get_block() charges the shared LIMIT as soon as rows pass the 
filters. Once peer scanners
+        // exhaust it, the context may finish without running this scanner 
again, so rows held in
+        // _padding_block would be charged but never returned. Do not pad 
under a shared LIMIT.
+        if (_shared_scan_limit || !_can_merge_padding_blocks(_padding_block, 
_origin_block)) {
             DORIS_CHECK(_padding_block.empty())

Review Comment:
   Valid, fixed in 846a19ba33c.
   
   `Scanner::get_block_after_projects()` pads again under a shared LIMIT. 
Instead of skipping padding, it keeps the shared counter equal to "LIMIT minus 
returned rows" between calls: rows still held in `_padding_block` when a call 
returns are refunded to the counter, and they are charged again when the next 
call starts. During a call every row the scanner holds is charged 
(`get_block()` still charges as rows pass the filters), so the read loop stops 
as soon as the LIMIT is reached and small LIMITs keep the same read bound as 
before. An exhausted counter therefore always means that at least LIMIT rows 
were returned, which is what the admission guard and the EOS check in 
`ScannerContext` rely on, while a selective projected scan under a large LIMIT 
still coalesces source blocks up to half a batch per call.
   
   Tests:
   - `ScannerProjectionTest.shared_limit_keeps_padding_small_blocks`: LIMIT 
100, six single-row source blocks, batch size 8; one call returns four rows as 
one block and the counter reads 96. It fails on the previous head, which 
returned one row per call.
   - `ScannerProjectionTest.shared_limit_refunds_rows_held_in_padding`: a held 
row is not charged after the call returns (counter 5 = 10 - 5 returned); when 
peers bring the counter to 1, the next call charges the row, sees the counter 
exhausted and flushes it with eos instead of reading on. It fails on master, 
which keeps the held row charged.
   - 
`ScannerProjectionTest.shared_limit_reports_eos_when_scanner_rows_exhaust_it` 
is unchanged.
   - Local regression: `query_p0/limit` suites and 
`correctness_p0/test_shared_scan_limit_pending_tasks` pass with the thread pool 
scheduler as default.
   



##########
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:
   Valid, fixed in 95f931a121e.
   
   `multiget_data_v2` now calls the new 
`RowIdStorageReader::read_by_rowids_for_rpc()`, which runs `read_by_rowids()` 
under `ASSIGN_STATUS_IF_CATCH_EXCEPTION` and serializes the resulting status 
into the response. A `doris::Exception` thrown while decoding the request (for 
example `SlotDescriptor(pslot)` on an unknown `PTypeNode.type`) becomes an 
error status on the response, which is written before `ClosureGuard` completes 
the RPC, so no exception leaves the ThreadPool worker. This mirrors the 
exception conversion the external-row path already has in 
`submit_external_scan_tasks()`.
   
   Test: 
`RowIdStorageReaderTest.RpcHelperReportsThrownExceptionOnThreadPoolPath` 
registers a file mapping for the query, sends a slot whose type node kind is 
unknown to this BE, and runs the helper on a real 
`ThreadPoolSimplifiedScanScheduler`. It asserts that the response carries a 
non-zero status code whose message names the invalid type node. Without the fix 
the exception escapes the worker and the test process aborts. 
`ExecEnv::set_id_manager()` was added so the test can install the `IdManager` 
the reader looks up.
   



##########
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:
   Not changed in this PR, for the record on why.
   
   This thread describes `ThreadPool::do_submit()` as it exists on master, in 
code this PR does not touch: the task is queued before the first worker start 
is attempted, and a failed start with no running worker returns an error while 
the task stays queued. The ThreadPool rework that addressed it (withdrawing the 
task, waiting for pending starts, keeping the last worker) was removed from 
this PR on purpose to keep it to the default switch and the two shared-LIMIT 
fixes, and is kept on a separate branch for its own PR, where the generic 
ThreadPool change can be reviewed with its own tests. Every ThreadPool user is 
exposed to it in the same way; this PR only makes the scan schedulers use the 
pool that was the default before the task executor existed. Reaching the path 
also requires a pool with zero running workers whose `pthread_create` fails, 
while the scan pools start `min_threads` workers at init.
   



##########
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:
   Not changed in this PR, for the record on why.
   
   This thread describes `ThreadPool::do_submit()` as it exists on master, in 
code this PR does not touch: the task is queued before the first worker start 
is attempted, and a failed start with no running worker returns an error while 
the task stays queued. The ThreadPool rework that addressed it (withdrawing the 
task, waiting for pending starts, keeping the last worker) was removed from 
this PR on purpose to keep it to the default switch and the two shared-LIMIT 
fixes, and is kept on a separate branch for its own PR, where the generic 
ThreadPool change can be reviewed with its own tests. Every ThreadPool user is 
exposed to it in the same way; this PR only makes the scan schedulers use the 
pool that was the default before the task executor existed. Reaching the path 
also requires a pool with zero running workers whose `pthread_create` fails, 
while the scan pools start `min_threads` workers at init.
   



-- 
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