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


##########
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:
   Follow-up for the scheduling case from the review of 3abe4030fe5 (A holds 
the last row and waits in the pending stack, B returns an empty block, and the 
context runs B again): valid, fixed in e7600d25a9b and 06adaacea6e.
   
   **Context scheduling (e7600d25a9b).** `ScannerContext` now reads the 
buffered-row counter of the scan operator. While the buffered rows cover the 
remaining LIMIT, admission runs the most recently used pending scanner that 
holds padding rows, instead of the top of the pending stack. In the reported 
sequence the pending scanners are [A, B] after the operator consumed B's empty 
block, the context runs A, A emits the row it holds without reading, and the 
LIMIT is exhausted. The TaskExecutor path had the same order, because 
`_pull_next_scan_task()` resubmitted the consumed scanner first; both paths now 
share `_pop_pending_scan_task()`.
   
   In all other cases the order stays LIFO. This includes the case where no 
pending scanner holds rows (the holder is running, its block is not consumed 
yet, or a retired scanner left its rows counted), so a scanner without rows 
still progresses through its range. The pending scanners that hold rows are 
kept in a second vector, so nothing is searched when there is no such scanner.
   
   **Empty reads inside one `get_block()` call (06adaacea6e).** The loop in 
`Scanner::get_block()` reads on while the block is empty and is bounded by the 
rows the reader returns. A reader that filters rows itself returns empty blocks 
that are not counted, so B could still read its whole range in one call. The 
loop now also ends when the buffered rows cover the remaining LIMIT. It is a 
do/while loop, so every call still reads once.
   
   **Tests**
   - 
`ScannerProjectionTest.shared_limit_context_runs_scanner_holding_buffered_rows` 
is the Context-level test: two scanners run through `ScannerContext` and a real 
`ThreadPoolSimplifiedScanScheduler`, the operator side consumes with 
`get_block_from_queue()`. A emits four rows and holds one while B's first read 
is in flight, then the context is limited to one scanner. The test asserts that 
B reads one block of its filtered tail (9 of 10 blocks are left) and that the 
scan returns 5 rows. Without the fix B reads all 10 blocks and the test fails.
   - `ScannerContextTest.admission_prefers_scanner_holding_buffered_rows` 
covers both admission paths and the LIFO fallbacks.
   - The test scanner no longer counts a full read of rows for an empty block, 
so the existing `shared_limit_*` tests now cover the uncounted empty reads.
   - 4dcf5397fe9 also makes the two thread pool tests for failed thread starts 
deterministic: they force the interleaving with debug point callbacks and 
latches instead of sleeps.
   
   **Verification (local)**
   - BE UT: `ScannerProjectionTest.*`, `ScannerContextTest.*`, 
`ScannerLateArrivalRfTest.*`, `WorkloadGroupManagerTest.*`, `ThreadPoolTest.*` 
and the file/olap scanner tests pass (154 tests).
   - Regression on a local cluster: `query_p0/limit` and 
`correctness_p0/test_shared_scan_limit_pending_tasks` pass with the thread pool 
scheduler and with `enable_task_executor_in_*_table = true`.
   
   Known limit: if the scanner that holds the rows belongs to the context of 
another parallel instance, this context keeps its LIFO order until that context 
emits the rows.
   



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