mrhhsg commented on code in PR #68610:
URL: https://github.com/apache/doris/pull/68610#discussion_r4132689880
##########
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 scheduling case from the review of 3abe4030fe5 (A holds
the last row and waits in the pending stack, B returns an empty block, and the
context runs B again): valid, fixed in e7600d25a9b and 06adaacea6e.
**Context scheduling (e7600d25a9b).** `ScannerContext` now reads the
buffered-row counter of the scan operator. While the buffered rows cover the
remaining LIMIT, admission runs the most recently used pending scanner that
holds padding rows, instead of the top of the pending stack. In the reported
sequence the pending scanners are [A, B] after the operator consumed B's empty
block, the context runs A, A emits the row it holds without reading, and the
LIMIT is exhausted. The TaskExecutor path had the same order, because
`_pull_next_scan_task()` resubmitted the consumed scanner first; both paths now
share `_pop_pending_scan_task()`.
In all other cases the order stays LIFO. This includes the case where no
pending scanner holds rows (the holder is running, its block is not consumed
yet, or a retired scanner left its rows counted), so a scanner without rows
still progresses through its range. The pending scanners that hold rows are
kept in a second vector, so nothing is searched when there is no such scanner.
**Empty reads inside one `get_block()` call (06adaacea6e).** The loop in
`Scanner::get_block()` reads on while the block is empty and is bounded by the
rows the reader returns. A reader that filters rows itself returns empty blocks
that are not counted, so B could still read its whole range in one call. The
loop now also ends when the buffered rows cover the remaining LIMIT. It is a
do/while loop, so every call still reads once.
**Tests**
-
`ScannerProjectionTest.shared_limit_context_runs_scanner_holding_buffered_rows`
is the Context-level test: two scanners run through `ScannerContext` and a real
`ThreadPoolSimplifiedScanScheduler`, the operator side consumes with
`get_block_from_queue()`. A emits four rows and holds one while B's first read
is in flight, then the context is limited to one scanner. The test asserts that
B reads one block of its filtered tail (9 of 10 blocks are left) and that the
scan returns 5 rows. Without the fix B reads all 10 blocks and the test fails.
- `ScannerContextTest.admission_prefers_scanner_holding_buffered_rows`
covers both admission paths and the LIFO fallbacks.
- The test scanner no longer counts a full read of rows for an empty block,
so the existing `shared_limit_*` tests now cover the uncounted empty reads.
- 4dcf5397fe9 also makes the two thread pool tests for failed thread starts
deterministic: they force the interleaving with debug point callbacks and
latches instead of sleeps.
**Verification (local)**
- BE UT: `ScannerProjectionTest.*`, `ScannerContextTest.*`,
`ScannerLateArrivalRfTest.*`, `WorkloadGroupManagerTest.*`, `ThreadPoolTest.*`
and the file/olap scanner tests pass (154 tests).
- Regression on a local cluster: `query_p0/limit` and
`correctness_p0/test_shared_scan_limit_pending_tasks` pass with the thread pool
scheduler and with `enable_task_executor_in_*_table = true`.
Known limit: if the scanner that holds the rows belongs to the context of
another parallel instance, this context keeps its LIFO order until that context
emits the rows.
--
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]