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


##########
be/src/exec/scan/scanner.cpp:
##########
@@ -88,6 +89,38 @@ 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) {
+        _shared_scan_limit->fetch_sub(block->rows(), 
std::memory_order_acq_rel);
+        *eos = *eos || _shared_scan_limit->load(std::memory_order_acquire) <= 
0;
+    }
+    // After the charge above, so peers never see the emitted rows missing 
from both counters.
+    _publish_padding_rows();
+    return Status::OK();
+}
+
+void Scanner::_publish_padding_rows() {
+    if (!_shared_scan_limit) {
+        return;
+    }
+    const auto padding_rows = cast_set<int64_t>(_padding_block.rows());
+    _shared_scan_buffered_rows->fetch_add(padding_rows - 
_published_padding_rows,

Review Comment:
   Valid, fixed in 03980e13e00.
   
   `Scanner::_publish_padding_rows()` now returns before the atomic operation 
when the number of padding rows equals the number this scanner has already 
published. A scanner without projection never holds padding rows, so ordinary 
LIMIT scans no longer issue a read-modify-write on the shared counter for every 
block. The counter is still updated whenever the padding rows actually change.
   
   The values seen by peers are unchanged: this function is the only writer of 
the counter and of `_published_padding_rows`, and `fetch_add(0)` never changed 
the value. No reader relies on the removed RMW for synchronization. The context 
reads a scanner's `buffered_rows()` under `_transfer_lock` after its task 
completes.
   
   Test: new 
`ScannerProjectionTest.shared_limit_scanner_without_projection_publishes_no_padding_rows`
 covers the non-projected path (only the emitted rows are charged; the 
buffered-row counter published by a peer stays untouched). 
`ScannerProjectionTest.*` / `ScannerContextTest.*` pass, and the 
`query_p0/limit` suites 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