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


##########
be/src/exec/scan/scanner.cpp:
##########
@@ -101,6 +101,16 @@ Status Scanner::get_block_after_projects(RuntimeState* 
state, Block* block, bool
             RETURN_IF_ERROR(get_block(state, &_origin_block, eos));
             return _do_projections(&_origin_block, block);
         }
+        // get_block() charges the shared LIMIT as soon as rows pass the 
filters, but rows that stay
+        // in _padding_block when this call returns are not returned by it. 
Peer scanners may
+        // exhaust the LIMIT before this scanner runs again, and the context 
then finishes without
+        // running it, so those rows would be charged but never returned. Keep 
the counter equal
+        // to "LIMIT minus returned rows" between calls: charge the held rows 
again on entry, so
+        // every row this scanner holds during the call is charged and the 
read loop below stops
+        // as soon as the LIMIT is reached, and refund whatever is still held 
on exit.
+        if (_shared_scan_limit && !_padding_block.empty()) {
+            _shared_scan_limit->fetch_sub(_padding_block.rows(), 
std::memory_order_acq_rel);
+        }
         const auto min_batch_size = std::max(state->batch_size() / 2, 1);

Review Comment:
   Valid, fixed in 64b2915c764.
   
   `ScanLocalState::shared_scan_limit_ptr()` now decides from the operator's 
SQL LIMIT (`limit() < 0`) whether to return the counter, instead of from the 
counter's current value. That is also how 
`ScannerContext::_is_shared_scan_limit_exhausted()` decides it. A scanner whose 
local state opens while the counter is briefly negative therefore still gets 
the pointer: its `get_block()` reports eos while the counter is at or below 
zero, and once peers restore it the scanner charges its rows and stops when the 
LIMIT is reached, as before. Without a SQL LIMIT the pointer stays null, so the 
no-LIMIT and key TopN paths are unchanged.
   
   `MockScanOperatorX::set_limit_for_test()` sets the LIMIT and the counter the 
way the operator constructor does; the existing shared-limit tests use it.
   
   Tests:
   - 
`ScannerProjectionTest.shared_limit_observed_by_scanner_initialized_under_negative_counter`:
 LIMIT 8, the counter is -1 while a peer still holds 2 charged padding rows, 
the scanner is initialized at that moment, the peer refunds its 2 rows (counter 
1), and the scanner charges its first block, reports eos and leaves its second 
block unread. On the previous head the scanner gets a null pointer and the 
first assertion fails.
   - `ScannerProjectionTest.no_shared_limit_without_sql_limit`: without a SQL 
LIMIT the pointer stays null whatever the counter holds.
   - `ScannerProjectionTest.*`, `ScannerContextTest.*`, 
`ScannerLateArrivalRfTest.*`, `WorkloadGroupManagerTest.*` and 
`RowIdStorageReaderTest.*` pass (109 tests). `query_p0/limit` (19 suites) and 
`correctness_p0/test_shared_scan_limit_pending_tasks` pass on a local cluster 
with the thread pool scheduler as default. The branch is rebased onto current 
master.
   



##########
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 on the pending-holder / empty-reader case that the latest review 
body still carries: not changed in this PR, for the record on why.
   
   That case is read amplification only. While a pending scanner holds refunded 
padding rows, an active scanner may read on through a filtered tail before the 
LIMIT is satisfied, but the returned rows stay complete: between calls the 
counter equals LIMIT minus returned rows, so an exhausted counter always means 
at least LIMIT rows were returned, and the operator cuts the excess. This PR 
deliberately does not cap what the scan reads at the LIMIT. The buffered-row 
counter and the pending-queue reordering that addressed this case in earlier 
heads were removed on purpose, to keep the PR to the default switch and the 
correctness fixes and to avoid cross-scanner bookkeeping on every block for a 
bound on extra reads. The extra read is the same kind a non-projected LIMIT 
scan does when one scanner is slower than its peers.
   



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