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


##########
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 remaining case from the review of d29c2133dbb (a scanner 
that holds no padding rows while a peer holds the rows the LIMIT still needs): 
valid, fixed in 3abe4030fe5.
   
   The padding loop in `Scanner::_get_block_after_projects()` now also checks 
the buffered-row counter after each read that was merged into the padding 
block, whether or not the scanner holds rows itself. In the reported sequence 
(projected, filtered `LIMIT 5`, batch size 8, A holds 1 row and emitted 4), B 
reads once, sees that the buffered rows (1) cover the remaining LIMIT (1), and 
returns an empty block without eos. Its scan task completes, A is scheduled 
again and emits the row it holds without reading, the counter reaches 0, and B 
reports eos at its next `get_block()`.
   
   The check is placed after the read on purpose. A scanner always reads once 
per call, so it still progresses to the end of its range when a retired peer 
left its rows in the counter. An empty block without eos is not a new shape: 
the non-projected path already returns it when all rows of one `get_block()` 
call are filtered.
   
   Tests:
   - 
`ScannerProjectionTest.shared_limit_stops_padding_when_peer_buffered_rows_cover_it`
 models the reported sequence and asserts that B reads one block of its 
filtered tail instead of all of them. It fails on d29c2133dbb.
   - 
`ScannerProjectionTest.shared_limit_scanner_progresses_when_buffered_rows_are_stale`
 asserts that a scanner under a stale counter still emits its rows and reaches 
eos.
   - `ScannerProjectionTest.*`, `ScannerContextTest.*` and 
`ScannerLateArrivalRfTest.*` pass locally (56 tests). `query_p0/limit` and 
`correctness_p0/test_shared_scan_limit_pending_tasks` pass on a local cluster.
   



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