mrhhsg commented on code in PR #67846:
URL: https://github.com/apache/doris/pull/67846#discussion_r4226815573
##########
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:
Fixed in 7e69ac0ae0491f9488d534f8a231370d053f610b.
`revoke_memory_from_other_groups_()` now returns the value of
`WorkloadGroup::revoke_memory()`, i.e. the memory of the queries that were
actually cancelled (plus the memory of queries cancelled within
`revoke_memory_max_tolerance_ms` that is still being released), instead of the
10% target. When nothing can be cancelled it returns 0, so the below-min route
falls through to the hard-limit / timeout fallback instead of setting
`revoking_memory_from_other_query_` and resuming the query with a fresh timer
on the next pass.
Added
`process_mem_exceeded_other_group_not_reclaimable_cancels_after_timeout`: a
peer workload group 140 MiB over its minimum made of 8 queries of 30 MiB each
(all excluded by `EXCLUDE_IS_SMALL`). It asserts
`revoke_memory_from_other_groups_() == 0`, that the paused query stays paused
across passes with `revoking_memory_from_other_query_` false, that it is
cancelled once the timeout elapses, and that the peer queries are not
cancelled. The three `ProcessMemoryNotEnough` expectations were updated to the
actual released amount (300 MiB / 500 MiB / 500 MiB).
##########
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:
Fixed in 7e69ac0ae0491f9488d534f8a231370d053f610b. The below-min route now
dispatches to `handle_single_query_()` when
`GlobalMemoryArbitrator::is_exceed_hard_mem_limit()` is true or the timeout has
elapsed (same as master's `handle_process_memory_exceeded_()` after
0df5f800262), so a paused query at the hard limit is spilled or cancelled on
the next pass without depending on memory GC.
Added `process_mem_exceeded_below_min_memory_cancels_at_hard_limit`
(`total_mem_used() <= min_memory_limit()`, no other workload group, hard limit
exceeded), which asserts immediate cancellation with `exceed hard limit: true`.
##########
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:
Not changed, because the early return does not block hard-limit handling
indefinitely and the alternative would cancel queries that can spill:
- The `has_running_task` return is pre-existing on branch-4.1 and gated the
previous immediate cancellation in exactly the same way; this PR does not
change when that path is reached.
- The running state is transient. `PipelineTask::execute()` evaluates
`_is_blocked()` at the top of every loop iteration, and `_is_blocked()`
includes `_memory_sufficient_dependency->is_blocked_by(...)`. Once the query is
marked memory-insufficient, every sibling task returns from `execute()` after
the operator call it is currently in, and `TaskScheduler` clears `_running` in
its `Defer` right after `execute()` returns. So `has_running_task` is bounded
by one operator call, and the next maintenance pass
(`memory_maintenance_sleep_time_ms`, 50 ms) reaches the hard-limit cancellation.
- Cancelling before that guard would cancel a query with revocable memory
while a sibling is finishing its call, although `handle_single_query_()` would
spill it on the next pass. The guard exists precisely so that the spill path
waits until no task is running.
The same question is open on the master PR (#67845) and is kept consistent
with it here.
##########
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:
Fixed in 7e69ac0ae0491f9488d534f8a231370d053f610b. Added
`MemInfo::set_sys_mem_available_for_test()` (BE_TEST only, returns the previous
value), and the process-pressure cases now use a `ScopedProcessMemLimits` RAII
helper that pins `_s_sys_mem_available` to a large value for the duration of
the case and restores the mem limit, soft limit and available memory
afterwards. Both `is_exceed_soft_mem_limit()` and `is_exceed_hard_mem_limit()`
therefore depend only on the limits each case sets. When asserting resume, the
soft limit is set to the same large value instead of being restored to the
host-derived original, so the `exceed hard limit: false` and resume assertions
no longer depend on the runner's memory state.
--
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]