mrhhsg commented on code in PR #67845:
URL: https://github.com/apache/doris/pull/67845#discussion_r4228101792
##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -678,32 +715,64 @@ 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& [snapshot_exceeded_memory, wg] : exceeded_wgs) {
+ if (snapshot_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(snapshot_exceeded_memory,
TUnit::BYTES)
+ << " less than 128MB, no need to revoke memory";
+ break;
+ }
+ // The snapshot above was taken before the earlier workload groups
were scanned, during
+ // which this one may have released memory (its queries finished or
were cancelled).
+ // Decide from its current usage: the min memory stays reserved for a
workload group that
+ // no longer exceeds it, and the amount to revoke is what it exceeds
by now.
+ const int64_t exceeded_memory = wg->total_mem_used() -
wg->min_memory_limit();
+ if (exceeded_memory < 1 << 27) {
+ LOG(INFO) << "The workload group " << wg->memory_debug_string()
+ << " exceeded its min memory by "
+ << PrettyPrinter::print(snapshot_exceeded_memory,
TUnit::BYTES)
+ << " when the walk started, now by "
+ << PrettyPrinter::print(exceeded_memory, TUnit::BYTES)
+ << " less than 128MB, skip it";
+ continue;
+ }
+ auto need_free_mem = static_cast<int64_t>((double)exceeded_memory *
0.1);
+ // Revoke 10% of memory from the workload group that exceed most memory
+ freed_mem = wg->revoke_memory(need_free_mem, "exceed_memory",
profile.get());
Review Comment:
Fixed in b1e7d9c1280. `revoke_memory_from_other_groups_()` now takes the
requestor's recorded `reserve_size`: after at least one peer has been tried, it
re-checks `is_exceed_soft_mem_limit(reserve_size)` before touching the next
peer and stops the walk (returning 0) as soon as the reservation fits, so no
query of a later peer is cancelled for pressure that is gone.
`handle_process_memory_exceeded_()` then re-checks the reservation after a zero
walk and resumes the requestor through the new `resume_process_paused_query_()`
instead of falling through to the hard-limit/timeout fallback. Added
`process_mem_exceeded_below_min_memory_resumes_when_relieved_during_peer_scan`:
the first peer only holds a query whose cancellation exceeded
`revoke_memory_max_tolerance_ms` (frees nothing) and the pressure is relieved
from its `is_cancelled()` hook while it is scanned; the second peer's 250 MiB
query stays, `revoking_memory_from_other_query_` stays false and the requestor
is resumed. Without t
he source change the case fails on `second_peer_query->is_cancelled()`.
##########
be/test/runtime/workload_group/workload_group_manager_test.cpp:
##########
@@ -153,9 +175,127 @@ class WorkloadGroupManagerTest : public testing::Test {
}
}
+ // Helpers for the PROCESS_MEMORY_EXCEEDED cases.
+
+ // A workload group whose min memory limit is 100 MiB, so that a query
consuming more than
+ // that routes into handle_single_query_ directly, and a smaller one is
routed through the
+ // other workload groups first.
+ std::shared_ptr<WorkloadGroup> _create_wg_with_min_memory(uint64_t id) {
+ MemInfo::set_mem_limit_for_test(kProcessMemLimitForWg);
+ WorkloadGroupInfo wg_info {.id = id,
+ .memory_limit = kProcessMemLimitForWg,
+ .min_memory_percent = 10,
+ .max_memory_percent = 100};
+ auto wg = _wg_manager->get_or_create_workload_group(wg_info);
+ EXPECT_EQ(wg->min_memory_limit(), 1024L * 1024 * 100);
+ return wg;
+ }
+
+ // A query on `wg` that consumes `memory` bytes until TearDown.
+ std::shared_ptr<QueryContext>
_create_query_with_memory(std::shared_ptr<WorkloadGroup>& wg,
+ int64_t memory) {
+ auto query = _generate_on_query(wg);
+ query->query_mem_tracker()->consume(memory);
+ _consumed_memory.emplace_back(query, memory);
+ wg->refresh_memory_usage();
+ return query;
+ }
+
+ // Release `memory` bytes of `query` while a case has lowered the process
memory limit.
+ // refresh_memory_usage() also derives the workload group limits from the
process memory
+ // limit, so it runs under the limit the workload group was created with.
+ void _release_query_memory(const std::shared_ptr<WorkloadGroup>& wg,
+ const std::shared_ptr<QueryContext>& query,
int64_t memory) {
+ query->query_mem_tracker()->consume(-memory);
+ _consumed_memory.emplace_back(query, -memory);
+ const int64_t mem_limit = MemInfo::mem_limit();
+ MemInfo::set_mem_limit_for_test(kProcessMemLimitForWg);
+ wg->refresh_memory_usage();
+ MemInfo::set_mem_limit_for_test(mem_limit);
+ ASSERT_EQ(wg->min_memory_limit(), 1024L * 1024 * 100);
+ }
+
+ // Process soft memory limit is exceeded, hard memory limit is not.
+ static void _exceed_process_soft_mem_limit() {
+ MemInfo::set_mem_limit_for_test(kLargeProcessMemLimit);
+ MemInfo::set_soft_mem_limit_for_test(1);
+
ASSERT_TRUE(GlobalMemoryArbitrator::is_exceed_soft_mem_limit(kProcessPausedReserveSize));
+ ASSERT_FALSE(GlobalMemoryArbitrator::is_exceed_hard_mem_limit());
Review Comment:
Fixed in b1e7d9c1280. The fixture now controls both sides of the predicates:
`SetUp` saves the runner's `_s_sys_mem_available` and pins it far above the
water marks (`TearDown` restores it and resets
`refresh_interval_memory_growth`). The hard limit is no longer simulated by
`mem_limit = 0` (which also made `refresh_memory_usage()` derive a zero
workload group min memory) but by setting
`GlobalMemoryArbitrator::refresh_interval_memory_growth` to the process limit,
so the process usage (vm_rss is never refreshed in a unit test) exceeds
`mem_limit` while the workload group limits stay intact;
`_exceed_process_soft_mem_limit()` keeps `mem_limit` and sets the soft limit to
1 with the growth reset; `_relieve_process_mem_limit()` resets the growth,
gives both process limits the explicit 1000 MiB headroom and sets the system
available memory back to the pinned value. The system-boundary case lowers the
system available memory to the warning mark + 16 MiB on top of that explicit
process h
eadroom and no longer saves/restores locally.
##########
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:
Follow-up in b1e7d9c1280: the re-check in the walk read `total_mem_used()`,
which is the usage cached by the last maintenance refresh, so a peer that
released memory while an earlier peer was scanned still showed its snapshot
usage. The walk now uses `wg->refresh_memory_usage()` right before deciding on
the peer (the same refresh `revoke_memory()` already does under the same lock).
The test no longer refreshes the cache from its hook: it only releases the
memory, asserts that the cached usage still exceeds the min memory at that
point, and asserts after the round that the second peer's query is not
cancelled and its usage was refreshed to 70 MiB with the min memory still at
100 MiB. Without the source change it fails on
`second_peer_query->is_cancelled()`.
##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -602,14 +603,33 @@ bool WorkloadGroupMgr::handle_process_memory_exceeded_(
return false;
}
+ // 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.
+ // Test the recorded reservation itself: this is the predicate the failed
reservation used,
+ // so the query is resumed exactly when its request fits now.
+ if
(!GlobalMemoryArbitrator::is_exceed_soft_mem_limit(query_it->reserve_size_)) {
+ LOG(INFO) << "Query: " <<
print_id(resource_ctx->task_controller()->task_id())
Review Comment:
Follow-up in b1e7d9c1280: keeping only the largest pending reservation still
woke the query when every request fit individually but not together (64 MiB + 4
MiB with 65 MiB headroom), after which the smaller one failed again and
restarted the wait. `add_paused_query()` now adds each pending reservation to
the recorded one. Each call site in `pipeline_task.cpp` stops the task
(`_spilling = true`) until the query is resumed, so one pause records exactly
one request per blocked task and the sum is the exact "all pending requests fit
at once" predicate. `process_mem_exceeded_keeps_all_pending_reservations` keeps
the query paused with 65 MiB headroom (each request fits on its own) and
resumes it at 69 MiB;
`process_mem_exceeded_pending_reservations_cancel_after_timeout` has headroom
for either request but not both and is cancelled at the timeout. Both fail
without the source change.
--
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]