github-actions[bot] commented on code in PR #67846:
URL: https://github.com/apache/doris/pull/67846#discussion_r4226614891


##########
be/test/runtime/workload_group/workload_group_manager_test.cpp:
##########
@@ -1360,4 +1362,250 @@ TEST_F(WorkloadGroupManagerTest, 
AdaptiveFlushRegistrationSurvivesIdChangeAndReu
     controller->adjust_once();
 }
 
+// Helpers for the process memory exceeded cases. The workload group memory 
limit is derived from
+// the process memory limit, so the process memory limit is set before the 
workload group is
+// created, and only the soft limit is lowered afterwards to simulate process 
memory pressure.
+static constexpr int64_t kProcessMemLimitForWg = 1024L * 1024 * 1000;
+static constexpr int64_t kLargeProcessMemLimit = 
std::numeric_limits<int64_t>::max() / 2;
+
+// When the process memory is exceeded and the paused query has no revocable 
memory, the query
+// should be kept paused instead of being cancelled immediately, so that it 
can be resumed once
+// other queries release memory.
+TEST_F(WorkloadGroupManagerTest, 
process_mem_exceeded_keeps_paused_and_resumes) {
+    const int64_t original_mem_limit = MemInfo::mem_limit();
+    const int64_t original_soft_mem_limit = MemInfo::soft_mem_limit();
+    Defer restore_mem_limit {[&]() {
+        MemInfo::set_mem_limit_for_test(original_mem_limit);
+        MemInfo::set_soft_mem_limit_for_test(original_soft_mem_limit);
+    }};
+    MemInfo::set_mem_limit_for_test(kProcessMemLimitForWg);
+    WorkloadGroupInfo wg_info {.id = 1,
+                               .memory_limit = kProcessMemLimitForWg,
+                               .min_memory_percent = 10,
+                               .max_memory_percent = 100};
+    auto wg = _wg_manager->get_or_create_workload_group(wg_info);
+
+    auto query = _generate_on_query(wg);
+    // Let the workload group use more than its min memory limit, so that the 
paused query is
+    // handled by handle_single_query_ directly instead of waiting for other 
workload groups.
+    query->query_mem_tracker()->consume(1024L * 1024 * 128);
+    Defer release_memory {[&]() { query->query_mem_tracker()->consume(-1024L * 
1024 * 128); }};
+    wg->refresh_memory_usage();
+    ASSERT_EQ(wg->min_memory_limit(), 1024L * 1024 * 100);
+    ASSERT_GT(wg->total_mem_used(), wg->min_memory_limit());
+
+    // Process soft memory limit is exceeded, hard memory limit is not.
+    MemInfo::set_mem_limit_for_test(kLargeProcessMemLimit);
+    MemInfo::set_soft_mem_limit_for_test(1);
+    _wg_manager->add_paused_query(query->resource_ctx(), 1024L,
+                                  
Status::Error(ErrorCode::PROCESS_MEMORY_EXCEEDED, "test"));
+
+    config::spill_in_paused_queue_timeout_ms = 60 * 1000;
+    for (int i = 0; i < 3; ++i) {
+        _wg_manager->handle_paused_queries();
+        ASSERT_FALSE(query->is_cancelled());
+        std::unique_lock<std::mutex> lock(_wg_manager->_paused_queries_lock);
+        ASSERT_EQ(_wg_manager->_paused_queries_list[wg].size(), 1);
+    }
+    ASSERT_TRUE(query->resource_ctx()
+                        ->task_controller()
+                        ->paused_reason()
+                        .is<ErrorCode::PROCESS_MEMORY_EXCEEDED>());
+
+    // Process memory pressure is relieved, the query should be resumed.
+    MemInfo::set_soft_mem_limit_for_test(original_soft_mem_limit);

Review Comment:
   [P2] Control available memory before asserting resume. Restoring 
`original_soft_mem_limit` here (and in the below-min resume test) does not make 
`is_exceed_soft_mem_limit(32 MiB)` false when the runner's system-available 
memory is below the warning watermark plus 32 MiB. `MemInfo::init()` reads that 
value from the host/cgroup, and the fixture never controls it. The manager then 
correctly keeps the query paused, so both new resume assertions fail on a 
loaded runner; the same input also changes the other tests' expected hard-limit 
flag. Mock or establish both arbitrator preconditions before asserting wakeup.



##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -585,9 +585,47 @@ void WorkloadGroupMgr::handle_paused_queries() {
                     ++query_it;
                     continue;
                 }
