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]