This is an automated email from the ASF dual-hosted git repository.

Gabriel39 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 8a1ec48b7c4 [fix](be) Let adaptive scan concurrency rise to its memory 
ceiling (#68712)
8a1ec48b7c4 is described below

commit 8a1ec48b7c46fea546c35b3f67258586decb6dba
Author: Mingyu Chen (Rayner) <[email protected]>
AuthorDate: Sun Oct 4 21:55:27 2026 +0800

    [fix](be) Let adaptive scan concurrency rise to its memory ceiling (#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 master, one
    BE (Release), `SELECT *` as a dry run:
    
    | Table | adaptive scan on (default) | adaptive scan off |
    |---|---|---|
    | internal table, 16 tablets | 1 scanner at a time | 4 scanners |
    | fluss log table, 10M rows, 16 buckets | 8.7 s, 1 scanner | 3.5 s, 16
    scanners |
    | paimon table through the paimon catalog, 30M rows, 71 parquet splits |
    11.5 s, 1 scanner | 7.9 s, 16 scanners |
    
    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. On the TaskExecutor
    path, `_pull_next_scan_task()` now also takes a ceiling of zero as zero:
    it read `expected_scanners == 0` as `_max_scan_concurrency`, which made
    no difference while an instance never ran more than one scanner, but
    would let an instance whose allocation falls to zero keep every scanner
    it holds. Like `can_admit_scan_task()` on the thread-pool path, it still
    admits one task when nothing is occupied. 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 (both builds also
    carry the fluss connection changes of #68711); the scanners running at
    once (`RunningScannerPeak`) in brackets.
    
    Scans without aggregation (dry run):
    
    | Query | before | after | adaptive scan off, before / after |
    |---|---|---|---|
    | internal table, 16 tablets, `SELECT *` | 1 scanner | 4 scanners | 4
    scanners |
    | fluss log table (10M rows, 16 buckets), `SELECT *` | 8742 ms (1) |
    3411 ms (16) | 3496 / 3459 ms |
    | the same, `SELECT url, payload` | 4454 ms | 2455 ms | 2404 / 2437 ms |
    | paimon table (30M rows, 71 splits), `SELECT *` | 11460 ms (1) | 7821
    ms (16) | 7887 / 7634 ms |
    | the same, `SELECT url, payload` | 8509 ms | 6148 ms | 6214 / 6173 ms |
    
    Aggregations and `LIMIT` are not slower:
    
    | Query | fluss log table, before / after | paimon table, before / after
    |
    |---|---|---|
    | `COUNT(*)` | 261 / 140 ms | 158 / 147 ms |
    | `SUM(qty), SUM(price)` | 324 / 269 ms | 286 / 267 ms |
    | `GROUP BY category` with `COUNT`, `SUM`, `AVG` | 445 / 369 ms | 384 /
    354 ms |
    | `SELECT * ... LIMIT 10` | 19 / 18 ms | 133 / 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.
    
    The per-query file cache limit (`file_cache_query_limit_percent`) feels
    the same change. It is best-effort: once a query reaches it,
    `BlockFileCache::try_reserve()` evicts only the query's own blocks that
    no reader holds and admits the new block anyway when the rest are still
    being read, so the more scanners read at once, the further a limited
    query can go over its limit. `test_file_cache_query_limit` asserts the
    limit as a hard bound: at up to 16 scanners per instance its limited
    query ended with 54 blocks (56.6 MB) against a 21.5 MB limit on External
    Regression, while runs with one scanner per instance end with 20 blocks
    (20.2 MB). The suite now runs its scans with
    `max_file_scanners_concurrency=1`, as the condition cache suites next to
    it do.
    
    **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::_pull_next_scan_task()` (TaskExecutor path)
    (changed): caps the running scanners by `expected_scanners`, zero
    included; one task still runs when nothing is occupied.
    - `ScannerContext::_get_margin()` (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)
                            _pull_next_scan_task(): expected 0 caps at 0 (was 
_max_scan_concurrency),
                            one task when none is occupied
        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.
    #68610 merges with this PR without conflicts. #66838 already conflicts
    with master in `scanner_context.h` and the test; `scanner_context.cpp`
    merges cleanly, though both PRs edit `_pull_next_scan_task()` (it adds
    an `only_existing_split` argument, this PR takes a zero allocation as
    zero).
    
    **Merge order.** No code dependency. Please merge it after #68710, 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 #68710 that
    crashes BE. #68711 (fluss connection pooling) is best merged before it
    too, and #68713 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. -->
        - [x] 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` (three new): on both scheduler
    paths a scan with memory to spare is admitted up to
    `_max_scan_concurrency`; on the thread-pool path it follows the ceiling
    down when the budget shrinks and back up when it grows; on the
    TaskExecutor path an instance running two scanners whose allocation
    falls to zero puts the scanner whose block was consumed back to pending
    instead of resubmitting it, and still runs one scanner once nothing is
    occupied (this test fails without the change to
    `_pull_next_scan_task()`). All 44 tests pass (ASAN, macOS, clang 22).
    
    **Regression test.** `test_file_cache_query_limit` runs its scans with
    `max_file_scanners_concurrency=1`, for the reason given under Results.
    It needs the hive docker environment and a single backend, so it runs on
    External Regression.
    
    **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 -->
    
    ---------
    
    Co-authored-by: Claude Opus 5.5 (1M context) <[email protected]>
---
 be/src/exec/scan/scanner_context.cpp               |  30 ++--
 be/test/exec/scan/scanner_context_test.cpp         | 200 +++++++++++++++++++++
 .../cache/test_file_cache_query_limit.groovy       |   5 +
 3 files changed, 222 insertions(+), 13 deletions(-)

diff --git a/be/src/exec/scan/scanner_context.cpp 
b/be/src/exec/scan/scanner_context.cpp
index c4180cfcc2a..944cef09d8d 100644
--- a/be/src/exec/scan/scanner_context.cpp
+++ b/be/src/exec/scan/scanner_context.cpp
@@ -146,10 +146,8 @@ int ScannerContext::_available_pickup_scanner_count() {
         return _max_scan_concurrency;
     }
 
-    int min_scanners = std::max(1, _min_scan_concurrency);
     int max_scanners = _scanner_mem_limiter->available_scanner_count(_ins_idx);
     max_scanners = std::min(max_scanners, _max_scan_concurrency);
-    min_scanners = std::min(min_scanners, max_scanners);
     if (_ins_idx == 0) {
         // Adjust memory limit via memory share arbitrator
         
_adjust_scan_mem_limit(_scanner_mem_limiter->get_arb_scanner_mem_bytes(),
@@ -166,14 +164,17 @@ int ScannerContext::_available_pickup_scanner_count() {
     P.adjust_scanners_last_timestamp = now;
     auto old_scanners = P.expected_scanners;
 
-    scanners = std::max(min_scanners, scanners);
-    scanners = std::min(max_scanners, scanners);
+    // The memory limiter is what adapts: its ceiling shrinks when blocks are 
estimated larger or
+    // the query's scan budget is shared by more scan nodes, and grows back 
when they are not. Take
+    // the whole ceiling. Clamping the previous value into [minimum, ceiling] 
instead never raises
+    // it above the minimum, since expected_scanners starts at zero, so every 
scan would keep a
+    // single scanner however much memory there is. The ceiling still wins 
over the minimum, and a
+    // scheduler without slack still holds the Context at its minimum in 
_get_margin() and
+    // can_admit_scan_task().
+    scanners = max_scanners;
     VLOG_DEBUG << fmt::format(
-            "_available_pickup_scanner_count. context = {}, old_scanners = {}, 
scanners = {} "
-            ", min_scanners: {}, max_scanners: {}",
-            debug_string(), old_scanners, scanners, min_scanners, 
max_scanners);
-
-    // TODO(gabriel): Scanners are scheduled adaptively based on the memory 
usage now.
+            "_available_pickup_scanner_count. context = {}, old_scanners = {}, 
scanners = {}",
+            debug_string(), old_scanners, scanners);
     return scanners;
 }
 
@@ -838,12 +839,15 @@ std::shared_ptr<ScanTask> 
ScannerContext::_pull_next_scan_task(
         std::shared_ptr<ScanTask> current_scan_task, int32_t 
current_concurrency) {
     int32_t effective_max_concurrency = _max_scan_concurrency;
     if (_enable_adaptive_scanners) {
-        effective_max_concurrency = _adaptive_processor->expected_scanners > 0
-                                            ? 
_adaptive_processor->expected_scanners
-                                            : _max_scan_concurrency;
+        // _get_margin() has just refreshed expected_scanners, so zero is a 
real allocation here, as
+        // in can_admit_scan_task(), not a missing one. Reading it as 
_max_scan_concurrency let an
+        // instance whose allocation fell to zero keep every scanner it held.
+        effective_max_concurrency = _adaptive_processor->expected_scanners;
     }
 
-    if (current_concurrency >= effective_max_concurrency) {
+    // Keep one task progressing even when the adaptive limit is zero. 
Otherwise no worker can
+    // publish a result and wake the operator to make another scheduling 
decision.
+    if (current_concurrency > 0 && current_concurrency >= 
effective_max_concurrency) {
         VLOG_DEBUG << fmt::format(
                 "ScannerContext {} current concurrency {} >= 
effective_max_concurrency {}, skip "
                 "pull",
diff --git a/be/test/exec/scan/scanner_context_test.cpp 
b/be/test/exec/scan/scanner_context_test.cpp
index ec1d0092f00..a6b0e975ce3 100644
--- a/be/test/exec/scan/scanner_context_test.cpp
+++ b/be/test/exec/scan/scanner_context_test.cpp
@@ -828,6 +828,206 @@ TEST_F(ScannerContextTest, 
thread_pool_admission_keeps_zero_adaptive_allocation)
     EXPECT_EQ(scanner_context->try_get_next_scan_task(transfer_lock), nullptr);
 }
 
+TEST_F(ScannerContextTest, thread_pool_admission_follows_memory_ceiling) {
+    const int parallel_tasks = 4;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < 6; ++i) {
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+
+    TUniqueId query_id = state->get_query_ctx()->query_id();
+    const int64_t query_mem_limit = 1024LL * 1024 * 1024;
+    auto arbitrator = MemShareArbitrator::create_shared(query_id, 
query_mem_limit, 0.3);
+    auto limiter = MemLimiter::create_shared(query_id, parallel_tasks, false,
+                                             
static_cast<int64_t>(query_mem_limit * 0.3));
+    // 1GB budget with 1MB estimated blocks: 1024 scanners across four 
instances, far more than the
+    // four this Context may run. ins_idx = 1 keeps 
_available_pickup_scanner_count() away from the
+    // arbitrator-driven limit adjustment, which would overwrite the budgets 
set below.
+    limiter->update_open_tasks_count(1);
+    limiter->update_mem_limit(1024LL * 1024 * 1024);
+    limiter->reestimated_block_mem_bytes(1024LL * 1024);
+
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
false, scanners, -1,
+            scan_dependency, &shared_limit, arbitrator, limiter, 1, true, 
parallel_tasks);
+    std::unique_ptr<MockSimplifiedScanScheduler> scheduler =
+            std::make_unique<MockSimplifiedScanScheduler>(cgroup_cpu_ctl);
+    EXPECT_CALL(*scheduler, 
get_active_threads()).WillRepeatedly(testing::Return(0));
+    EXPECT_CALL(*scheduler, 
get_queue_size()).WillRepeatedly(testing::Return(0));
+    scanner_context->_scanner_scheduler = scheduler.get();
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 20;
+    scanner_context->_min_scan_concurrency = 1;
+
+    std::unique_lock<std::mutex> 
transfer_lock(scanner_context->transfer_lock());
+    ASSERT_EQ(scanner_context->_max_scan_concurrency, parallel_tasks);
+
+    // Memory allows more than the minimum, so admission runs up to 
_max_scan_concurrency instead
+    // of stopping at the single scanner the minimum asks for.
+    for (int i = 0; i < parallel_tasks; ++i) {
+        ASSERT_NE(scanner_context->try_get_next_scan_task(transfer_lock), 
nullptr) << i;
+    }
+    EXPECT_EQ(scanner_context->_adaptive_processor->expected_scanners, 
parallel_tasks);
+    EXPECT_FALSE(scanner_context->can_admit_scan_task(transfer_lock));
+
+    // A smaller budget lowers the ceiling at the next refresh: 8MB / 1MB = 8 
scanners, two per
+    // instance. With two of the four still in flight, nothing more is 
admitted.
+    limiter->update_mem_limit(8LL * 1024 * 1024);
+    scanner_context->_adaptive_processor->adjust_scanners_last_timestamp = 0;
+    scanner_context->_in_flight_tasks_num = 2;
+    EXPECT_EQ(scanner_context->try_get_next_scan_task(transfer_lock), nullptr);
+    EXPECT_EQ(scanner_context->_adaptive_processor->expected_scanners, 2);
+
+    // The ceiling grows back with the budget.
+    limiter->update_mem_limit(1024LL * 1024 * 1024);
+    scanner_context->_adaptive_processor->adjust_scanners_last_timestamp = 0;
+    EXPECT_NE(scanner_context->try_get_next_scan_task(transfer_lock), nullptr);
+    EXPECT_EQ(scanner_context->_adaptive_processor->expected_scanners, 
parallel_tasks);
+}
+
+TEST_F(ScannerContextTest, adaptive_margin_takes_memory_ceiling) {
+    const int parallel_tasks = 4;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < 6; ++i) {
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+
+    TUniqueId query_id = state->get_query_ctx()->query_id();
+    const int64_t query_mem_limit = 1024LL * 1024 * 1024;
+    auto arbitrator = MemShareArbitrator::create_shared(query_id, 
query_mem_limit, 0.3);
+    auto limiter = MemLimiter::create_shared(query_id, parallel_tasks, false,
+                                             
static_cast<int64_t>(query_mem_limit * 0.3));
+    // Same budget as above: memory allows far more scanners than 
_max_scan_concurrency.
+    limiter->update_open_tasks_count(1);
+    limiter->update_mem_limit(1024LL * 1024 * 1024);
+    limiter->reestimated_block_mem_bytes(1024LL * 1024);
+
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
false, scanners, -1,
+            scan_dependency, &shared_limit, arbitrator, limiter, 1, true, 
parallel_tasks);
+    std::unique_ptr<MockSimplifiedScanScheduler> scheduler =
+            std::make_unique<MockSimplifiedScanScheduler>(cgroup_cpu_ctl);
+    EXPECT_CALL(*scheduler, 
get_active_threads()).WillRepeatedly(testing::Return(0));
+    EXPECT_CALL(*scheduler, 
get_queue_size()).WillRepeatedly(testing::Return(0));
+    scanner_context->_scanner_scheduler = scheduler.get();
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 20;
+    scanner_context->_min_scan_concurrency = 1;
+
+    std::mutex transfer_mutex;
+    std::unique_lock<std::mutex> transfer_lock(transfer_mutex);
+    std::shared_mutex scheduler_mutex;
+    std::unique_lock<std::shared_mutex> scheduler_lock(scheduler_mutex);
+
+    // The scheduler has slack and memory allows four scanners, so the 
TaskExecutor path submits
+    // the whole ceiling at once rather than the single scanner the minimum 
asks for.
+    EXPECT_EQ(scanner_context->_get_margin(transfer_lock, scheduler_lock), 
parallel_tasks);
+    EXPECT_EQ(scanner_context->_adaptive_processor->expected_scanners, 
parallel_tasks);
+    // _pull_next_scan_task() stops at the same ceiling.
+    EXPECT_NE(scanner_context->_pull_next_scan_task(nullptr, parallel_tasks - 
1), nullptr);
+    EXPECT_EQ(scanner_context->_pull_next_scan_task(nullptr, parallel_tasks), 
nullptr);
+}
+
+TEST_F(ScannerContextTest, task_executor_keeps_zero_adaptive_allocation) {
+    const int parallel_tasks = 2;
+    auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
+                                                             parallel_tasks, 
TQueryCacheParam {});
+    auto olap_scan_local_state =
+            OlapScanLocalState::create_unique(state.get(), 
scan_operator.get());
+
+    OlapScanner::Params scanner_params;
+    scanner_params.state = state.get();
+    scanner_params.profile = profile.get();
+    scanner_params.limit = -1;
+    scanner_params.key_ranges = std::vector<OlapScanRange*>();
+    std::shared_ptr<Scanner> scanner =
+            OlapScanner::create_shared(olap_scan_local_state.get(), 
std::move(scanner_params));
+
+    std::list<std::shared_ptr<ScannerDelegate>> scanners;
+    for (int i = 0; i < 5; ++i) {
+        scanners.push_back(std::make_shared<ScannerDelegate>(scanner));
+    }
+
+    TUniqueId query_id = state->get_query_ctx()->query_id();
+    const int64_t query_mem_limit = 1024LL * 1024 * 1024;
+    auto arbitrator = MemShareArbitrator::create_shared(query_id, 
query_mem_limit, 0.3);
+    auto limiter = MemLimiter::create_shared(query_id, parallel_tasks, false,
+                                             
static_cast<int64_t>(query_mem_limit * 0.3));
+    // 1GB budget with 1MB estimated blocks: instance 1 may run both scanners 
this Context allows.
+    // ins_idx = 1 keeps _available_pickup_scanner_count() away from the 
arbitrator-driven limit
+    // adjustment, which would overwrite the budgets set below.
+    limiter->update_open_tasks_count(1);
+    limiter->update_mem_limit(1024LL * 1024 * 1024);
+    limiter->reestimated_block_mem_bytes(1024LL * 1024);
+
+    auto scanner_context = ScannerContext::create_shared(
+            state.get(), olap_scan_local_state.get(), output_tuple_desc, 
false, scanners, -1,
+            scan_dependency, &shared_limit, arbitrator, limiter, 1, true, 
parallel_tasks);
+    std::unique_ptr<MockSimplifiedScanScheduler> scheduler =
+            std::make_unique<MockSimplifiedScanScheduler>(cgroup_cpu_ctl);
+    EXPECT_CALL(*scheduler, 
get_active_threads()).WillRepeatedly(testing::Return(0));
+    EXPECT_CALL(*scheduler, 
get_queue_size()).WillRepeatedly(testing::Return(0));
+    scanner_context->_scanner_scheduler = scheduler.get();
+    scanner_context->_min_scan_concurrency_of_scan_scheduler = 20;
+    scanner_context->_min_scan_concurrency = 1;
+
+    std::mutex transfer_mutex;
+    std::unique_lock<std::mutex> transfer_lock(transfer_mutex);
+    std::shared_mutex scheduler_mutex;
+    std::unique_lock<std::shared_mutex> scheduler_lock(scheduler_mutex);
+
+    // Memory allows both scanners, and the TaskExecutor path starts both.
+    ASSERT_TRUE(scanner_context->schedule_scan_task(nullptr, transfer_lock, 
scheduler_lock).ok());
+    EXPECT_EQ(scanner_context->_adaptive_processor->expected_scanners, 
parallel_tasks);
+    EXPECT_EQ(scanner_context->_in_flight_tasks_num, parallel_tasks);
+
+    // Larger blocks leave room for one scanner in the whole node: it goes to 
instance 0, and this
+    // instance is allocated none.
+    limiter->update_mem_limit(1024LL * 1024);
+    scanner_context->_adaptive_processor->adjust_scanners_last_timestamp = 0;
+
+    // The operator consumes one scanner's block while the other scanner is 
still in flight. The
+    // margin still asks for a task, but zero is the ceiling: the consumed 
scanner goes back to
+    // pending instead of being resubmitted.
+    scanner_context->_in_flight_tasks_num = 1;
+    const size_t pending_tasks = scanner_context->_pending_tasks.size();
+    auto consumed_task = std::make_shared<ScanTask>(scanners.back());
+    ASSERT_TRUE(
+            scanner_context->schedule_scan_task(consumed_task, transfer_lock, 
scheduler_lock).ok());
+    EXPECT_EQ(scanner_context->_adaptive_processor->expected_scanners, 0);
+    EXPECT_EQ(scanner_context->_in_flight_tasks_num, 1);
+    EXPECT_EQ(scanner_context->_pending_tasks.size(), pending_tasks + 1);
+
+    // With nothing occupied, one scanner still runs so the scan keeps moving.
+    scanner_context->_in_flight_tasks_num = 0;
+    ASSERT_TRUE(scanner_context->schedule_scan_task(nullptr, transfer_lock, 
scheduler_lock).ok());
+    EXPECT_EQ(scanner_context->_in_flight_tasks_num, 1);
+}
+
 TEST_F(ScannerContextTest, 
thread_pool_admission_holds_minimum_when_pool_saturated) {
     const int parallel_tasks = 4;
     auto scan_operator = std::make_unique<OlapScanOperatorX>(obj_pool.get(), 
tnode, 0, *descs,
diff --git 
a/regression-test/suites/external_table_p0/cache/test_file_cache_query_limit.groovy
 
b/regression-test/suites/external_table_p0/cache/test_file_cache_query_limit.groovy
index e498101a547..72fc9950546 100644
--- 
a/regression-test/suites/external_table_p0/cache/test_file_cache_query_limit.groovy
+++ 
b/regression-test/suites/external_table_p0/cache/test_file_cache_query_limit.groovy
@@ -50,6 +50,11 @@ suite("test_file_cache_query_limit", 
"p0,external,nonConcurrent") {
 
     sql """set enable_file_cache=true"""
     sql """set disable_file_cache=false"""
+    // The per-query file cache limit is best-effort: once a query reaches it, 
BlockFileCache::try_reserve()
+    // evicts only the query's own blocks that no reader holds, and admits the 
new block anyway when the rest
+    // are still being read. The bound asserted below holds only while few 
blocks are held at once, so the
+    // scans run with one scanner per instance; with up to 16 the limited 
query caches well over its limit.
+    sql """set max_file_scanners_concurrency=1"""
 
     // Note: This test case assumes a single backend scenario. Testing with 
single backend is logically equivalent
     // to testing with multiple backends having identical configurations, but 
simpler in logic.


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to