mrhhsg commented on code in PR #67845:
URL: https://github.com/apache/doris/pull/67845#discussion_r4226850835
##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -823,40 +848,93 @@ bool WorkloadGroupMgr::handle_single_query_(const
std::shared_ptr<ResourceContex
requestor->task_controller()->cancel(error_status);
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)) {
+ }
+}
+
+void WorkloadGroupMgr::spill_query_(const std::shared_ptr<ResourceContext>&
requestor) {
+ SCOPED_ATTACH_TASK(requestor);
+ auto status = requestor->task_controller()->revoke_memory();
+ if (!status.ok()) {
+ requestor->task_controller()->cancel(status);
+ }
+}
+
+bool WorkloadGroupMgr::resolve_process_memory_exceeded_query_(
+ const std::shared_ptr<ResourceContext>& requestor, size_t
size_to_reserve,
+ int64_t time_in_queue, size_t memory_usage, bool has_running_task) {
+ const auto query_id = print_id(requestor->task_controller()->task_id());
+ const auto wg = requestor->workload_group();
+
+ // The caller observed the pressure before it inspected the pipeline tasks
and tried to
+ // revoke memory from other workload groups, during which other queries or
the cache may have
+ // released memory. Re-check the recorded reservation right before acting,
so that the query
+ // is never cancelled from a stale observation.
+ if (!GlobalMemoryArbitrator::is_exceed_soft_mem_limit(size_to_reserve)) {
+ 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;
+ }
+
+ const bool exceed_hard_mem_limit =
GlobalMemoryArbitrator::is_exceed_hard_mem_limit();
+ if (has_running_task) {
+ // Spilling needs every task of the query to be idle. Below the hard
limit wait for the
+ // running task to yield, the query is handled again in the next round.
+ if (!exceed_hard_mem_limit) {
Review Comment:
Fixed in 7cd9a4c0582. `resolve_process_memory_exceeded_query_()` no longer
returns early for a running task: the task only prevents collecting revocable
tasks (so no spill is started while it runs), and the query then goes through
the same `!exceed_hard_mem_limit && time_in_queue <
spill_in_paused_queue_timeout_ms` wait as a query without revocable memory. At
the timeout or at the hard limit it is cancelled, which is safe while the task
runs; the cancel message and the keep-paused log report `has running task`.
`process_mem_exceeded_running_task_waits_below_hard_limit` was replaced by
`process_mem_exceeded_running_task_cancels_after_timeout`: with a running task
and revocable memory the query waits below the timeout without
`revoke_memory()` being called, and is cancelled once the timer is backdated
past the timeout.
##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -620,11 +635,11 @@ bool WorkloadGroupMgr::handle_process_memory_exceeded_(
// TODO revoke from memtable
- // Fallback: if we have waited too long and cannot revoke from anywhere,
- // cancel the query or disable reserve memory to let it proceed.
- if (query_it->elapsed_time() > config::spill_in_paused_queue_timeout_ms) {
- // Cannot spill (no revocable memory), cannot revoke from other WGs,
- // and process memory is still exceeded. Cancel the query to protect
the system.
+ // Fallback: if the process reaches the hard limit or we have waited too
long and cannot
+ // revoke from anywhere, let this query spill or cancel it to protect the
process.
Review Comment:
The remaining case (a peer cancelled on a previous pass that still holds its
memory) is bounded, and I would keep it:
`MemoryReclamation::revoke_tasks_memory()` counts a query that is already being
cancelled as freed memory only while `now - cancelled_time() <=
revoke_memory_max_tolerance_ms` (3s by default, `memory_reclamation.cpp`, the
`keep_wait_cancelling_tasks` branch); after that the query is skipped and not
counted. So while the peer's cancellation is in flight, the requestor waits for
that release instead of being cancelled, which is the same semantic the process
full GC uses, and the right outcome: cancelling the small below-minimum
requestor during those milliseconds would not free the memory the process
needs, while the peer's release will. Once the tolerance is over,
`revoke_memory_from_other_groups_()` returns 0 and the requestor reaches the
hard-limit/timeout fallback. The resume between passes comes from the existing
phase-2 rule (`resume_paused_queries_after_revoke
_` resumes everything when no cancelled query is in the paused list); it is
not introduced here and does not extend the wait beyond the tolerance.
Added `process_mem_exceeded_below_min_memory_waits_for_cancelling_peer` in
7cd9a4c0582 to pin this down: at the hard limit the first pass cancels the 300
MiB peer and keeps the requestor paused; after the phase-2 resume and a
re-pause, the still-resident peer is counted again and the requestor keeps
waiting; after the peer's `cancelled_time` is moved past
`revoke_memory_max_tolerance_ms`, revoking returns 0 and the requestor is
cancelled at the hard limit.
##########
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
+ // process memory is still exceeded here.
Review Comment:
Follow-up in 7cd9a4c0582: the soft-limit recheck in
`resolve_process_memory_exceeded_query_()` now runs after
`get_revocable_tasks()` (the last inspection of the pipeline tasks) and
immediately before the spill/timeout/hard-limit decision, so there is no
further scan between the recheck and the action. The mock hook moved to
`get_revocable_tasks()` and
`process_mem_exceeded_recovery_before_timeout_decision_resumes` relieves the
pressure from there with the timer already expired.
--
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]