morningman opened a new pull request, #68712:
URL: https://github.com/apache/doris/pull/68712

   ### What problem does this PR solve?
   
   Issue Number: None
   
   Related PR: #61271 (adaptive scan concurrency), #53854 (one instance per BE 
for simple queries)
   
   Problem Summary:
   
   **In short.** With `enable_adaptive_scan` on, the default, every scan has 
run at its minimum concurrency since #61271 - one scanner per instance by 
default - however much memory it was allowed. Simple queries, which the FE runs 
as one instance per BE, therefore read a BE's whole share of a table with one 
scanner, internal and external tables alike: `SELECT *` over a 16-bucket fluss 
table took 8.7 s instead of 3.4 s, over 71 paimon splits 11.5 s instead of 7.8 
s. This PR lets the adaptive concurrency rise to the ceiling the memory limiter 
gives it, as #61271 describes.
   
   **Background**
   
   - A scan instance runs its scanners through `ScannerContext`. With adaptive 
scan on, `_available_pickup_scanner_count()` sets `expected_scanners`, and both 
scheduler paths cap the scanners running at once by it: the TaskExecutor path 
in `_get_margin()` and `_pull_next_scan_task()`, the thread-pool path in 
`can_admit_scan_task()`.
   - The memory limiter (`MemLimiter::available_scanner_count()`) is what 
adapts: its ceiling shrinks when blocks are estimated larger or the query's 
scan memory is shared by more scan nodes, and grows back when they are not.
   - The minimum comes from `min_scanners_concurrency` / 
`min_file_scanners_concurrency`, 1 when unset; the maximum from 
`max_scanners_concurrency` / `max_file_scanners_concurrency` (4 for internal 
tables, 16 for file scans).
   - A simple query (a result sink straight over a scan) runs as one instance 
