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]