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

morningman pushed a commit to branch branch-4.0
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.0 by this push:
     new baa67eadff1 branch-4.0: [fix](memory) should update memory more 
quickly when user changes cgroup directly (#65713)
baa67eadff1 is described below

commit baa67eadff1e58bee8e758658557198a0c1fed30
Author: yiguolei <[email protected]>
AuthorDate: Thu Jul 30 20:48:31 2026 +0800

    branch-4.0: [fix](memory) should update memory more quickly when user 
changes cgroup directly (#65713)
    
    ### What problem does this PR solve?
    
    Issue Number: close #xxx
    pick  #65695
---
 be/src/runtime/query_context.cpp                   | 24 ++++---
 be/src/runtime/workload_group/workload_group.cpp   | 37 ++++++++--
 .../workload_group/workload_group_manager.cpp      | 28 ++++----
 .../runtime/workload_management/memory_context.h   |  4 ++
 be/src/util/mem_info.h                             |  5 ++
 .../workload_group/workload_group_manager_test.cpp | 81 ++++++++++++++++++++--
 6 files changed, 148 insertions(+), 31 deletions(-)

diff --git a/be/src/runtime/query_context.cpp b/be/src/runtime/query_context.cpp
index e4bbeef60dd..9e39ff550cd 100644
--- a/be/src/runtime/query_context.cpp
+++ b/be/src/runtime/query_context.cpp
@@ -143,33 +143,38 @@ QueryContext::QueryContext(TUniqueId query_id, ExecEnv* 
exec_env,
 }
 
 void QueryContext::_init_query_mem_tracker() {
+    // If user not set query limit, will use default 1TB memory limit. It is 
large enough to cover most cases.
+    constexpr int64_t DEFAULT_QUERY_MEM_LIMIT = 1LL << 60;
     bool has_query_mem_limit = _query_options.__isset.mem_limit && 
(_query_options.mem_limit > 0);
-    int64_t bytes_limit = has_query_mem_limit ? _query_options.mem_limit : -1;
-    if (bytes_limit > MemInfo::mem_limit() || bytes_limit == -1) {
-        VLOG_NOTICE << "Query memory limit " << 
PrettyPrinter::print(bytes_limit, TUnit::BYTES)
+    int64_t user_set_mem_limit =
+            has_query_mem_limit ? _query_options.mem_limit : 
DEFAULT_QUERY_MEM_LIMIT;
+    int64_t adjusted_mem_limit = user_set_mem_limit;
+    if (adjusted_mem_limit > MemInfo::mem_limit()) {
+        VLOG_NOTICE << "Query memory limit "
+                    << PrettyPrinter::print(user_set_mem_limit, TUnit::BYTES)
                     << " exceeds process memory limit of "
                     << PrettyPrinter::print(MemInfo::mem_limit(), TUnit::BYTES)
-                    << " OR is -1. Using process memory limit instead.";
-        bytes_limit = MemInfo::mem_limit();
+                    << ". Using process memory limit instead.";
+        adjusted_mem_limit = MemInfo::mem_limit();
     }
     // If the query is a pure load task(streamload, routine load, group 
commit), then it should not use
     // memlimit per query to limit their memory usage.
     if (is_pure_load_task()) {
-        bytes_limit = MemInfo::mem_limit();
+        adjusted_mem_limit = MemInfo::mem_limit();
     }
     std::shared_ptr<MemTrackerLimiter> query_mem_tracker;
     if (_query_options.query_type == TQueryType::SELECT) {
         query_mem_tracker = MemTrackerLimiter::create_shared(
                 MemTrackerLimiter::Type::QUERY, fmt::format("Query#Id={}", 
print_id(_query_id)),
-                bytes_limit);
+                adjusted_mem_limit);
     } else if (_query_options.query_type == TQueryType::LOAD) {
         query_mem_tracker = MemTrackerLimiter::create_shared(
                 MemTrackerLimiter::Type::LOAD, fmt::format("Load#Id={}", 
print_id(_query_id)),
-                bytes_limit);
+                adjusted_mem_limit);
     } else if (_query_options.query_type == TQueryType::EXTERNAL) { // 
spark/flink/etc..
         query_mem_tracker = MemTrackerLimiter::create_shared(
                 MemTrackerLimiter::Type::QUERY, fmt::format("External#Id={}", 
print_id(_query_id)),
-                bytes_limit);
+                adjusted_mem_limit);
     } else {
         LOG(FATAL) << "__builtin_unreachable";
         __builtin_unreachable();
@@ -187,6 +192,7 @@ void QueryContext::_init_query_mem_tracker() {
     
query_mem_tracker->set_enable_check_limit(!(_query_options.__isset.enable_reserve_memory
 &&
                                                 
_query_options.enable_reserve_memory));
     _resource_ctx->memory_context()->set_mem_tracker(query_mem_tracker);
+    
_resource_ctx->memory_context()->set_user_set_mem_limit(user_set_mem_limit);
 }
 
 void QueryContext::_init_resource_context() {
diff --git a/be/src/runtime/workload_group/workload_group.cpp 
b/be/src/runtime/workload_group/workload_group.cpp
index f289d29b2b3..8c371579a7a 100644
--- a/be/src/runtime/workload_group/workload_group.cpp
+++ b/be/src/runtime/workload_group/workload_group.cpp
@@ -154,10 +154,7 @@ void WorkloadGroup::check_and_update(const 
WorkloadGroupInfo& wg_info) {
         return;
     }
     std::lock_guard<std::shared_mutex> wl {_mutex};
-    // In serverless mode, user may modify cgroup's memory limit directly and 
workload group's config
-    // is not changed. So that we should update workload group's config ignore 
version.
-    if (wg_info.version > _version ||
-        (wg_info.version == _version && _memory_limit != 
wg_info.memory_limit)) {
+    if (wg_info.version > _version) {
         _name = wg_info.name;
         _version = wg_info.version;
         _min_cpu_percent = wg_info.min_cpu_percent;
@@ -183,6 +180,38 @@ void WorkloadGroup::check_and_update(const 
WorkloadGroupInfo& wg_info) {
 
 // MemtrackerLimiter is not removed during query context release, so that 
should remove it here.
 int64_t WorkloadGroup::refresh_memory_usage() {
+    {
+        // In serverless mode, user may modify cgroup's memory limit directly 
and workload group's config
+        // is not changed. So that we should update workload group's config 
ignore version.
+        std::lock_guard<std::shared_mutex> wl {_mutex};
+        const int max_memory_percent = 
_max_memory_percent.load(std::memory_order_relaxed);
+        const std::string mem_limit_str = fmt::format("{}%", 
max_memory_percent);
+        bool is_percent = true;
+        const int64_t new_memory_limit =
+                ParseUtil::parse_mem_spec(mem_limit_str, -1, 
MemInfo::mem_limit(), &is_percent);
+        DCHECK(is_percent) << "mem_limit_str: " << mem_limit_str;
+        if (new_memory_limit != _memory_limit.load(std::memory_order_relaxed)) 
{
+            LOG(INFO) << fmt::format(
+                    "Workload group id:{}, name:{}, "
+                    "memory_limit changed from {} to {}",
+                    _id, _name,
+                    
PrettyPrinter::print(_memory_limit.load(std::memory_order_relaxed),
+                                         TUnit::BYTES),
+                    PrettyPrinter::print(new_memory_limit, TUnit::BYTES));
+
+            _memory_limit.store(new_memory_limit, std::memory_order_relaxed);
+            if (max_memory_percent == 0) {
+                _min_memory_limit.store(0, std::memory_order_relaxed);
+            } else {
+                const int min_memory_percent = 
_min_memory_percent.load(std::memory_order_relaxed);
+                const int64_t new_min_memory_limit =
+                        
static_cast<int64_t>(static_cast<double>(new_memory_limit) *
+                                             min_memory_percent / 
max_memory_percent);
+                _min_memory_limit.store(new_min_memory_limit, 
std::memory_order_relaxed);
+            }
+        }
+    }
+
     int64_t fragment_used_memory = 0;
     {
         std::shared_lock<std::shared_mutex> r_lock(_mutex);
diff --git a/be/src/runtime/workload_group/workload_group_manager.cpp 
b/be/src/runtime/workload_group/workload_group_manager.cpp
index d226120a1ad..6ad643c719a 100644
--- a/be/src/runtime/workload_group/workload_group_manager.cpp
+++ b/be/src/runtime/workload_group/workload_group_manager.cpp
@@ -231,16 +231,14 @@ void 
WorkloadGroupMgr::refresh_workload_group_memory_state() {
     for (auto& [wg_id, wg] : _workload_groups) {
         all_workload_groups_mem_usage += wg->refresh_memory_usage();
     }
-    if (all_workload_groups_mem_usage <= 0) {
-        return;
+    if (all_workload_groups_mem_usage > 0) {
+        std::string debug_msg = fmt::format(
+                "\nProcess Memory Summary: {}, {}, all workload groups memory 
usage: {}",
+                
doris::GlobalMemoryArbitrator::process_memory_used_details_str(),
+                doris::GlobalMemoryArbitrator::sys_mem_available_details_str(),
+                PrettyPrinter::print(all_workload_groups_mem_usage, 
TUnit::BYTES));
+        LOG_EVERY_T(INFO, 60) << debug_msg;
     }
-
-    std::string debug_msg =
-            fmt::format("\nProcess Memory Summary: {}, {}, all workload groups 
memory usage: {}",
-                        
doris::GlobalMemoryArbitrator::process_memory_used_details_str(),
-                        
doris::GlobalMemoryArbitrator::sys_mem_available_details_str(),
-                        PrettyPrinter::print(all_workload_groups_mem_usage, 
TUnit::BYTES));
-    LOG_EVERY_T(INFO, 60) << debug_msg;
     for (auto& wg : _workload_groups) {
         update_queries_limit_(wg.second, false);
     }
@@ -865,11 +863,13 @@ void 
WorkloadGroupMgr::update_queries_limit_(WorkloadGroupPtr wg, bool enable_ha
         // If the query is a pure load task, then should not modify its limit. 
Or it will reserve
         // memory failed and we did not hanle it.
         if (!resource_ctx->task_controller()->is_pure_load_task()) {
-            // If user's set mem limit is less than query weighted mem limit, 
then should not modify its limit.
-            // Use user settings.
-            if (resource_ctx->memory_context()->user_set_mem_limit() > 
query_weighted_mem_limit) {
-                
resource_ctx->memory_context()->set_mem_limit(query_weighted_mem_limit);
-            }
+            // The effective limit should be min(user_set_mem_limit, 
query_weighted_mem_limit).
+            // This ensures limits are both lowered under memory pressure and 
restored when
+            // pressure eases (e.g., when concurrent queries finish or WG 
memory drops below
+            // low watermark).
+            int64_t effective_limit = 
std::min(resource_ctx->memory_context()->user_set_mem_limit(),
+                                               query_weighted_mem_limit);
+            resource_ctx->memory_context()->set_mem_limit(effective_limit);
             resource_ctx->memory_context()->set_adjusted_mem_limit(
                     expected_query_weighted_mem_limit);
         }
diff --git a/be/src/runtime/workload_management/memory_context.h 
b/be/src/runtime/workload_management/memory_context.h
index 027f2fcb61d..494bbd79ed7 100644
--- a/be/src/runtime/workload_management/memory_context.h
+++ b/be/src/runtime/workload_management/memory_context.h
@@ -84,6 +84,10 @@ public:
         adjusted_mem_limit_ = mem_tracker_->limit();
     }
 
+    void set_user_set_mem_limit(int64_t user_set_mem_limit) {
+        user_set_mem_limit_ = user_set_mem_limit;
+    }
+
     // This method is called by workload group manager to set query's memlimit 
using slot
     // If user set query limit explicitly, then should use less one
     void set_mem_limit(int64_t new_mem_limit) const { 
mem_tracker_->set_limit(new_mem_limit); }
diff --git a/be/src/util/mem_info.h b/be/src/util/mem_info.h
index 4f0ddd2f57b..113b3352c25 100644
--- a/be/src/util/mem_info.h
+++ b/be/src/util/mem_info.h
@@ -83,6 +83,11 @@ public:
         DCHECK(_s_initialized);
         return _s_mem_limit.load(std::memory_order_relaxed);
     }
+#ifdef BE_TEST
+    static void set_mem_limit_for_test(int64_t mem_limit) {
+        _s_mem_limit.store(mem_limit, std::memory_order_relaxed);
+    }
+#endif
     static inline std::string mem_limit_str() {
         DCHECK(_s_initialized);
         return 
PrettyPrinter::print(_s_mem_limit.load(std::memory_order_relaxed), 
TUnit::BYTES);
diff --git a/be/test/runtime/workload_group/workload_group_manager_test.cpp 
b/be/test/runtime/workload_group/workload_group_manager_test.cpp
index f81488928c9..18e1af25742 100644
--- a/be/test/runtime/workload_group/workload_group_manager_test.cpp
+++ b/be/test/runtime/workload_group/workload_group_manager_test.cpp
@@ -35,6 +35,9 @@
 #include "runtime/query_context.h"
 #include "runtime/runtime_query_statistics_mgr.h"
 #include "runtime/workload_group/workload_group.h"
+#include "testutil/mock/mock_query_task_controller.h"
+#include "util/defer_op.h"
+#include "util/mem_info.h"
 #include "vec/spill/spill_stream_manager.h"
 
 namespace doris {
@@ -84,10 +87,13 @@ protected:
     }
 
 private:
-    std::shared_ptr<QueryContext> 
_generate_on_query(std::shared_ptr<WorkloadGroup>& wg) {
+    std::shared_ptr<QueryContext> 
_generate_on_query(std::shared_ptr<WorkloadGroup>& wg,
+                                                     int64_t mem_limit = 1024L 
* 1024 * 128,
+                                                     bool has_mem_limit = 
false) {
         TQueryOptions query_options;
         query_options.query_type = TQueryType::SELECT;
-        query_options.mem_limit = 1024L * 1024 * 128;
+        query_options.mem_limit = mem_limit;
+        query_options.__isset.mem_limit = has_mem_limit;
         query_options.query_slot_count = 1;
         TNetworkAddress fe_address;
         fe_address.hostname = "127.0.0.1";
@@ -125,6 +131,59 @@ TEST_F(WorkloadGroupManagerTest, 
get_or_create_workload_group) {
     ASSERT_EQ(wg->id(), 0);
 }
 
+TEST_F(WorkloadGroupManagerTest, refresh_memory_usage_updates_memory_limits) {
+    const int64_t original_mem_limit = MemInfo::mem_limit();
+    Defer restore_mem_limit {[&]() { 
MemInfo::set_mem_limit_for_test(original_mem_limit); }};
+    const int64_t initial_mem_limit = 1024L * 1024 * 1024;
+    MemInfo::set_mem_limit_for_test(initial_mem_limit);
+
+    WorkloadGroupInfo wg_info {.id = 1,
+                               .memory_limit = initial_mem_limit / 2,
+                               .min_memory_percent = 25,
+                               .max_memory_percent = 50};
+    auto wg = _wg_manager->get_or_create_workload_group(wg_info);
+
+    EXPECT_EQ(wg->memory_limit(), initial_mem_limit / 2);
+    EXPECT_EQ(wg->min_memory_limit(), initial_mem_limit / 4);
+
+    const int64_t updated_mem_limit = initial_mem_limit * 2;
+    MemInfo::set_mem_limit_for_test(updated_mem_limit);
+    wg->refresh_memory_usage();
+
+    EXPECT_EQ(wg->memory_limit(), updated_mem_limit / 2);
+    EXPECT_EQ(wg->min_memory_limit(), updated_mem_limit / 4);
+}
+
+TEST_F(WorkloadGroupManagerTest, 
refresh_restores_query_limit_after_cgroup_expands) {
+    const int64_t original_mem_limit = MemInfo::mem_limit();
+    Defer restore_mem_limit {[&]() { 
MemInfo::set_mem_limit_for_test(original_mem_limit); }};
+    const int64_t small_mem_limit = 1024L * 1024 * 20;
+    const int64_t large_mem_limit = 1024L * 1024 * 100;
+    MemInfo::set_mem_limit_for_test(small_mem_limit);
+
+    WorkloadGroupInfo wg_info {.id = 1,
+                               .memory_limit = small_mem_limit,
+                               .max_memory_percent = 100,
+                               .slot_mem_policy = TWgSlotMemoryPolicy::NONE};
+    auto wg = _wg_manager->get_or_create_workload_group(wg_info);
+    auto query_context = _generate_on_query(wg, large_mem_limit, true);
+    auto query_without_mem_limit = _generate_on_query(wg);
+
+    ASSERT_EQ(query_context->resource_ctx()->memory_context()->mem_limit(), 
small_mem_limit);
+    
ASSERT_EQ(query_context->resource_ctx()->memory_context()->user_set_mem_limit(),
+              large_mem_limit);
+    
ASSERT_EQ(query_without_mem_limit->resource_ctx()->memory_context()->mem_limit(),
+              small_mem_limit);
+
+    MemInfo::set_mem_limit_for_test(large_mem_limit);
+    _wg_manager->refresh_workload_group_memory_state();
+
+    ASSERT_EQ(wg->memory_limit(), large_mem_limit);
+    ASSERT_EQ(query_context->resource_ctx()->memory_context()->mem_limit(), 
large_mem_limit);
+    
ASSERT_EQ(query_without_mem_limit->resource_ctx()->memory_context()->mem_limit(),
+              large_mem_limit);
+}
+
 // Query is paused due to query memlimit exceed, after waiting in queue for  
spill_in_paused_queue_timeout_ms
 // it should be resumed
 TEST_F(WorkloadGroupManagerTest, query_exceed) {
@@ -208,8 +267,13 @@ TEST_F(WorkloadGroupManagerTest, wg_exceed2) {
 // query limit > workload group limit
 // query's limit will be set to workload group limit
 TEST_F(WorkloadGroupManagerTest, wg_exceed3) {
-    WorkloadGroupInfo wg_info {
-            .id = 1, .memory_limit = 1024L * 1024, .slot_mem_policy = 
TWgSlotMemoryPolicy::NONE};
+    const int64_t original_mem_limit = MemInfo::mem_limit();
+    Defer restore_mem_limit {[&]() { 
MemInfo::set_mem_limit_for_test(original_mem_limit); }};
+    MemInfo::set_mem_limit_for_test(1024L * 1024 * 100);
+    WorkloadGroupInfo wg_info {.id = 1,
+                               .memory_limit = 1024L * 1024,
+                               .max_memory_percent = 1,
+                               .slot_mem_policy = TWgSlotMemoryPolicy::NONE};
     auto wg = _wg_manager->get_or_create_workload_group(wg_info);
     auto query_context = _generate_on_query(wg);
 
@@ -250,6 +314,9 @@ TEST_F(WorkloadGroupManagerTest, wg_exceed3) {
 
 // TWgSlotMemoryPolicy::FIXED
 TEST_F(WorkloadGroupManagerTest, wg_exceed4) {
+    const int64_t original_mem_limit = MemInfo::mem_limit();
+    Defer restore_mem_limit {[&]() { 
MemInfo::set_mem_limit_for_test(original_mem_limit); }};
+    MemInfo::set_mem_limit_for_test(1024L * 1024 * 100);
     WorkloadGroupInfo wg_info {.id = 1,
                                .memory_limit = 1024L * 1024 * 100,
                                .memory_low_watermark = 80,
@@ -287,6 +354,9 @@ TEST_F(WorkloadGroupManagerTest, wg_exceed4) {
 
 // TWgSlotMemoryPolicy::DYNAMIC
 TEST_F(WorkloadGroupManagerTest, wg_exceed5) {
+    const int64_t original_mem_limit = MemInfo::mem_limit();
+    Defer restore_mem_limit {[&]() { 
MemInfo::set_mem_limit_for_test(original_mem_limit); }};
+    MemInfo::set_mem_limit_for_test(1024L * 1024 * 100);
     WorkloadGroupInfo wg_info {.id = 1,
                                .memory_limit = 1024L * 1024 * 100,
                                .min_memory_percent = 10,
@@ -406,6 +476,9 @@ TEST_F(WorkloadGroupManagerTest, query_released) {
 }
 
 TEST_F(WorkloadGroupManagerTest, ProcessMemoryNotEnough) {
+    const int64_t original_mem_limit = MemInfo::mem_limit();
+    Defer restore_mem_limit {[&]() { 
MemInfo::set_mem_limit_for_test(original_mem_limit); }};
+    MemInfo::set_mem_limit_for_test(1024L * 1024 * 1000);
     WorkloadGroupInfo wg1_info {.id = 1,
                                 .memory_limit = 1024L * 1024 * 1000,
                                 .min_memory_percent = 10,


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

Reply via email to