mrhhsg commented on code in PR #67845:
URL: https://github.com/apache/doris/pull/67845#discussion_r4228617093


##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -300,6 +302,19 @@ void WorkloadGroupMgr::add_paused_query(const 
std::shared_ptr<ResourceContext>&
         LOG(INFO) << "Insert one new paused query: "
                   << resource_ctx->task_controller()->debug_string()
                   << ", workload group: " << wg->debug_string();
+    } else {
+        // Another task of the same query already failed a reservation. The 
query is resumed as
+        // a whole and every blocked task then retries its own request, so 
record the sum of the
+        // pending requests: the query is woken up only once all of them fit 
at the same time.
+        // Waking it up as soon as one of them fits would fail the others 
again and start a new
+        // wait from scratch.
+        LOG(INFO) << "Query: " << 
print_id(resource_ctx->task_controller()->task_id())
+                  << " is already paused, add the pending reservation "
+                  << PrettyPrinter::print_bytes(reserve_size) << " to the 
recorded "
+                  << PrettyPrinter::print_bytes(it->reserve_size_);
+        auto node = queries_list.extract(it);
+        node.value().reserve_size_ += reserve_size;

Review Comment:
   Fixed in 65ba0b876f7. `add_paused_query()` records the largest pending 
reservation again (as in a428b3e787d), so the query is woken up once that one 
fits and the tasks reserve one after another (a task releases its reservation 
after each block, so the pending requests do not have to fit at the same time). 
The wait is now bounded across retries instead of per entry: 
`TaskController::start_process_memory_wait()` keeps the time of the first 
reservation that failed since the last successful one, `add_paused_query()` 
uses it as the `enqueue_at` of a `PROCESS_MEMORY_EXCEEDED` entry, and 
`PipelineTask::_try_to_reserve_memory()` ends the wait when a reservation 
succeeds. A query resumed to retry that fails again before any of its 
reservations succeeded therefore continues the same wait and is still cancelled 
at `spill_in_paused_queue_timeout_ms`; a smaller sibling that does not fit next 
to the larger request can no longer restart it. `PausedQuery::enqueue_at` is a 
`MonotonicMillis()` time
 stamp now so that both use the same clock. Tests: 
`process_mem_exceeded_keeps_largest_pending_reservation` (64 MiB + 4 MiB 
pending, 65 MiB headroom resumes the query although 68 MiB does not fit), 
`process_mem_exceeded_largest_pending_reservation_cancels_after_timeout`, 
`process_mem_exceeded_failed_retry_continues_the_wait` (resumed past the 
timeout, re-paused by the failed retry, cancelled in the next round; without 
the change the new entry starts at 0 and stays paused), 
`process_mem_exceeded_successful_reservation_starts_new_wait`, and 
`PipelineTaskTest.TEST_RESERVE_MEMORY_FAIL_SPILLABLE` asserts that a successful 
reservation clears the wait start (fails without the change).



##########
be/test/testutil/mock/mock_query_task_controller.h:
##########
@@ -35,6 +38,52 @@ struct MockQueryTaskController : public QueryTaskController {
     }
 
     void set_cancelled_time(int64_t ctime) { cancelled_time_ = ctime; }
+
+    // Memory reclamation asks every candidate query whether it is cancelled 
while it scans a
+    // workload group, so this is where a test can let another query release 
memory during that
+    // scan.
+    bool is_cancelled() const override {
+        if (on_is_cancelled_) {
+            on_is_cancelled_();
+        }
+        return QueryTaskController::is_cancelled();
+    }
+
+    // Pipeline state seen by WorkloadGroupMgr::handle_single_query_. The 
query context of a
+    // unit test has no fragments, so the real implementation never reports a 
running or a
+    // revocable task; these knobs simulate them.
+    // NOLINTNEXTLINE(readability-make-member-function-const): overrides a 
non-const virtual.
+    void get_revocable_info(size_t* revocable_size, size_t* memory_usage,
+                            bool* has_running_task) override {
+        QueryTaskController::get_revocable_info(revocable_size, memory_usage, 
has_running_task);
+        *has_running_task = has_running_task_;
+    }
+
+    // The manager only checks whether the list is empty before it calls 
revoke_memory(), which
+    // is mocked below, so a placeholder entry is enough to stand for a 
revocable task.
+    std::vector<PipelineTask*> get_revocable_tasks() override {
+        if (on_get_revocable_tasks_) {
+            on_get_revocable_tasks_();
+        }
+        return has_revocable_task_ ? std::vector<PipelineTask*> {nullptr}
+                                   : std::vector<PipelineTask*> {};
+    }
+
+    Status revoke_memory() override {

Review Comment:
   Fixed in 65ba0b876f7. `MockQueryTaskController::revoke_memory()` now 
completes the spill the way the real implementation does: it clears 
`has_revocable_task_` (the spilled tasks have nothing left to revoke) and calls 
`set_memory_sufficient(true)`, which readies the memory dependency and resets 
the paused reason. `process_mem_exceeded_spills_revocable_tasks_at_hard_limit` 
asserts that the dependency is blocked after the pause, that the first round 
calls `revoke_memory()` once, resumes the query (`paused_reason().ok()`, 
dependency ready) and removes the entry, then pauses it again as the failed 
retry (dependency blocked again, no revocable task left) and checks that the 
next round cancels it at the hard limit without another `revoke_memory()` call.



##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -678,32 +702,50 @@ int64_t 
WorkloadGroupMgr::revoke_memory_from_other_groups_() {
                 // then not revoke memory from it.
                 continue;
             }
-            if (total_used_memory - min_memory_limit > max_exceeded_memory) {
-                max_wg = workload_group.second;
-                max_exceeded_memory = total_used_memory - min_memory_limit;
-            }
+            exceeded_wgs.emplace_back(total_used_memory - min_memory_limit, 
workload_group.second);
         }
     }