+
+                // The process memory pressure may have been relieved by cache 
reclamation or by
+                // other queries that finished. Check it before routing the 
query below, otherwise
+                // a query in a workload group that uses less than its min 
memory limit has to
+                // wait for the timeout.
+                const size_t test_memory_size =
+                        std::max<size_t>(query_it->reserve_size_, 32L * 1024 * 
1024);
+                if 
(!GlobalMemoryArbitrator::is_exceed_soft_mem_limit(test_memory_size)) {
+                    LOG(INFO) << "Query: " << 
print_id(resource_ctx->task_controller()->task_id())
+                              << ", process limit not exceeded now, resume 
this query"
+                              << ", process memory info: "
+                              << 
GlobalMemoryArbitrator::process_memory_used_details_str()
+                              << ", wg info: " << wg->debug_string();
+                    
resource_ctx->task_controller()->set_memory_sufficient(true);
+                    query_it = queries_list.erase(query_it);
+                    continue;
+                }
+
                 // If workload group's memory usage > min memory, then it 
means the workload group use too much memory
                 // in memory contention state. Should just spill
-                if (wg->total_mem_used() > wg->min_memory_limit()) {
+                bool handle_query_now = wg->total_mem_used() > 
wg->min_memory_limit();
+                if (!handle_query_now) {
+                    // Other workload groups many use a lot of memory, should 
revoke memory from other workload groups
+                    // by cancelling their queries.
+                    int64_t revoked_size = revoke_memory_from_other_groups_();
+                    if (revoked_size > 0) {
+                        // Revoke memory from other workload groups will 
cancel some queries, wait them cancel finished
+                        // and then check it again.
+                        revoking_memory_from_other_query_ = true;
+                        return;
+                    }
+
+                    // TODO revoke from memtable
+
+                    // Fallback: the query has waited too long and memory 
cannot be revoked from
+                    // anywhere, let handle_single_query_ spill or cancel it 
to protect the system.
+                    handle_query_now =

Review Comment:
   [P1] Check the hard process limit before waiting on the below-min timeout. 
When this WG is at or below its minimum and no other group can be revoked, 
`handle_query_now` stays false until `spill_in_paused_queue_timeout_ms` 
elapses, so the new `is_exceed_hard_mem_limit()` check inside 
`handle_single_query_()` never runs. A paused query holding memory can 
therefore remain blocked for the default 60 seconds after the process crosses 
its hard limit when memory GC is disabled. Include the hard-limit condition in 
this dispatch gate and test a below-min WG at the hard limit.



##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -796,38 +822,45 @@ bool WorkloadGroupMgr::handle_single_query_(const 
std::shared_ptr<ResourceContex
             return true;
         }
     } else {
-        // Should not consider about process memory. For example, the query's 
limit is 100g, workload
-        // group's memlimit is 10g, process memory is 20g. The query reserve 
will always failed in wg
-        // limit, and process is always have memory, so that it will resume 
and failed reserve again.
-        const size_t test_memory_size = std::max<size_t>(size_to_reserve, 32L 
* 1024 * 1024);
-        if 
(!GlobalMemoryArbitrator::is_exceed_soft_mem_limit(test_memory_size)) {
-            LOG(INFO) << "Query: " << query_id
-                      << ", process limit not exceeded now, resume this query"
-                      << ", process memory info: "
-                      << 
GlobalMemoryArbitrator::process_memory_used_details_str()
-                      << ", wg info: " << wg->debug_string();
-            requestor->task_controller()->set_memory_sufficient(true);
-            return true;
-        } else {
-            // if cannot find any memory to release, then let the query 
continue to run as far as possible
-            // or cancelled by gc if memory is really not enough.
-            Status error_status = Status::MemoryLimitExceeded(
-                    "Query {} process memory is exceeded"
-                    ", and there is no cache now. And could not find task to 
spill, disable "
-                    "reserve memory and resume it. "
-                    "Query memory usage: {}, limit: {}, reserved "
-                    "size: {}, try to reserve: {}, wg info: {}."
-                    " Maybe you should set the workload group's limit to a 
lower value. {}",
-                    query_id, PrettyPrinter::print_bytes(memory_usage),
-                    PrettyPrinter::print_bytes(limit), 
PrettyPrinter::print_bytes(reserved_size),
-                    PrettyPrinter::print_bytes(size_to_reserve), 
wg->memory_debug_string(),
-                    doris::ProcessProfile::instance()
-                            ->memory_profile()
-                            ->process_memory_detail_str());
-            LOG_LONG_STRING(INFO, error_status.to_string());
-            requestor->task_controller()->cancel(error_status);
-            return true;
+        // PROCESS_MEMORY_EXCEEDED. The caller (handle_paused_queries) has 
already resumed
+        // the query if the process is no longer above the soft memory limit, 
so the process
+        // memory is still exceeded here.
+        const bool exceed_hard_mem_limit = 
GlobalMemoryArbitrator::is_exceed_hard_mem_limit();

Review Comment:
   [P1] Apply hard-limit cancellation before the running-task return. In a 
multi-task query, one task can pause on process reservation while a sibling 
remains in an operator call. `get_revocable_info()` then marks 
`has_running_task`, so `handle_single_query_()` returns at line 740 before this 
new hard-limit check on every pass, even after the timeout. With memory GC 
disabled, the query can keep holding memory at the hard process limit 
indefinitely. Keep the running-task guard for spill, but cancel a 
process-paused query at the hard limit before that return.



##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -585,9 +585,47 @@ void WorkloadGroupMgr::handle_paused_queries() {
                     ++query_it;
                     continue;
                 }
+
+                // The process memory pressure may have been relieved by cache 
reclamation or by
+                // other queries that finished. Check it before routing the 
query below, otherwise
+                // a query in a workload group that uses less than its min 
memory limit has to
+                // wait for the timeout.
+                const size_t test_memory_size =
+                        std::max<size_t>(query_it->reserve_size_, 32L * 1024 * 
1024);
+                if 
(!GlobalMemoryArbitrator::is_exceed_soft_mem_limit(test_memory_size)) {
+                    LOG(INFO) << "Query: " << 
print_id(resource_ctx->task_controller()->task_id())
+                              << ", process limit not exceeded now, resume 
this query"
+                              << ", process memory info: "
+                              << 
GlobalMemoryArbitrator::process_memory_used_details_str()
+                              << ", wg info: " << wg->debug_string();
+                    
resource_ctx->task_controller()->set_memory_sufficient(true);
+                    query_it = queries_list.erase(query_it);
+                    continue;
+                }
+
                 // If workload group's memory usage > min memory, then it 
means the workload group use too much memory
                 // in memory contention state. Should just spill
-                if (wg->total_mem_used() > wg->min_memory_limit()) {
+                bool handle_query_now = wg->total_mem_used() > 
wg->min_memory_limit();
+                if (!handle_query_now) {
+                    // Other workload groups many use a lot of memory, should 
revoke memory from other workload groups
+                    // by cancelling their queries.
+                    int64_t revoked_size = revoke_memory_from_other_groups_();
+                    if (revoked_size > 0) {

Review Comment:
   [P1] Use the actual reclaimed amount before deferring the timeout. 
`revoke_memory_from_other_groups_()` returns its 10% target even when 
`WorkloadGroup::revoke_memory()` cancels nothing (for example, a group with 
>128 MiB excess spread across <=32 MiB queries). This positive result sets 
`revoking_memory_from_other_query_` and returns before the new timeout check; 
the next pass wakes the paused query, whose failed reservation requeues it with 
a fresh timer. Under sustained process pressure it can cycle indefinitely 
instead of reaching the promised cancellation. Return/check the actual 
reclaimed amount, then run the timeout fallback when it is zero.



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