per BE (#53854), so that one instance's scanners are all the parallelism the 
scan gets on that BE.
   
   **The problem, and what it cost**
   
   `_available_pickup_scanner_count()` refreshed `expected_scanners` by 
clamping its previous value into `[max(1, minimum), memory ceiling]`. It starts 
at zero and nothing ever raised it (a TODO in the code marked the missing 
step), so every scan stayed at its minimum, one scanner by default.
   
   It costs most on scans without aggregation over a BE's data, which run as 
one instance: they read every split one after another. On one BE (Release), 
`SELECT *` as a dry run, scanners running at once (`RunningScannerPeak`) in 
brackets:
   
   - an internal table of 16 tablets: 1 scanner; 4 with 
`enable_adaptive_scan=false`;
   - a 10M-row fluss log table of 16 buckets: 8.7 s (1); 3.5 s with adaptive 
scan off;
   - a 30M-row paimon table of 71 parquet splits through the paimon catalog: 
11.5 s (1); 7.9 s with adaptive scan off.
   
   Scans with fewer ranges than instances are not hit (FE marks the scan serial 
and BE multiplies the minimum by the instance count), and aggregations run many 
instances, so the slowdown shows on large `SELECT *`-style reads and exports.
   
   **How this PR fixes it**
   
   `_available_pickup_scanner_count()` takes the memory limiter's ceiling 
(still capped at `_max_scan_concurrency`) as the expected concurrency instead 
of clamping the previous value into it. Nothing else changes:
   
   - the ceiling still wins over the minimum, and still follows the memory 
budget down and back up;
   - a scheduler without slack still holds a context at its minimum in 
`_get_margin()` and `can_admit_scan_task()`;
   - with adaptive scan off, nothing changes.
   
   That is the concurrency scans had before #61271, now bounded by the memory 
limiter, and the behaviour #61271's description lays out: 
`_available_pickup_scanner_count()` takes its count from 
`ScannerMemLimiter::available_scanner_count(ins_idx)`, the scan memory limit 
over the estimated block size, divided among the instances. Raising the 
concurrency step by step instead would need thresholds and would not reach 
short queries of a few hundred milliseconds.
   
   **Results**
   
   Same BE (Release, macOS), same data, before and after; `SELECT *` as a dry 
run, `RunningScannerPeak` in brackets:
   
   | Table | before | after | adaptive scan off (before / after) |
   |---|---|---|---|
   | internal table, 16 tablets | 1 scanner | 4 scanners | 4 |
   | fluss log table, 10M rows, 16 buckets | 8.7 s (1) | 3.4 s (16) | 3.5 s / 
3.5 s |
   | paimon table, 30M rows, 71 splits | 11.5 s (1) | 7.8 s (16) | 7.9 s / 7.6 
s |
   
   `COUNT(*)`, `SUM`, `GROUP BY` and `LIMIT 10` on the same tables are not 
slower (the fluss log table's `COUNT(*)`: 261 -> 140 ms).
   
   The trade-off is the one before #61271: a scan without aggregation may now 
run up to 16 file scanners (4 for internal tables) per instance at once, as far 
as memory and the scheduler allow. For JNI readers that keep much in BE's JVM 
heap (paimon merge reads of uncompacted primary-key tables, fluss primary-key 
bucket reads with large change logs) that is enough to run the default 2 GB 
heap out, as `enable_adaptive_scan=false` already does today; see the merge 
order below.
   
   **Classes, and how they call each other**
   
   - `ScannerContext::_available_pickup_scanner_count()` (changed): 
`expected_scanners = min(memory ceiling, _max_scan_concurrency)`, refreshed as 
before.
   - `ScannerContext::_get_margin()`, `_pull_next_scan_task()` (TaskExecutor 
path) and `can_admit_scan_task()` (thread-pool path) (untouched): cap the 
running scanners by `expected_scanners`, and fall back to the minimum when the 
scheduler has no slack.
   - `MemLimiter::available_scanner_count()` (untouched): the ceiling.
   
   ```
   ScanOperator instance --> ScannerContext
       _available_pickup_scanner_count()
           ceiling = min(MemLimiter.available_scanner_count(ins_idx), 
_max_scan_concurrency)
           before: expected = clamp(previous expected, max(1, min), ceiling)   
-> starts at 0, stays at min
           after:  expected = ceiling
       TaskExecutor path:  _get_margin() / _pull_next_scan_task()  <= expected  
 (min when no slack)
       thread-pool path:   can_admit_scan_task()                   <= expected  
 (min when no slack)
   ```
   
   **Related open PRs.** #68374 raises the default of 
`min_scanners_concurrency` from 1 to 4 on the FE side, which lifts the floor 
scans were stuck at; this PR removes the reason they were stuck on it. #68610 
and #66838 also touch `scanner_context.cpp` and its test; the overlap is in the 
test file only.
   
   **Merge order.** No code dependency. Please merge it after #PR1, which keeps 
fluss and paimon threads that die of an `OutOfMemoryError` from exiting BE: 
with up to 16 JNI readers per instance again, running BE's JVM heap out becomes 
more likely by default, and without #PR1 that crashes BE. #PR2 (fluss 
connection pooling) is best merged before it too, and #PR4 adds an opt-in 
admission for heavy JNI readers and an error that says how to give BE's JVM 
more heap.
   
   ### Release note
   
   Fix scans running with a single scanner per instance when 
`enable_adaptive_scan` is on (the default), which made queries without 
aggregation read each backend's data serially.
   
   ### Check List (For Author)
   
   - Test <!-- At least one of them must be included. -->
       - [ ] Regression test
       - [x] Unit Test
       - [x] Manual test (add detailed scripts or steps below)
       - [ ] No need to test or manual test. Explain why:
           - [ ] This is a refactor/code format and no logic has been changed.
           - [ ] Previous test can cover this change.
           - [ ] No code files have been changed.
           - [ ] Other reason <!-- Add your reason?  -->
   
     **Unit tests.** `ScannerContextTest` (two new): on both scheduler paths a 
scan with memory to spare is admitted up to `_max_scan_concurrency`, follows 
the ceiling down when the budget shrinks and back up when it grows. All 43 
tests pass (ASAN, macOS, clang 22).
   
     **Manual.** The before/after runs above, on one Release BE with the same 
data.
   
   - Behavior changed:
       - [ ] No.
       - [x] Yes. <!-- Explain the behavior change --> With 
`enable_adaptive_scan` on, a scan may run up to `max_scanners_concurrency` / 
`max_file_scanners_concurrency` scanners per instance again, bounded by the 
memory limiter and by the scheduler's slack.
   
   - Does this need documentation?
       - [x] No.
       - [ ] Yes. <!-- Add document PR link here. eg: 
https://github.com/apache/doris-website/pull/1214 -->
   
   ### Check List (For Reviewer who merge this PR)
   
   - [ ] Confirm the release note
   - [ ] Confirm test cases
   - [ ] Confirm document
   - [ ] Add branch pick label <!-- Add branch pick label that this PR should 
merge into -->
   


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