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]