-    if (max_wg == nullptr) {
-        return 0;
-    }
-    if (max_exceeded_memory < 1 << 27) {
-        LOG(INFO) << "The workload group that exceed most memory is :"
-                  << max_wg->memory_debug_string() << ", max_exceeded_memory: "
-                  << PrettyPrinter::print(max_exceeded_memory, TUnit::BYTES)
-                  << " less than 128MB, no need to revoke memory";
-        return 0;
-    }
-    int64_t freed_mem = static_cast<int64_t>((double)max_exceeded_memory * 
0.1);
-    // Revoke 10% of memory from the workload group that exceed most memory
-    max_wg->revoke_memory(freed_mem, "exceed_memory", profile.get());
-    std::stringstream ss;
-    profile->pretty_print(&ss);
-    LOG(INFO) << fmt::format(
-            "[MemoryGC] process memory not enough, revoke memory from 
workload_group: {}, "
-            "free memory {}. cost(us): {}, details: {}",
-            max_wg->memory_debug_string(), 
PrettyPrinter::print_bytes(freed_mem),
-            watch.elapsed_time() / 1000, ss.str());
+    std::sort(exceeded_wgs.begin(), exceeded_wgs.end(),
+              [](const auto& lhs, const auto& rhs) { return lhs.first > 
rhs.first; });
+
+    int64_t freed_mem = 0;
+    size_t tried_wgs = 0;
+    for (const auto& [exceeded_memory, wg] : exceeded_wgs) {
+        if (exceeded_memory < 1 << 27) {
+            // The remaining workload groups exceed even less.
+            LOG(INFO) << "The workload group that exceed most memory among the 
untried ones is :"
+                      << wg->memory_debug_string() << ", exceeded_memory: "
+                      << PrettyPrinter::print(exceeded_memory, TUnit::BYTES)
+                      << " less than 128MB, no need to revoke memory";
+            break;
+        }

Review Comment:
   On the variant carried in the latest review (a peer falling below its 
minimum between the walk's `refresh_memory_usage()` and 
`WorkloadGroup::revoke_memory()`): the window left is the few calls between 
that refresh and `memory_used()` at the start of `revoke_memory()`, with no 
waiting in between; the window that mattered (scanning the earlier peers for up 
to `revoke_memory_max_tolerance_ms`) is what the refresh in the walk closes. 
The same check-then-act window exists inside 
`MemoryReclamation::revoke_tasks_memory()` itself, which snapshots the 
consumption and then cancels, and `WorkloadGroup::revoke_memory()` is shared 
with the memory gc callers; adding a min-memory eligibility check inside it 
would change their semantics, so I am leaving this as is in this PR.



##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -824,38 +839,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_process_memory_exceeded_) has already
+        // resumed the query if the process is no longer above the soft memory 
limit, so the

Review Comment:
   On the QUERY-first variant carried in the latest review: 
`TaskController::update_paused_reason()` keeps `QUERY_MEMORY_EXCEEDED` over a 
later `PROCESS_MEMORY_EXCEEDED` failure of a sibling task, which is 
pre-existing precedence that this PR does not change. A QUERY reason means the 
query is over its own limit, and the QUERY route either spills it or resumes it 
with `disable_reserve_memory()`; from then on every allocation goes through 
`Allocator::sys_memory_check()`, which checks `is_exceed_hard_mem_limit()` 
directly and fails the allocation regardless of `disable_memory_gc`. The 
running-task exit on that route lasts for the current operator call only. The 
mixed-reason handling belongs to the existing QUERY route rather than to the 
PROCESS route this PR changes, so no change here.



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